ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

四层映射:从800连接到高并发的长连接网关路由架构实践

四层映射:从800连接到高并发的长连接网关路由架构实践 1. 项目背景与问题定义800 个连接背后的路由难题先交代一下这个项目的背景。我手上这套系统叫 WeClaw本质上是一个长连接网关服务负责把各种业务消息在服务端和设备端、服务端和前端页面之间做实时转发。最开始连接数只有几十个的时候路由逻辑随便写写就能跑无非就是一个 Map 存 session来消息了按 key 找 session 发出去。但连接数涨到 800 的时候事情开始不对劲消息偶尔延迟飙升、连接无征兆断开、内存占用肉眼可见地涨最头疼的是出现了消息发给了错误的连接这种低级但致命的 bug。你可能会问800 个连接对很多分布式系统来说根本不算大但关键在于这个网关跑在单机、内存受限的环境里而且消息类型极其复杂——有设备上报的数据、有客户端请求的应答、有服务端主动下发的指令还有广播通知。每一种消息的转发目标都不一样有的要发给单一连接有的要发给某个用户的所有连接有的要发给一组设备。如果每次转发都靠遍历 800 个连接去筛选目标CPU 扛不住如果只靠一个全局锁去保护共享状态并发又上不去。这时候就需要一套结构化的路由方案。BridgeConnectionManager 就是我在这个节点上设计出来的核心组件。它的思路很直接把连接管理这件事拆成四层映射每一层负责一个维度的路由查询。为什么是四层而不是三层或者五层这个后面展开。先记住一个核心结论这四层映射解决的是怎么用最短路径从消息到达目标连接的问题而 800 个连接只是这个方案的第一个压力验证场景方案本身是为更高并发预留的。这套东西适合谁看如果你在做 WebSocket 长连接服务、IM 系统、物联网设备接入网关或者任何需要维护大量实时连接的中间件这个四层映射的思路可以直接抄。即使你的连接规模只有一两百理解这套分层也能帮你提前避开不少坑。2. 四层映射的整体设计从连接到业务语义的逐级收敛先说一个通用的认知任何路由系统本质上都是在做从消息特征到目标地址的映射。WebSocket 场景的特殊之处在于连接是动态的随时有连接建立、断开、迁移而业务消息则是海量的、高频的。所以路由表必须支持高频查询和低频更新同时发生而且查询要快更新要稳。2.1 第一层连接注册表——最底层的物理连接索引这一层是任何连接管理器的地基。它维护的是连接唯一标识connectionId到 WebSocket Session 对象的映射。在 Java 实现里最自然的选型是ConcurrentHashMapString, WebSocketSession。为什么用 ConcurrentHashMap因为 JDK 8 以后的 ConcurrentHashMap 在并发读场景下基本无锁在写场景下只会锁住 hash 桶的头节点吞吐量足够支撑几千个连接的并发注册和注销。相比之下给整个 Map 加一把全局锁的Collections.synchronizedMap在连接数上来之后锁竞争会非常严重800 个连接同时心跳续期时就能明显感觉到延迟。这一层的 key 设计有讲究。我见过有人直接用用户 ID 当 key这在小规模 demo 里没问题但一旦一个用户开了多个设备比如手机 浏览器 桌面客户端后登录的就会挤掉先登录的消息永远只发到最后一个连接。所以这里必须用连接级唯一 ID比如用 UUID 或者设备ID 自增序列拼接保证每个 WebSocket 连接都有独立身份。// 连接注册核心逻辑 public class BridgeConnectionManager { // 第一层映射connectionId - WebSocketSession private final ConcurrentHashMapString, WebSocketSession sessionRegistry new ConcurrentHashMap(); public void registerConnection(String connectionId, WebSocketSession session) { sessionRegistry.put(connectionId, session); } public WebSocketSession getSession(String connectionId) { return sessionRegistry.get(connectionId); } public void unregisterConnection(String connectionId) { sessionRegistry.remove(connectionId); } }这层映射本身不复杂但它是后面三层的基础。所有查询最终都要落到这一层拿到 Session 才能发消息所以它的查询性能直接决定了端到端的延迟。后期我还为这一层加了一个最近活跃时间的辅助字段方便做心跳清理这个后面单独说。2.2 第二层身份映射表——把业务身份路由到连接集合第一层解决的是给我 connectionId我告诉你 Session 是哪条但业务代码往往不关心 connectionId它只关心给用户 10086 发的消息应该走哪些连接。这就是第二层映射要做的事业务身份到连接 ID 集合的映射。这里的业务身份可以是用户 ID、设备 ID、客户端类型具体取决于你的业务语义。在我的场景里一个用户可能持有多个连接所以 value 必须是一个连接集合。这里的数据结构选择很关键如果用普通的List连接注册和注销时需要遍历列表判断是否存在O(n) 的复杂度在连接频繁进出时不可接受如果用CopyOnWriteArrayList读性能好但写性能差因为每次添加删除都会复制整个数组。我的方案是外层ConcurrentHashMapString, SetStringvalue 使用ConcurrentHashMap.newKeySet()。这个 set 底层就是 ConcurrentHashMap读写都是分段的既能保证并发安全又不会像 CopyOnWrite 那样有写放大问题。// 第二层映射userId - SetconnectionId private final ConcurrentHashMapString, SetString identityIndex new ConcurrentHashMap(); public void bindIdentity(String identityId, String connectionId) { SetString connSet identityIndex.computeIfAbsent(identityId, k - ConcurrentHashMap.newKeySet()); connSet.add(connectionId); } public SetString getConnectionIdsByIdentity(String identityId) { SetString connSet identityIndex.get(identityId); return connSet null ? Collections.emptySet() : connSet; }这一层还有一个容易被忽略的点同一身份的多连接必须区分主备。比如用户在手机和 PC 上同时登录某些消息需要全端推送某些消息只推给其中一个端还有一些消息要求如果 PC 在线就一定走 PC。我的做法是在第二层映射里额外维护一个主连接标识绑定连接时根据客户端类型和登录时间判断谁当主连接路由时优先走主连接主连接不可用再降级到其他连接。2.3 第三层业务订阅表——从点对点到组播/广播的跃迁第二层解决的是点对点路由但真实业务大量存在组播需求给某个房间的所有在线用户推送消息、给某个设备的全部订阅方推送状态变更、给某个群组的所有成员广播通知。如果只在第二层做路由每条组播消息都要先查出所有成员 ID再逐个查他们的连接集合这相当于两次查询叠加而且业务层要自己维护成员有哪些这个信息违背了网关层做路由抽象的初衷。第三层映射解决这个问题频道/群组标识到连接 ID 集合的映射。我在实现中称之为订阅表subscription registry。当客户端发来一条订阅频道 xxx的消息时BridgeConnectionManager 会把当前 connectionId 加入该频道的订阅集合退订时反向操作。这层的设计要点有两个。第一订阅关系必须和连接生命周期强绑定——连接断开时不仅要清理第一层和第二层的记录还要把所有订阅表中含该 connectionId 的记录一并移除否则会出现幽灵订阅消息永远发向一个已经断开的连接。第二广播场景下要避免重复发送——如果同一个连接同时通过多个渠道比如既在群里又单独订阅了某个设备命中同一条消息必须在到达第三层时做一次去重判断否则客户端会收到重复消息。// 第三层映射channelId - SetconnectionId private final ConcurrentHashMapString, SetString channelSubscriptions new ConcurrentHashMap(); public void subscribe(String channelId, String connectionId) { channelSubscriptions.computeIfAbsent(channelId, k - ConcurrentHashMap.newKeySet()).add(connectionId); } public SetString getSubscribers(String channelId) { SetString subs channelSubscriptions.get(channelId); return subs null ? Collections.emptySet() : subs; } // 连接断开时的联动清理 public void onConnectionClosed(String connectionId) { // 遍历所有频道移除该连接反向索引优化后可以不用全表扫 channelSubscriptions.values().forEach(set - set.remove(connectionId)); }这段代码里有个明显的性能隐患连接断开时要遍历所有频道做移除频道多的时候 O(n*m) 是不可接受的。优化方案是加一个反向索引——每个 connectionId 记录它订阅了哪些频道断开时直接查这个反向索引只清理涉及的频道。这个优化在连接数 800 的场景就已经能看出效果连接频繁上下线时尤其明显。2.4 第四层消息路由层——组装前三层决策的路由规则前三层是存储结构它们回答了数据怎么存第四层是决策引擎它回答消息怎么走。这一层主要负责把消息头里的目标描述翻译成具体的目标连接列表是一个按优先级逐级匹配的过程。我的实现是把第四层设计成一组路由规则链每条规则匹配一种目标类型TargetType.SINGLE命中第一层直接把 connectionId 转成 Session。TargetType.IDENTITY命中第二层一个身份映射到多个连接。TargetType.CHANNEL命中第三层一个频道映射到一个连接集合。TargetType.BROADCAST直接全量遍历 sessionRegistry。消息进入网关时先解析出它的 targetType 和 targetId然后按类型走对应的查询逻辑。这里有一个关键的设计决策为什么我不把 targetType 的解析放在业务层而是收拢到网关层因为这样的话业务方只需要声明这条消息发给用户 10086 的在线连接网关负责所有映射细节业务方不需要知道任何连接管理内部结构。这个抽象边界让网关可以作为独立组件被多个业务复用。第四层还负责一个重要的兜底逻辑目标不可达时的降级策略。比如某个用户的所有连接都离线了这个消息是丢弃、进入离线消息队列还是返回错误我采用的是可配置策略默认丢进一个待投递队列等该用户重新上线时由连接建立事件触发补发。这一点也解释了为什么第二层映射不能只存在线连接——它天然可以衔接离线存储的逻辑。3. 核心实现细节同步机制、心跳清理与消息转发路径四层映射的存储结构定下来之后真正的挑战在于把这些结构接入真实的并发场景。这一章讲三个我在实现中反复打磨的细节并发读写的一致性、死连接的识别和清理、以及消息从进入到发出的完整路径。3.1 并发一致性为什么用 CopyOnWrite 思想处理连接集合的快照WebSocket 服务是典型的高并发读写场景连接随时被 local thread 注册/注销同时多个业务线程在读路由表发送消息。如果读操作拿到的是一个正在被修改的 Set直接遍历时会出现ConcurrentModificationException或者漏发消息的问题。一种粗暴的做法是所有读写都加锁但实测在 800 连接、每秒数千条消息的负载下锁竞争会导致 P99 延迟从 3ms 飙到 40ms。我的做法是读写分离的快照机制发送消息时先从路由表复制出一个只读快照然后在快照上遍历发送。虽然复制集合有一点开销但相比锁的代价小得多而且快照能保证一次广播中所有目标看到的是同一时刻的连接状态避免消息发送过程中连接列表发生变化导致的行为不一致。具体实现上我用List.copyOf()或者Set.copyOf()对目标集合做一次性快照然后遍历发送。快照的复制时间在这个规模下是微秒级的完全可接受。// 消息发送的统一入口快照语义 public void dispatch(OutboundMessage message) { SetString targets resolveTargets(message); // 复制快照避免发送过程中集合被修改 ListString snapshot new ArrayList(targets); for (String connId : snapshot) { WebSocketSession session sessionRegistry.get(connId); if (session ! null session.isOpen()) { sendMessage(session, message); } } }3.2 心跳机制与死连接清理800 连接下如何保持路由表干净长连接服务有一件绕不开的事TCP 层的连接断开不一定会被服务端立刻感知。客户端拔网线、断电、杀进程很多情况下服务端要等到发送数据时收到 RST 或者等到 TCP 超时才能发现。如果路由表里堆积了大量死连接每次组播都要对死连接做一次无效的写操作浪费 CPU 和内存严重时还会拖慢消息队列。我实现的方案是客户端心跳 服务端踢出的双向机制。客户端每 30 秒发送一次协议层心跳包不是 TCP keepalive而是应用层心跳服务端记录每个连接的最后心跳时间。然后有一个后台扫描任务每 10 秒扫一次全量连接把超过 90 秒没有心跳的连接判定为死亡执行完整的四层清理。这里有一个很重要的细节清理必须走统一的关闭流程而不是直接从第一层 remove。因为连接对象可能还持有未发送完的消息缓冲、注册过的事件监听器、以及第二三层映射的引用。只 remove 第一层会导致第二三层留下垃圾数据垃圾数据积累多了会引发更隐蔽的问题比如给已断开的连接发送消息发送失败后触发重试风暴。我的统一关闭流程是public void kickOff(String connectionId, CloseReason reason) { WebSocketSession session sessionRegistry.remove(connectionId); if (session null) return; // 清理第二层身份索引 identityIndex.forEach((identity, connSet) - connSet.remove(connectionId)); // 清理第三层频道订阅用反向索引优化 SetString subscribedChannels reverseChannelIndex.get(connectionId); if (subscribedChannels ! null) { for (String channel : subscribedChannels) { SetString subs channelSubscriptions.get(channel); if (subs ! null) subs.remove(connectionId); } } // 关闭连接 try { session.close(new CloseStatus(reason.getCode(), reason.getReason())); } catch (IOException e) { // 记录日志连接关闭失败通常意味着底层 TCP 已经断开忽略即可 } }心跳间隔的取值也要结合业务调整间隔太短会增加无谓的包开销太长则会让死连接在路由表中存活过久。我在这个项目里用的 30 秒心跳 90 秒超时是综合了客户端省电需求和服务器容忍度之后的选择。如果你们的客户端对电量不敏感也可以把心跳压到 15 秒这样死连接的发现时间能缩短一半。3.3 毫秒级转发路径一次完整消息的全过程拆解很多人在对比不同网关性能时只盯着单条消息的发送耗时但实际体验是端到端延迟——从客户端 A 发出消息到客户端 B 收到消息这中间要经过多少个环节我把自己系统的完整链路拆开来看客户端 A 发送消息WebSocket 协议解析netty 的 decoder 负责微秒级。消息进入业务处理线程池反序列化为领域对象约 0.1ms。调用 BridgeConnectionManager 的dispatch()方法解析 targetType 和 targetId约 0.02ms。根据 targetType 走对应映射查询得到目标连接 ID 集合快照约 0.05ms。遍历快照对每个 Session 执行sendMessage()写入 netty channel 的写出缓冲约 0.1ms 到 0.3ms取决于网络状况和写缓冲状态。netty 的 event loop 把缓冲数据真正写到 socket。从 3 到 5 是 BridgeConnectionManager 的职责我压测的结果是这部分的 P99 耗时在 1ms 以内通常在 0.3~0.6ms 之间。这里的耗时大头其实不在映射查询而在第 5 步对每个 session 的写出操作——因为每个 session 都关联到不同的 event loop 线程跨线程提交写任务有一定的调度开销。优化的一个关键技巧是批量写如果同一批消息要发给同一个连接的多条消息比如频道内连续两条通知可以在业务层合并成一次消息发送而不是逐条调用sendMessage()。这样既减少了跨线程调度次数也减少了 TCP 小包的数量一举两得。4. 800 连接压测实录从 63ms 到 3ms 的调优过程理论设计归理论真刀真枪的压测最能暴露问题。我搭了一个测试环境单机部署网关服务用 800 个模拟客户端建立 WebSocket 连接每个连接随机订阅 1~5 个频道同时往不同频道注入混合消息点对点、组播、广播各占一定比例观察延迟和吞吐。第一轮结果非常打脸P99 延迟 63ms平均延迟 8ms完全达不到预期的毫秒级标准。4.1 瓶颈一session 写操作的锁竞争用 JFRJava Flight Recorder采了一段时间的样本发现热点集中在 netty 写操作相关的锁上。原因是我的模拟客户端全部跑在同一台机器的不同线程里而服务端的 event loop 线程数默认是 CPU 核数测试机是 8 核所以只有 8 个 event loop 线程。800 个连接均匀分布在这 8 个线程上每个线程要处理 100 个连接的读写消息量大时写操作开始排队锁竞争自然严重。解决的方案不是简单地增加 event loop 线程数——线程多了上下文切换成本也会上去。我做了两件事一是把 event loop 线程数调到 16等于 CPU 超线程数二是把消息序列化操作移出 event loop全部放到业务线程池做让 event loop 只负责网络读写的原始字节流转。这个改动立竿见影P99 从 63ms 降到了 17ms。4.2 瓶颈二路由表查询中的无意遍历第二轮的性能数据是平均延迟 4ms、P99 17ms还算能看但离毫秒级还是差一口气。我继续采样发现耗时不光在写操作上还有一个隐藏较深的点频道广播消息的resolveTargets()里因为要组装多个频道的结果并对连接 ID 做去重我用了HashSet.addAll()这个操作本身没问题。问题出在我为了记录每个连接订阅了哪些频道的反向索引每次查询频道订阅集合时都顺带更新了访问时间这个更新操作有写锁导致并发查询时互相阻塞。这个案例很典型读路径中混入写操作是最隐蔽的性能杀手。解决办法是把统计访问频率这个需求和路由查询完全解耦——访问热度统计放到另一个独立的 TTL 缓存组件里做路由表本身只做纯粹的快照读。改完之后 P99 稳定在 5ms 以内大多数场景 P99 在 3ms 左右。4.3 最终结果与扩展性验证调优完成后的压测数据指标调优前调优后平均延迟8ms0.8msP99 延迟63ms3.2ms最大吞吐4200 条/秒15000 条/秒连接数800800内存占用210MB165MB另外我还做了一组 2000 连接的扩展验证虽然不在原需求范围内但四层映射的结构没有因为连接数翻倍而出现明显衰减P99 仍在 8ms 以内。这说明把 800 连接作为初始目标的设计容量是够的。5. 实操中的高频故障与排查技巧这一章是纯经验向的内容。四层映射的思路看着简单真正接入业务之后各种边界情况层出不穷。我整理了这段时间碰到的高频问题和对应的排查套路后面做类似系统的时候可以少走不少弯路。5.1 连接幽灵残留清理链路不完整现象是客户端明明已经断开但服务端监控面板上连接数迟迟不降广播消息也会莫名发给不存在的连接。排查链路先看第一层 sessionRegistry 的大小如果持续不减基本可以判断是关闭事件没走统一清理流程。比如某些异常路径直接调用session.close()而没有通知 BridgeConnectionManager或者 netty 的 channelInactive 事件在某些异常场景没有触发到业务层的监听器。我的处理方式是双保险第一在网关的抽象层包装一个GuardedSession它的close()方法强制走统一清理流程第二后台扫描定时任务除了检查心跳超时还要检查 session 的底层 channel 状态如果 channel 已经 inactive 但注册表还没清理立即补做清理。5.2 消息乱序与重复投递WebSocket 本身是顺序传输的但网关层一旦引入多线程和快照复制顺序就可能被破坏。比如同一个频道有三条消息线程池里的三个 worker 同时取到并并发发送由于每个 session 发送完成的时机不同客户端可能观察到第 2 条先于第 1 条到达。要根治乱序需要在消息进入网关时给每个连接维护一个发送序列号并且让同一个连接的消息串行化发送。我的方案是按连接维度做轻量级队列消息解析后不直接提交到共享线程池而是放入按 connectionId 哈希分桶的队列每个桶由一个独立线程消费保证同一个连接的消息顺序。至于重复投递最常见的原因是业务重试和断线重连后的快照不一致。收到重复消息的客户端如果幂等性做得不好会出现数据被重复写入。网关层的缓解措施是给消息加全局唯一 ID客户端可以用它做去重但这不是网关的职责边界建议业务层自行处理。5.3 反向索引的内存泄漏前面提到我用 reverseChannelIndex 来优化连接断开时的频道清理但没想到这个索引本身成了泄漏点。原因是在客户端发来退订消息和连接断开这两个场景中我的处理逻辑不一致退订消息走的是unsubscribe(channelId, connectionId)只从 channelSubscriptions 里移除却忘记更新 reverseChannelIndex而断开清理只依赖 reverseChannelIndex导致某些 connectionId 的反向记录残留越堆越多。定位到这个问题是在压测内存曲线时发现的——连接数没涨内存却稳步上升。修复合起来很简单在任何修改 channelSubscriptions 的地方都要同时维护 reverseChannelIndex两者的一致性必须在一个事务性方法里完成这个教训后来被我写进了代码评审 checklist。5.4 压测工具自身的瓶颈最后提醒一个容易被忽略的坑你用的 WebSocket 压测客户端本身可能成为瓶颈。我用开源工具跑 800 连接时客户端进程占用的 CPU 比服务端还高导致延迟数据不稳定。后来改成用轻量级脚本配合自研客户端打流数据才可信。压测时务必先确认客户端连接数确实建立满了、消息收发计数是准确的再开始分析服务端性能数据否则你调的可能不是服务端而是压测客户端。6. 一些值得记住的经验总结回看这个项目我觉得最有价值的部分不是四层映射本身而是把路由问题拆成多个独立维度的思考过程。每一层映射都只回答一个问题第一层回答连接在哪第二层回答身份对应哪些连接第三层回答频道如何组织连接第四层回答消息如何选择前三层。这种分层让每一层的代码都足够简单可以单独优化出了问题也能快速定位——消息延迟高先查第四层规则链连接断线清理慢先查反向索引更新逻辑而不是在几百行纠缠不清的路由代码里大海捞针。另一个体会是在写这类高并发组件时一定要把读路径写路径分离当成第一原则。我翻过两次车一次是把统计逻辑混进读路径导致读写竞争一次是清理逻辑和订阅逻辑没有保持一致性导致内存泄漏。都是因为某个小功能的便利性诱惑破坏了核心结构的纯粹性。后来我给自己定了一条规矩路由表的读路径不允许有任何写操作所有统计、审计、更新逻辑要么放在消息进入前要么放在消息发送完成后的独立环节。最后分享一个小技巧如果你也要做类似的长连接网关建议从设计第一天就给所有映射表加上完整的 metrics 埋点——每个 key 的查询次数、耗时分布、集合大小、命中率全部暴露成监控指标。这些数据在你扩容、调优、排查问题时是唯一的决策依据。我在 BridgeConnectionManager 里加了大约 20 个监控指标后期定位 5.3 节的内存泄漏时就是靠reverseChannelIndex 的条目数 vs 实际连接数这个指标的分叉发现的。这套四层映射方案目前已经稳定跑了一段时间后续我还在尝试把它扩展成支持多机部署的版本基本思路是在每层映射的外部再套一层一致性哈希把连接和订阅关系分布到多台机器上同时用消息总线同步路由表的变更事件。如果你也在做类似的方向欢迎在实际项目中验证这套分层思路有问题可以在评论区一起讨论。
返回列表