1. 项目背景与问题定义:800 个连接背后的路由难题
先交代一下这个项目的背景。我手上这套系统叫 WeClaw,本质上是一个长连接网关服务,负责把各种业务消息在服务端和设备端、服务端和前端页面之间做实时转发。最开始连接数只有几十个的时候,路由逻辑随便写写就能跑,无非就是一个 Map 存 session,来消息了按 key 找 session 发出去。但连接数涨到 800 的时候,事情开始不对劲:消息偶尔延迟飙升、连接无征兆断开、内存占用肉眼可见地涨,最头疼的是出现了"消息发给了错误的连接"这种低级但致命的 bug。
你可能会问,800 个连接对很多分布式系统来说根本不算大,但关键在于这个网关跑在单机、内存受限的环境里,而且消息类型极其复杂——有设备上报的数据、有客户端请求的应答、有服务端主动下发的指令,还有广播通知。每一种消息的转发目标都不一样,有的要发给单一连接,有的要发给某个用户的所有连接,有的要发给一组设备。如果每次转发都靠遍历 800 个连接去筛选目标,CPU 扛不住;如果只靠一个全局锁去保护共享状态,并发又上不去。这时候就需要一套结构化的路由方案。
BridgeConnectionManager 就是我在这个节点上设计出来的核心组件。它的思路很直接:把"连接管理"这件事拆成四层映射,每一层负责一个维度的路由查询。为什么是四层而不是三层或者五层?这个后面展开。先记住一个核心结论:这四层映射解决的是"怎么用最短路径从消息到达目标连接"的问题,而 800 个连接只是这个方案的第一个压力验证场景,方案本身是为更高并发预留的。
这套东西适合谁看?如果你在做 WebSocket 长连接服务、IM 系统、物联网设备接入网关,或者任何需要维护大量实时连接的中间件,这个四层映射的思路可以直接抄。即使你的连接规模只有一两百,理解这套分层也能帮你提前避开不少坑。
2. 四层映射的整体设计:从连接到业务语义的逐级收敛
先说一个通用的认知:任何路由系统,本质上都是在做"从消息特征到目标地址"的映射。WebSocket 场景的特殊之处在于,连接是动态的,随时有连接建立、断开、迁移,而业务消息则是海量的、高频的。所以路由表必须支持高频查询和低频更新同时发生,而且查询要快,更新要稳。
2.1 第一层:连接注册表——最底层的物理连接索引
这一层是任何连接管理器的地基。它维护的是连接唯一标识(connectionId)到 WebSocket Session 对象的映射。在 Java 实现里,最自然的选型是ConcurrentHashMap<String, 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 ConcurrentHashMap<String, 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,读性能好但写性能差,因为每次添加删除都会复制整个数组。
我的方案是:外层ConcurrentHashMap<String, Set<String>>,value 使用ConcurrentHashMap.newKeySet()。这个 set 底层就是 ConcurrentHashMap,读写都是分段的,既能保证并发安全,又不会像 CopyOnWrite 那样有写放大问题。
// 第二层映射:userId -> Set<connectionId> private final ConcurrentHashMap<String, Set<String>> identityIndex = new ConcurrentHashMap<>(); public void bindIdentity(String identityId, String connectionId) { Set<String> connSet = identityIndex.computeIfAbsent(identityId, k -> ConcurrentHashMap.newKeySet()); connSet.add(connectionId); } public Set<String> getConnectionIdsByIdentity(String identityId) { Set<String> connSet = identityIndex.get(identityId); return connSet == null ? Collections.emptySet() : connSet; }这一层还有一个容易被忽略的点:同一身份的多连接必须区分主备。比如用户在手机和 PC 上同时登录,某些消息需要全端推送,某些消息只推给其中一个端,还有一些消息要求"如果 PC 在线就一定走 PC"。我的做法是在第二层映射里额外维护一个"主连接"标识,绑定连接时根据客户端类型和登录时间判断谁当主连接,路由时优先走主连接,主连接不可用再降级到其他连接。
2.3 第三层:业务订阅表——从"点对点"到"组播/广播"的跃迁
第二层解决的是点对点路由,但真实业务大量存在组播需求:给某个房间的所有在线用户推送消息、给某个设备的全部订阅方推送状态变更、给某个群组的所有成员广播通知。如果只在第二层做路由,每条组播消息都要先查出所有成员 ID,再逐个查他们的连接集合,这相当于两次查询叠加,而且业务层要自己维护"成员有哪些"这个信息,违背了网关层做路由抽象的初衷。
第三层映射解决这个问题:频道/群组标识到连接 ID 集合的映射。我在实现中称之为订阅表(subscription registry)。当客户端发来一条"订阅频道 xxx"的消息时,BridgeConnectionManager 会把当前 connectionId 加入该频道的订阅集合;退订时反向操作。
这层的设计要点有两个。第一,订阅关系必须和连接生命周期强绑定——连接断开时,不仅要清理第一层和第二层的记录,还要把所有订阅表中含该 connectionId 的记录一并移除,否则会出现幽灵订阅,消息永远发向一个已经断开的连接。第二,广播场景下要避免重复发送——如果同一个连接同时通过多个渠道(比如既在群里又单独订阅了某个设备)命中同一条消息,必须在到达第三层时做一次去重判断,否则客户端会收到重复消息。
// 第三层映射:channelId -> Set<connectionId> private final ConcurrentHashMap<String, Set<String>> channelSubscriptions = new ConcurrentHashMap<>(); public void subscribe(String channelId, String connectionId) { channelSubscriptions.computeIfAbsent(channelId, k -> ConcurrentHashMap.newKeySet()).add(connectionId); } public Set<String> getSubscribers(String channelId) { Set<String> 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) { Set<String> targets = resolveTargets(message); // 复制快照,避免发送过程中集合被修改 List<String> 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)); // 清理第三层:频道订阅(用反向索引优化) Set<String> subscribedChannels = reverseChannelIndex.get(connectionId); if (subscribedChannels != null) { for (String channel : subscribedChannels) { Set<String> 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 写操作的锁竞争
用 JFR(Java 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 最终结果与扩展性验证
调优完成后的压测数据:
| 指标 | 调优前 | 调优后 |
|---|---|---|
| 平均延迟 | 8ms | 0.8ms |
| P99 延迟 | 63ms | 3.2ms |
| 最大吞吐 | 4200 条/秒 | 15000 条/秒 |
| 连接数 | 800 | 800 |
| 内存占用 | 210MB | 165MB |
另外我还做了一组 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 实际连接数"这个指标的分叉发现的。
这套四层映射方案目前已经稳定跑了一段时间,后续我还在尝试把它扩展成支持多机部署的版本,基本思路是在每层映射的外部再套一层一致性哈希,把连接和订阅关系分布到多台机器上,同时用消息总线同步路由表的变更事件。如果你也在做类似的方向,欢迎在实际项目中验证这套分层思路,有问题可以在评论区一起讨论。