前阵子有个朋友的公司做在线客服系统,单机跑了一年多相安无事。结果有次渠道推广爆量,WebSocket 在线连接数直接冲上五位数,服务开始频繁卡顿、内存飙高。他们第一反应是加机器,结果加了机器不但没缓解,反而冒出更诡异的 bug:用户明明在线,消息却经常推送不到,发一条消息别人要隔好几秒才收到,甚至干脆收不到。
问题就出在 WebSocket 的本质上——它是一条有状态的 TCP 长连接。HTTP 请求处理完就断开了,随便负载均衡到哪台机器都行;WebSocket 一旦握手成功,这个连接就"长"在了某一台节点上,后续所有消息都只能由那台节点转发。你可以在 Redis 里存用户的 session 数据,但没法把一条已经建立的 TCP 连接瞬移到另一台机器上。这就是所有 WebSocket 集群方案的出发点。
这篇文章我会从单机服务的容量边界讲起,再逐步拆解三种主流的集群方案:Sticky Session 粘连、Redis Pub/Sub 消息路由、MQ 推送服务化架构,最后给出一套可以直接落地的代码骨架和我在生产环境踩过的坑。不管你是做在线客服、消息推送、聊天室还是实时协作,这套思路都能直接套用。
1. WebSocket 的"有状态"本质:所有集群麻烦的根源
1.1 一条 TCP 长连接不是一份可以随意路由的报文
WebSocket 的握手过程和 HTTP 很像,本质是借助 HTTP Upgrade 机制完成协议升级:
GET /ws/chat HTTP/1.1 Host: im.example.com Upgrade: websocket Connection: Upgrade Sec-WebSocket-Key: x3JJHMbDL1EzLkh9GBhXDw== Sec-WebSocket-Version: 13关键是这次握手发生在哪台机器上,后续一整条双向通道就跟死在哪台机器上。因为 WebSocket 是长连接,服务端内存里保存着这个连接的 session 对象、读写缓冲区、各种状态标记,这些都不是 Redis 里存一个字符串就能搬走的东西。TCP socket 四元组绑定的是某一台机器的某个端口,数据报文只会被投递到那个 socket 上,其他节点根本没有这个连接的任何信息。
这就导致了一个很尴尬的现状:你可以在 Redis 里存用户的登录态,没法把连接本身"同步"给所有节点。负载均衡可以把 HTTP 请求均匀分发到每台机器,但对 WebSocket 来说,连接一旦建立,它的归属就固定了。
1.2 HTTP 无状态与 WebSocket 有状态的对比
HTTP 的每个请求都是独立的,服务端处理完就丢,不存任何客户端的状态。所以 Nginx 后面挂 10 台机器和挂 1 台机器,对业务代码来说没有区别,随便怎么轮询都行。
WebSocket 完全不同,它在一次 TCP 连接上建立了全双工通道,而且这个通道是长久的。服务端必须要记住"这个 userId 对应哪个 session,这个 session 连在哪个 socket 上",否则收到消息不知道往哪里发。这就是状态,而状态是分布式系统最棘手的敌人。
| 对比维度 | HTTP 短连接 | WebSocket 长连接 |
|---|---|---|
| 连接生命周期 | 请求结束即断开 | 一直保持,直到双方关闭 |
| 服务端是否保存连接状态 | 一般不保存 | 必须保存 session 引用 |
| 负载均衡策略 | 随便轮询、随机、加权 | 连接一旦建立就不能迁移 |
| 故障转移 | 请求重发即可 | 连接断开需要客户端重连 |
| 集群复杂度 | 低,天然横向扩展 | 高,需要额外的路由机制 |
很多人第一次做 WebSocket 集群时,会下意识地按 HTTP 的思路去设计:前边挂 Nginx,后边挂一堆应用节点,结果一上线就发现消息发不出去,然后才开始理解"有状态"意味着什么。
1.3 先定义业务场景再选后续方案
做技术选型前,我建议先想清楚业务场景,因为不同的场景对集群方案的要求差别很大。我见过最典型的几类:
- 服务端主动推送:比如库存变动通知、订单状态推送,数据源在服务端,客户端被动接收。这类场景广播和点对点都要用,但对消息可靠性要求相对宽松。
- 在线客服 / IM 聊天:消息是双向的,用户和客服之间的消息要准确投递,不能丢,顺序还不能乱。这类场景对点对点路由要求很高,通常要配套离线消息。
- 聊天室 / 直播弹幕:重点是广播能力,一个房间的消息要推给房间内所有人,而且量大、实时性要求高。这类场景要特别注意广播风暴问题。
- 实时协同编辑:比如白板、在线文档,消息频率高、延迟敏感,而且要求多端状态一致。
不同场景决定了你后面是用 Redis Pub/Sub 就够了,还是必须上 MQ,甚至需要单独做一个推送网关。这些决策在单机阶段看不出来,但等到集群阶段就全是债。
2. 单机服务:先把容量边界和连接管理做扎实
2.1 一台机器到底能扛多少连接
很多人的第一个误区是"WebSocket 很重,一台机器扛不了多少连接"。其实恰恰相反,WebSocket 的协议开销非常小,真正吃资源的是每个连接占用的文件描述符、内核 socket 缓冲区和应用层 buffer。
Linux 下 WebSocket 服务底层走的是 epoll 模型,百万并发连接在理论上是可以做到的,但实际业务环境远达不到。我通常这样估算:
- 每个空闲连接在应用层大约占用 20KB~50KB 内存(包括 session 对象、读缓冲区、写缓冲区)。
- 10 万在线连接大约需要 2GB~5GB 内存,这是纯连接的消耗,还不算业务对象。
- CPU 消耗主要来自心跳包的编解码和消息的序列化,空闲连接几乎不占 CPU。
所以一台 8C16G 的机器,跑 5 万在线连接、每秒几千条消息,一般绰绰有余;跑到 20 万以上就要认真调优了。
单机部署前一定要改几个 Linux 内核参数,这是最容易被忽略的:
# 调整文件描述符上限 ulimit -n 1048576 # 内核层面提升连接队列长度 net.core.somaxconn = 65535 net.ipv4.tcp_max_syn_backlog = 65535 # 加大本地端口范围,防止大量短连接耗尽端口 net.ipv4.ip_local_port_range = 1024 65535 # 加快 TIME_WAIT 回收 net.ipv4.tcp_fin_timeout = 15不调文件描述符上限的话,默认 1024 的 ulimit 会直接卡死在连接数上,业务代码写得再漂亮也没用。
2.2 连接管理的核心数据结构
单机模式下,连接管理的核心就是在内存里维护一张"用户 ID 到 WebSocketSession"的映射表。Java 里用 ConcurrentHashMap,Go 里用 sync.Map,本质都一样:
@Component public class SessionRegistry { // userId -> WebSocketSession private final ConcurrentHashMap<String, WebSocketSession> sessions = new ConcurrentHashMap<>(); public void register(String userId, WebSocketSession session) { sessions.put(userId, session); } public void unregister(String userId) { sessions.remove(userId); } public WebSocketSession get(String userId) { return sessions.get(userId); } public int count() { return sessions.size(); } }注意这里有个细节:注册的 key 是业务用户 ID,不是 session ID。因为你的消息投递是面向用户的,而不是面向连接的。如果同一个用户开多个标签页、多台设备,那就需要建立 userId 到一组 session 的映射,也就是 Map<String, Set >,广播给这个用户所有端。这个在 IM 场景下属于刚需,做单机的时候就要留好这个设计余地。
2.3 心跳与死连接清理
单机模式最容易踩的坑是"连接泄漏"。客户端断网、拔网线、电脑休眠,TCP 层不一定能及时感知,尤其是有中间 NAT 设备的情况下,连接会一直挂在那里,变成半开连接(half-open)。如果不做处理,这些死连接会一直占着内存和文件描述符,最终把服务拖垮。
解决方式就是应用层心跳。我常用的方案是:服务端每 30 秒下发一个 Ping 帧,客户端收到后回 Pong 帧;服务端如果连续 3 次(90 秒)没收到某个连接的 Pong,就判定它已死亡,主动关闭并清理 session。
public void startHeartbeatCheck() { scheduledExecutor.scheduleAtFixedRate(() -> { long now = System.currentTimeMillis(); sessionRegistry.getAll().forEach((userId, session) -> { long idleTime = now - lastPongTime(userId); if (idleTime > 90_000) { session.close(CloseStatus.SESSION_NOT_RELIABLE); sessionRegistry.unregister(userId); } else if (idleTime > 30_000) { session.sendMessage(new PingMessage()); } }); }, 10, 10, TimeUnit.SECONDS); }这里间隔不是随便定的。间隔太短(比如 5 秒),心跳包会占用大量带宽和 CPU,10 万连接每秒光心跳就是几万条消息;间隔太长(比如 5 分钟),死连接清理不及时,连接数虚高。30 秒心跳、90 秒判定死亡是我在多个项目里验证过比较稳的参数。
2.4 单机模式最常见的坑
除了心跳,单机还会遇到几个高频问题:
- Nginx 默认超时:如果前面挂了 Nginx 做反向代理,默认 proxy_read_timeout 是 60 秒,WebSocket 连接空闲超过 60 秒就会被 Nginx 掐断。解决方式是显式设置较大的超时时间:proxy_read_timeout 3600s。
- 在事件循环里做阻塞操作:很多 WebSocket 框架是基于 Netty 或类似的事件循环模型,如果你在消息处理器里直接调用远程 API、查数据库,会阻塞事件线程,导致整个服务吞吐暴跌。正确做法是把耗时的操作丢到业务线程池,或者用异步方式处理。
- 单机依赖单点:连接全在一台机器上,进程崩溃、机器重启,所有连接瞬间全部断开。这个无解,只能靠集群解决。
单机做扎实的意义在于:它帮你把连接管理的细节(注册、心跳、清理、投递)都理清楚了,这些逻辑在集群模式下会被复用,而不会白白浪费。
3. 集群第一板斧:Sticky Session 粘滞与网关层配置
3.1 ip_hash 与 cookie 粘滞的原理
先说说最简单的一种集群思路:让同一个用户的连接总是被负载均衡到同一台后端节点。
Nginx 的 ip_hash 算法会根据客户端 IP 做哈希,同一个 IP 的请求会被分配到同一个 upstream 节点:
upstream ws_backend { ip_hash; server 192.168.1.10:8080; server 192.168.1.11:8080; server 192.168.1.12:8080; } server { listen 80; location /ws { proxy_pass http://ws_backend; # WebSocket 升级必需 proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; # 长连接超时放宽 proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }这样同一个客户端 IP 的所有 WebSocket 连接都会落在同一台节点上。对于用户量不大、节点不多的场景,这个方案能解决大部分问题,而且改动量最小。
还有一个变体是基于 cookie 的粘滞,Nginx 会下发一个带后端节点标识的 cookie,后续请求根据 cookie 直接路由到指定节点。相比之下 cookie 方案比 ip_hash 更精确,因为同一个 NAT 后面的多个用户 IP 相同,ip_hash 会把他们都砸到同一台节点上,而 cookie 方案能区分开。
3.2 粘滞方案的优势与天花板
Sticky Session 的优势非常明显:零业务改造,不需要额外引入 Redis 或 MQ,连接在哪个节点就是哪个节点,点对点消息直接查本地 session 表就能发。
但它的天花板也很低:
- 节点故障就是灾难:如果一台节点宕机,粘滞在这台机器上的所有连接全部断开,客户端需要重新握手,但此时 ip_hash 仍然会把它们路由到同一台(已经宕机的)节点,直到 Nginx 把该节点摘除。即使摘除了,其它节点上也没有这些用户的 session,必须靠客户端重新注册。
- 负载不均衡:某个 IP 段用户量大时,ip_hash 可能把大量连接堆在同一台节点上,其他节点空闲。这就是"加了机器反而没效果"的经典原因之一。
- 无法解决广播问题:要向所有用户广播消息时,需要遍历所有节点上的所有连接。粘滞方案里每个节点只知道自己本地的连接,你仍然需要一套机制把广播消息分发到每个节点。
所以我的结论是:Sticky Session 只适合作为"最小可用集群方案",一般在项目初期、用户量不大、对可用性要求不高的场景使用。它不能算真正意义上的集群解决方案,只能算负载均衡策略。
3.3 网关层必须处理的细节
不管用不用粘滞,网关层(Nginx)的 WebSocket 配置都有几个必须处理的细节,很多人在这里踩坑:
location /ws { proxy_pass http://ws_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_connect_timeout 60s; proxy_read_timeout 3600s; proxy_send_timeout 3600s; proxy_buffer_size 64k; proxy_buffers 8 64k; }几个关键点:
proxy_http_version 1.1必须设置,WebSocket Upgrade 依赖 HTTP/1.1 的持久连接特性。proxy_set_header Upgrade和Connection "upgrade"是协议升级的关键,不设置的话 Nginx 不会转发 Upgrade 头,WebSocket 握手直接失败。proxy_read_timeout和proxy_send_timeout必须调大,否则空闲连接会被 Nginx 掐掉。proxy_buffer_size关系到 WebSocket 帧的缓冲区大小,如果消息体比较大(比如超过 64KB),需要同步调大 buffer,否则会出现消息截断或报错。
另外,如果集群里走的是 HTTP/2,要注意 Nginx 对 WebSocket over HTTP/2 的支持在较老版本里不完善,生产环境建议 WebSocket 走独立的 HTTP/1.1 监听端口,和普通 HTTPS 业务分开。
4. 集群第二板斧:Redis Pub/Sub 做跨节点消息路由
4.1 核心思路:本地注册表 + 全局路由表
Sticky Session 解决了连接归属问题,但没有解决跨节点消息路由的问题。真正通用的做法是引入一层全局路由信息:用 Redis 保存"用户 ID 落在哪个节点"的映射关系,节点间通过 Redis Pub/Sub 互相通信。
这个方案的核心思路是:
- 每个节点在本地内存维护自己的 session 注册表,只保存连到本节点的连接。
- 每个节点启动时生成一个全局唯一的 nodeId,并把自己注册到 Redis。
- 用户连接建立时,把 userId -> nodeId 的映射写入 Redis,并定期续期。
- 节点间消息传递走 Redis Pub/Sub,发送节点把消息发布到指定节点的频道,目标节点的订阅者收到后查本地 session 表并推送。
这样每个节点都不需要知道其他节点的完整连接信息,只需要知道"目标用户在哪台节点",剩下的投递动作由目标节点本地完成。
Redis 的 key 设计我习惯这么搞:
ws:user:{userId} -> nodeId # 全局路由表,TTL 90 秒,心跳续期 ws:nodes -> Set<nodeId> # 存活节点列表 ws:msg:{nodeId} -> 消息 # 点对点消息频道,发往指定节点 ws:broadcast -> 消息 # 广播消息频道,所有节点订阅每个节点订阅两个频道:ws:msg:{自己的nodeId}和ws:broadcast。这样点对点和广播就用一套机制统一处理了。
4.2 点对点消息的完整链路
点对点消息的完整链路是这样的:
- 业务服务想给用户 U 推送一条消息。
- 查询 Redis
ws:user:U得到目标节点 nodeId。 - 如果 nodeId 就是本节点,直接从本地 session 表查连接并发送。
- 如果不是本节点,把消息发布到 Redis 频道
ws:msg:{nodeId}。 - 目标节点的订阅者收到消息,查询本地 session 表,找到连接后发送。
用 Java 伪代码表示大概是这个样子:
public void sendToUser(String userId, String payload) { String nodeId = stringRedisTemplate.opsForValue().get("ws:user:" + userId); if (nodeId == null) { // 用户不在线,走离线消息逻辑 handleOfflineMessage(userId, payload); return; } if (nodeId.equals(localNodeId)) { // 本节点直接投递 WebSocketSession session = sessionRegistry.get(userId); if (session != null && session.isOpen()) { session.sendMessage(new TextMessage(payload)); } } else { // 跨节点路由:发布到目标节点频道 WsRouteMessage routeMsg = new WsRouteMessage(userId, payload); stringRedisTemplate.convertAndSend("ws:msg:" + nodeId, routeMsg.toJson()); } }订阅端统一处理:
@Component public class WsMessageSubscriber extends AbstractMessageListener { @Override public void onMessage(Message message, byte[] pattern) { WsRouteMessage msg = JSON.parseObject(message.getBody(), WsRouteMessage.class); if (msg.isBroadcast()) { // 广播消息:遍历本地所有连接发送 sessionRegistry.getAll().forEach((userId, session) -> { if (session.isOpen()) { session.sendMessage(new TextMessage(msg.getPayload())); } }); return; } // 点对点消息:查本地 session WebSocketSession session = sessionRegistry.get(msg.getTargetUserId()); if (session != null && session.isOpen()) { session.sendMessage(new TextMessage(msg.getPayload())); } } }这套机制的核心好处是:每个节点只保存自己的连接,路由信息收敛到 Redis 里,节点可以随时水平扩展,新节点上线只需要订阅自己的频道即可。
4.3 广播消息的完整链路
广播消息的处理比点对点简单:发送方直接往ws:broadcast频道发布消息,所有节点订阅后各自往本地连接推送。
但这里有一个隐蔽的问题:如果广播的接收者是"某个聊天室的所有人"而不是"所有在线用户",那么每个节点在收到广播后,还需要判断哪些本地连接属于这个聊天室。这就需要在本地 session 表之外,再维护一个"聊天室 -> 成员连接"的映射表。
// 房间 -> userIds private final ConcurrentHashMap<String, Set<String>> roomMembers = new ConcurrentHashMap<>();连接建立时根据客户端带上来的参数(比如 URL query 里的 roomId)把 userId 加入对应房间;连接销毁时从所有房间移出。广播给房间时,节点拿到房间成员列表,再逐个查 session 表发送。这样广播的范围就被限制在目标房间内,而不是全量广播。
4.4 Redis Pub/Sub 的边界:消息丢失与不可回溯
Redis Pub/Sub 有个非常重要的特性:消息不持久化。发布者把消息发出去,如果此时某个订阅者恰好不在线(节点宕机、网络抖动),这条消息就永久丢失了。Redis 不会像 MQ 那样帮你把消息存起来等消费者恢复后再投递。
所以 Redis Pub/Sub 方案只适用于"消息实时投递、丢了也无所谓或可以从业务侧补偿"的场景。比如在线状态推送、心跳类通知、弹幕这类实时性消息。如果消息不能丢——比如聊天记录、订单通知——就必须在业务层做持久化和补偿,或者直接换用 MQ 方案。
另外,Redis Pub/Sub 的广播是 push 模型,如果某个节点处理消息过慢,会导致该节点的订阅者积压甚至断连。如果消息量非常大,建议结合以下做法:
- 把广播消息按业务维度切分到多个 channel,比如
ws:broadcast:room:{roomId},只有需要接收的节点才订阅,减少无效消息传输。 - 在节点本地用队列缓冲收到的消息,再由独立线程池发送,避免阻塞 Redis 订阅线程。这个是生产环境很重要的优化点,后面踩坑部分会细说。
5. 集群第三板斧:消息队列与推送服务化架构
5.1 什么时候必须上 MQ
Redis Pub/Sub 有消息丢失的硬伤,对于 IM、客服系统、交易通知这类要求不丢消息的场景,就需要引入消息队列(RabbitMQ、Kafka、RocketMQ 等)。
我判断是否需要上 MQ 的几条标准:
- 消息不能丢:比如用户聊天记录、支付结果通知,丢失会引发资损或客诉。
- 需要削峰填谷:比如秒杀场景,服务端瞬间产生大量推送消息,直接打到 WebSocket 连接上会把节点打挂;MQ 可以做流量缓冲。
- 需要离线消息:用户不在线时消息要持久化,等用户上线后再补推。
- 需要消息有序性:IM 场景里同一个聊天窗口的消息必须有序,Redis Pub/Sub 做不到精细化的顺序保证,而 MQ 可以按 key 分区保证局部有序。
5.2 事件驱动设计
引入 MQ 之后,架构从"节点间互相路由"演进成"事件驱动 + 独立的推送层"。整个链路变成:
业务服务 --生产--> MQ Exchange --路由--> Queue --消费--> WebSocket 节点 --推送--> 客户端业务服务不再直接关心目标用户在哪个节点,它只需要把消息投递到 MQ 对应的队列。WebSocket 节点作为消费者监听队列,拿到消息后查本地 session 表并推送。
这个设计最大的好处是业务逻辑和连接管理彻底解耦。订单服务不需要关心用户当前连在哪台机器上,它只管发消息;WebSocket 节点只管消费和推送。节点可以随时扩缩容,对业务方完全透明。
从 Redis Pub/Sub 迁移到 MQ 时,之前那套路由表依然有用——节点消费到消息后,仍然需要查ws:user:{userId}判断目标用户是否在本节点。不过这里有个优化点:可以让 MQ 按目标节点做分区,比如把消息路由到指定节点的专用队列,这样每个节点只消费自己需要处理的消息,减少无效消费。
5.3 离线消息与消息补偿
有了 MQ,离线消息就好处理了。用户不在线时,消息先落库(或者存储在 Redis 里),用户重新建立 WebSocket 连接后,服务端从存储中拉取该用户的离线消息补推。
这里有一个我趟过的坑:离线消息的补推不能一股脑全推。用户断线十分钟可能积压几百条消息,一次性推过去不仅客户端渲染卡顿,还会触发大量 ACK 回执,反而把刚刚恢复的连接打挂。正确做法是分批补推,比如每次推 20 条,等客户端确认后再推下一批。
另外,消息补偿机制要考虑幂等。客户端收到消息后可能会回执,服务端重推时要有去重逻辑,否则用户会看到重复消息。通常用消息 ID 做幂等键,客户端按 ID 去重,服务端按 ID 记录已推送游标。
MQ 方案也有新的问题要处理:消费者宕机恢复后从哪个 offset 开始消费?超时未 ACK 的消息是否会重复投递?这些属于 MQ 使用的基础问题,这里不展开,但一定要在设计方案时提前想好。
6. 实战拆解:一套可落地的 WebSocket 集群骨架
6.1 整体架构布局
把前面的方案结合起来,一套比较完整的 WebSocket 集群架构长这样:
- 接入层:Nginx / SLB 做负载均衡,不配置粘滞,WebSocket 握手随机分发到任意节点。因为引入路由层后,连接落在哪台节点已经无所谓了。
- 连接层:N 个 WebSocket 应用节点,每个节点维护自己的本地 session 表,并注册到 Redis。
- 路由层:Redis 保存 userId -> nodeId 映射;节点间通过 Pub/Sub 通信。
- 消息层:RabbitMQ 承担可靠消息投递,业务系统通过 MQ 解耦。
- 存储层:MySQL 存聊天记录等需要持久化的数据;Redis 同时承担路由表和热数据缓存。
这个架构的好处是每一层都可以独立扩展:连接多了加 WebSocket 节点,消息量大了扩 MQ 分区,路由表性能不够就升级 Redis 集群。
6.2 连接注册与注销流程
连接建立时的完整流程:
- 客户端发起 WebSocket 握手,Nginx 转发到任意一个应用节点。
- 节点在 onOpen 回调里拿到 userId(从 token 或 URL 参数解析)。
- 把 WebSocketSession 注册到本地 session 表。
- 把
ws:user:{userId}写入 Redis,值为当前节点 nodeId,TTL 90 秒。 - 开启该连接的心跳监控。
连接断开时的清理流程:
- 客户端主动关闭或心跳超时判定死亡后,触发 onClose 回调。
- 从本地 session 表移除该 userId。
- 删除 Redis 里的
ws:user:{userId}(先比对 nodeId 是否为本节点,防止误删)。 - 从所有房间成员表里移除该 userId。
- 把离线消息标记为待补推状态。
这里有个细节:删除 Redis 路由表时要带上 nodeId 做条件删除。因为可能用户刚断线,又在新节点上建立了新连接并写入了新的路由表,此时旧节点如果直接 del,会把新路由信息也删掉,导致消息路由失败。用 Lua 脚本或者 compare-and-delete 都能解决。
6.3 节点上下线与故障转移
节点的健康检查一般在 Redis 里做,每个节点启动时把自己的 nodeId 写进一个有序集合,并周期性地更新心跳时间戳:
// 节点心跳上报,每 30 秒一次 stringRedisTemplate.opsForZSet().add( "ws:nodes", localNodeId, System.currentTimeMillis() );其他节点或监控服务定期扫描这个有序集合,把心跳时间超过 90 秒的 nodeId 判定为宕机节点,从集合里移除,并帮他做善后工作:
- 广播"节点 xxx 已下线"的通知,各节点清理本地可能存在的该节点的关联信息。
- 该节点持有的所有用户连接被动断开,客户端通过重连机制落到其他节点。
- 如有必要,从存储层拉取这些用户的会话状态,重新路由。
节点故障时客户端重连是不可避免的,但一定要做重连保护,否则会触发重连风暴。我常用的策略是:指数退避 + 随机抖动。客户端第一次重连等 1 秒,之后 2 秒、4 秒、8 秒……最大 30 秒封顶,每次重连时间加一个 0~30% 的随机抖动,避免所有客户端同时重连压垮网关。
6.4 关键代码骨架
最后给一个完整的 WebSocket 连接注册 + 跨节点路由的最小骨架,基于 Spring Boot 和 Redis:
@ServerEndpoint("/ws/{userId}") @Component public class WsEndpoint { @OnOpen public void onOpen(Session session, @PathParam("userId") String userId) { // 1. 本地注册 SessionRegistry.register(userId, session); // 2. 全局路由表注册,TTL 90s,心跳续期 RedisUtil.set("ws:user:" + userId, LocalNode.getNodeId(), 90); // 3. 订阅当前节点的消息频道(只订阅一次) RedisSubscriber.subscribe("ws:msg:" + LocalNode.getNodeId()); // 4. 启动心跳任务 HeartbeatManager.start(session, userId); } @OnClose public void onClose(@PathParam("userId") String userId) { SessionRegistry.unregister(userId); RedisUtil.compareAndDelete("ws:user:" + userId, LocalNode.getNodeId()); RoomManager.removeFromAllRooms(userId); } @OnMessage public void onMessage(String message, @PathParam("userId") String userId) { // 业务消息处理,投递到 MQ 或直接路由 ChatService.handleUserMessage(userId, message); } @OnError public void onError(Session session, Throwable error) { // 记录日志,连接由心跳机制兜底清理 log.error("ws error", error); } }消息路由服务:
@Service public class MessageRouter { public void sendToUser(String userId, String payload) { String targetNode = RedisUtil.get("ws:user:" + userId); if (targetNode == null) { offlineMessageStore.save(userId, payload); return; } if (targetNode.equals(LocalNode.getNodeId())) { SessionRegistry.sendToUser(userId, payload); } else { RedisUtil.publish("ws:msg:" + targetNode, payload); } } public void broadcastToRoom(String roomId, String payload) { RedisUtil.publish("ws:broadcast:room:" + roomId, payload); } }这套骨架可以直接跑通单机和集群两种模式:单机运行时,所有连接都注册到唯一的节点上,sendToUser 直接走本地发送;集群运行时,靠 Redis 路由表自动切换到跨节点发布。业务层几乎不需要改动。
7. 生产环境踩坑记:心跳、连接泄漏与广播风暴
7.1 心跳间隔不当引发的"雪崩重连"
有次我在压测环境发现一个诡异现象:在线连接数稳定在 5 万左右,但每分钟都有大量连接断开重连,服务端日志里全是"session closed"和"new connection"。
排查了很久才发现问题出在心跳参数上。当时的配置是 10 秒发一次 Ping、30 秒判定死亡,但网关层 Nginx 的超时设置只有 20 秒。Nginx 在超过 20 秒没有收到任何数据时会主动关闭连接,而服务端的心跳判定周期是 30 秒,导致 Nginx 先于服务端把空闲连接关掉了。客户端发现连接断开后立刻重连,重连风暴把负载打得很高。
这个问题的根源是网关超时和心跳周期不匹配。后来我把整套链路的心跳节奏统一了:应用层 30 秒发心跳,Nginx 超时 3600 秒,服务端 90 秒判定死亡。至此再没出现过莫名重连。
教训很简单:心跳不是一个"应用层参数",而是一条完整链路上的协作参数。客户端、网关、应用节点、Redis TTL 四者的超时时间必须按"防火墙 >= Nginx >= 应用判定死亡 > 心跳间隔"的层级关系设置好,任何一环不匹配都会出怪问题。
7.2 连接泄漏:只会读不会关的客户端怎么处理
还有一次线上问题让我印象很深:某个老版本的 App 客户端在弱网环境下不会正常发送关闭帧,也不会响应 Pong。TCP 连接半死不活,服务端检测不到异常,连接数一直涨,最终内存被打满。
应用层心跳只能发现"不响应 Pong"的连接,但如果你只是每 30 秒发一次 Ping,且客户端永远不会回 Pong,那你最早也要等 90 秒才能清理掉它。如果 QPS 很高,90 秒就能积累大量死连接。
后来我做了两个优化:
- 在心跳检测时,除了发 Ping,还会检查连接最近一次"收到任何帧"的时间。只要客户端还在 TCP 层传输数据(哪怕是垃圾帧),就把它标记为活跃;超过 90 秒没有任何数据到达的,直接强杀。
- 对客户端主动断开但 TCP 层没有 FIN 的情况,依赖内核的 keepalive 来做最后兜底,应用层只负责更快的检测。
另外,我要强调一点:清理死连接时一定要在 finally 块里执行 unregister。否则连接关闭异常会导致 session 残留在注册表里,造成"幽灵连接",这是连接泄漏最隐蔽的形态。
7.3 广播风暴与消息重复
最后一个坑来自广播场景。做直播弹幕时,一开始广播消息直接走ws:broadcast全局频道,所有节点收到后向本地所有连接推送。等房间人数上到几千,问题就爆发了:每个节点推送队列积压,Redis 订阅线程卡死,消息延迟从毫秒级飙升到秒级。
问题本质是广播范围没有做细粒度控制。全局广播频道会让每个节点都收到所有消息,但一个弹幕只属于一个房间,节点收到后还要去查这个房间里有哪些本地连接,大部分查询都是空转。
我把广播频道从全局切到了房间维度:ws:broadcast:room:{roomId},节点按需订阅用户所在房间的频道。在此基础上,发送线程也从 Redis 订阅线程里拆了出来:订阅线程收到消息后只负责放进本地队列,由独立的发送线程池消费并推送给客户端。这样即使某个房间的消息量特别大,影响的也只是这个房间的发送线程,不会拖垮整个节点。
消息重复的问题也值得一提。Redis Pub/Sub 本身不会重复投递,但引入 MQ 后,消费者如果在发送推送后、ACK 之前宕机,重启后会重新消费这条消息,导致用户收到重复推送。解决方式是给每条消息生成唯一 ID,节点在本地维护一个最近处理过的消息 ID 缓存,收到重复 ID 直接丢弃。缓存窗口不用太长,30 秒即可,因为 MQ 的重复消费通常发生在很短时间内。
做 WebSocket 集群这几年,我最大的体会是:不要一上来就追求最复杂的架构。单机能把连接管理、心跳、推送这些基本功做扎实,是比上来就上分布式更重要的能力。集群方案的选择也遵循这个原则——用户量上来了,先做 Sticky Session 顶一阵;遇到跨节点路由需求了,再加 Redis Pub/Sub;消息可靠性要求高了,再引入 MQ 和服务化改造。每一层方案都解决上一层的痛点,但也会引入新的复杂度,只有在真正需要的时候才值得接住这份复杂度。