ARTICLE DETAIL

资讯详情

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

WebSocket集群方案详解:从Sticky Session到MQ推送架构的演进与实践

WebSocket集群方案详解:从Sticky Session到MQ推送架构的演进与实践 前阵子有个朋友的公司做在线客服系统单机跑了一年多相安无事。结果有次渠道推广爆量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 模型百万并发连接在理论上是可以做到的但实际业务环境远达不到。我通常这样估算每个空闲连接在应用层大约占用 20KB50KB 内存包括 session 对象、读缓冲区、写缓冲区。10 万在线连接大约需要 2GB5GB 内存这是纯连接的消耗还不算业务对象。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 里用 ConcurrentHashMapGo 里用 sync.Map本质都一样Component public class SessionRegistry { // userId - WebSocketSession private final ConcurrentHashMapString, 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 的映射也就是 MapString, 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 秒心跳包会占用大量带宽和 CPU10 万连接每秒光心跳就是几万条消息间隔太长比如 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 - SetnodeId # 存活节点列表 ws:msg:{nodeId} - 消息 # 点对点消息频道发往指定节点 ws:broadcast - 消息 # 广播消息频道所有节点订阅每个节点订阅两个频道ws:msg:{自己的nodeId}和ws:broadcast。这样点对点和广播就用一套机制统一处理了。4.2 点对点消息的完整链路点对点消息的完整链路是这样的业务服务想给用户 U 推送一条消息。查询 Redisws: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 ConcurrentHashMapString, SetString 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 什么时候必须上 MQRedis 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值为当前节点 nodeIdTTL 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 秒封顶每次重连时间加一个 030% 的随机抖动避免所有客户端同时重连压垮网关。6.4 关键代码骨架最后给一个完整的 WebSocket 连接注册 跨节点路由的最小骨架基于 Spring Boot 和 RedisServerEndpoint(/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 和服务化改造。每一层方案都解决上一层的痛点但也会引入新的复杂度只有在真正需要的时候才值得接住这份复杂度。
返回列表