ARTICLE DETAIL

资讯详情

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

WebSocket分布式集群实战:从单机连接到百万级在线消息推送

WebSocket分布式集群实战:从单机连接到百万级在线消息推送 前阵子折腾一个实时通知类项目WebSocket 单机服务从测试环境里几千连接突然被放到了线上几万甚至几十万连接面前。单机版本跑起来倒是轻快可真要顶着在线量做分布式集群才发现要解决的全都是单机时代根本不存在的痛点连接存在哪台机器上别人怎么知道、消息怎么从业务服务转发到正确的连接、连接断了怎么恢复。这篇文章就把这段时间踩过的坑、用过的方案、最终落地的一些判断依据完整分享出来项目里直接能抄。1. 单机方案看起来够用实际藏了一堆雷很多项目最开始做 WebSocket 都是这个套路Spring Boot 或者 Netty 起一个服务客户端连上来服务端往一个 ConcurrentHashMap 里塞 connection收到消息就遍历 map 找到目标连接直接 send。业务简单代码量少联调测试也完全没问题。我最初那版就是这么干的一个 JVM 里跑个 Nginx 做端口转发后端起两个 WebSocket 服务实例靠 Nginx 轮询把连接打散表面上是集群实际上两个实例各管各的连接互相看得见对方的连接不存在消息推给 A 实例的用户如果用户连在 B 实例上就永远收不到。单机方案真正的风险不在代码逻辑而是它把“连接”当成了一种无限资源在消耗。每条 WebSocket 连接在服务器上都不只是套接字文件描述符它背后还有 TCP 收发缓冲区、应用层的心跳定时器、加密握手上下文、业务侧的用户会话对象。我后来按实际内存占用估算过一条空闲连接光服务端对象和缓冲区加起来大概要吃掉 10KB 到 30KB 的内存这还不算框架自身的封装。一台 8G 内存的云主机跑默认 JVM 堆配置就算文件描述符放开到十万连接撑到三四万的时候 GC 就已经开始肉眼可见地不健康了。单机方案的另一个硬伤是发布和重启。只要服务需要发新版、改配置、调参数就必然会出现一段窗口期所有连接同时切断。客户端如果没有重连机制用户在那一瞬间感受到的就是“页面突然收不到消息了”。如果项目里还挂着定时任务在凌晨触发重启那一下正好赶上消息推送窗口丢消息基本是板上钉钉的事。所以等到在线量真正上来或者业务上开始要求发布不中断、单节点故障不影响全局的时候就不是“要不要上分布式”的问题而是“必须怎么上”的问题。只要做了这个决定后面所有的架构调整都会围绕一个核心目标来组织让“连接在哪台机器上”这个状态变成集群内所有节点都能感知并查询到的信息同时让消息的转发路径可以被设计、被监控、被兜底。2. 分布式集群要先看清瓶颈在哪里在动手拆集群之前先把单机的瓶颈拆清楚特别重要否则很容易把“分布式”做成本地 map 换成了 Redis map连接还在单机上只是把会话信息搬了个地方该挂还是挂。先列一下影响连接规模的主要因素文件描述符数量操作系统默认一般 1024改到 65535 是很基础的一步百万连接需要评估单机端口和 fd 的极限。内存每连接内存成本上面提过。Netty 的内存池设计会省一些但业务对象的开销省不掉。CPU心跳包处理、消息编解码、业务回调连接多了 CPU 会成为瓶颈尤其 JSON 序列化这种。单点故障无论单机配置多高宕机、断网、发布都意味着服务不可用或连接全断。框架自身的扩展边界Tomcat 原生 WebSocket 到一定连接数后线程模型开始吃力Netty 要好很多但也不是无上限。分布式集群要解决的不只是“加的机器更多”而是让连接状态在集群里有一个统一的视图。这句话是整篇文章的题眼。你要知道任何一个用户 ID 对应的连接在哪个节点上才能把消息准确地送过去你要知道集群里所有节点的健康状态才能决定新连接往哪里分你要知道某个节点挂了之后它身上的连接要怎么重新分配才能做故障转移。这个阶段最容易犯的错是把单机的MapString, Session直接替换成MapString, String存到 Redis然后所有节点都去 Redis 查“这个用户连在哪”。方向是对的但只做了会话共享消息路由、连接生命周期管理、节点摘除通知这些配套全没跟上最后 Redis 成了单点Redis 一抖动整个推送链路就跟着颤。所以在设计分布式方案时我建议从三个维度去拆会话注册与发现连接状态怎么维护、怎么查消息路由与转发消息怎么到达正确节点、到达正确连接生命周期管理连接断开、节点上下线、客户端重连怎么协同下面几节就按这个框架逐个展开。3. 会话注册与发现不要把所有连接都塞进 Redis先说结论连接本身是进程内的资源不可能也不应该序列化到 Redis。Redis 里存的只是“路由信息”也就是这个用户当前连接在哪台机器上。连接对应的 WebSocket Session 对象永远只存在于某个节点的本地内存里。这样一来分布式环境下的会话模型就变成了两层第一层本地会话表存在于每个节点内存。就是一个 Mapkey 是 userId 或者连接 IDvalue 是本地 WebSocket Session。第二层全局路由表存放在 Redis。key 也是 userId/连接 IDvalue 是节点标识比如ws-node-002同时设置 TTL 做自动过期。连接建立流程大致是这样客户端通过 Nginx 或被网关路由到某个 WebSocket 节点。节点建立连接后往本地会话表写入userId - Session。节点往 Redis 写入SET user:1001 ws-node-002 EX 90。节点启动一个定时任务每隔 30 秒为所有活跃连接续期也就是重新EXPIRE user:1001 90。连接关闭时删除本地会话并删除 Redis 中的对应 key。这样做的好处是整个集群的任意节点都能通过一次 Redis 查询知道目标用户连接在哪台机器上。而且 Redis 只存字符串键值对容量和性能都足够支撑在线用户量几倍甚至十几倍的会话记录。我后来在十万在线用户的压测环境里Redis 这个表的 QPS 峰值也就几千完全不是瓶颈。这里要补一个关键判断TTL 为什么要设成 90 秒而不是更长或更短这个值要同时满足两个条件一是比心跳上报间隔大至少 3 倍以上二是节点宕机后路由信息需要尽快自动失效。如果心跳是 30 秒一次TTL 设 90 秒就是留出了 3 次续期的容错空间。网络闪断、Redis 连接池暂时拿不到连接这类的偶发情况不至于让会话被误删。但如果 TTL 设太长比如 10 分钟一个节点宕机之后别的节点依然会往这个死节点转发消息客户端重连后又会遇到新连接写入和旧路由残留的冲突。选择 Redis 的时候要注意一个细节不要用默认的单实例。会话表是核心链路的一部分Redis 挂了对推送来说是致命的。至少用主从 哨兵有条件就上 Redis Cluster。这个表的写入是SET/EXPIRE/DEL这种简单操作Cluster 模式下按 key 的 hash slot 分散就好不需要担心事务问题。会话表还有一个容易被忽略的用途顶号处理。同一个用户 ID 在新设备上登录新连接建立后正常逻辑应该把旧连接踢下线。做法是新节点先写 Redis 会话表GET之前的节点标识如果存在且不是自己就向那个节点发一个内部命令让它把旧连接主动关闭。旧节点再通过本地会话表找到对应 Session 执行close同时把 Redis 里的 key 覆盖成新节点。这个顺序不能反先踢旧连接再覆盖路由否则可能一瞬间两个设备都认为自己是有效连接。4. 消息路由与推送定向推送和广播推送要分开设计会话表解决了“找到连接”的问题接下来要解决“把消息送过去”的问题。这两个问题不能混在一起讨论因为定向推送和广播推送在后端实现上走的完全是两条路。4.1 定向推送查路由表 内部调用一对一的消息比如聊天私信、单用户通知流程最简单业务服务收到消息。通过 userId 查 Redis 会话表拿到节点标识。如果节点标识不是自己走内部 HTTP/RPC 调用目标节点的一个接口比如POST /internal/push把 userId 和消息体带过去。目标节点收到请求后从本地会话表拿 Session执行sendMessage。这里有一个绕不开的组件内部调用通道。不管你是用 OpenFeign 调 HTTP 接口还是直接发 Netty 自定义协议甚至是走 Dubbo/gRPC都得有一套节点间通信的能力。我建议不要走公网地址而是把节点间的内网地址写到服务发现配置里用独立端口来做内部 API方便做鉴权和隔离。这个方案最大的风险点是如果调用链路上没有兜底Redis 里查到节点但请求发过去的时候连接刚好断了消息就会在这个环节丢。所以我在内部推送接口里通常加一个返回值设计目标节点收到推送请求后先查本地会话如果本地也没有这个连接返回一个明确的状态码调用方拿到这个状态码后就可以触发离线消息保存逻辑。这样不会把错误留到调用链的最末端才发现。4.2 广播推送用 Pub/Sub 还是 MQ广播场景比如直播弹幕、系统公告、全站广播如果每个节点都去 Redis 查一遍全量会话再逐个推送会有两个问题一是全量查询压力大二是消息到达不同节点的时间不一致尾延迟不可控。更常见的做法是让每个节点都订阅一个同一个消息通道广播消息发到这个通道上所有节点同时收到各自负责推送自己本地连接的这部分用户。这个“通道”选型我在 Redis Pub/Sub 和 MQ 之间纠结了挺久最终建议按场景分开维度Redis Pub/SubMQRabbitMQ / RocketMQ消息堆积不支持消费慢就丢支持可以缓冲消费状态确认无有 ACK可重试持久化无有适合场景弹幕、实时通知、可容忍少量丢失系统公告、订单通知、不允许丢消息运维成本低复用 Redis高需部署维护如果项目里有现成的 RocketMQ 或者 RabbitMQ优先用 MQ因为订阅通知、死信、消费重试这些能力都是直接自带。如果纯粹是轻量级实时推送不想引入额外中间件Redis Pub/Sub 也能跑前提是你能接受客户端重连后漏掉几条消息以及 Redis 本身没有持久化堆积能力这一点。我在线上项目里最终是两类通道都用了广播公告类走 RocketMQ弹幕这类高频低价值消息走 Redis Pub/Sub。这样既保证重要消息不丢又不让高吞吐场景把 MQ 的消费压力拉满。广播推送时还需要注意一点消息体尽量用二进制或紧凑 JSON避免在每个节点上做重复的序列化。最好是业务服务只发一次序列化之后的字节数组直接丢到通道里各节点消费时直接 send 给客户端。我见过一个项目广播消息里嵌套了两层 JSON 字符串客户端拿到之后自己又做了一次JSON.parse一个消息白白浪费了两次序列化开销高并发下 CPU 损耗非常明显。4.3 顺序与去重同一个用户的多个消息别乱序广播场景下消息顺序问题不敏感但定向推送场景下同一个用户的多条消息如果经过多个线程或多次内部调用很可能出现后发的先到。解决办法有两种一种是在业务层给消息分配单调递增的序列号客户端根据序列号做排序另一种是在推送链路里保证同一个 userId 的消息走同一个线程池或者用单连接顺序发送避免并发写同一 Session。Redis 会话表在这种场景下还能再发挥一个作用给每个用户记录一个lastMsgId每次推送前先自增拿 ID客户端带上这个 ID 做幂等处理。如果消息重复推送客户端可以根据 ID 去重不会出大问题。5. 两种落地结构从轻量网关到独立接入层会话模型和路由方案确定之后整个分布式 WebSocket 系统的骨架就能看到全貌了。根据团队规模和项目复杂度我实践下来有两种比较顺手的落地结构。5.1 轻量接入结构适合中小规模这种结构的核心是把 WebSocket 接入和业务逻辑放在同一个进程里通过 Nginx 做负载均衡和粘滞路由会话表存 Redis广播走 MQ。架构长这样Nginx 四层负载均衡监听 8080后端是多个 WebSocket 节点。客户端连接请求进入 NginxNginx 根据用户 ID 做一致性哈希保证同一个用户的所有连接都打到同一个节点。每个节点启动时向注册中心上报自己的节点标识和地址。节点内部的本地会话表 Redis 路由表同时维护。业务服务需要推消息时先查 Redis 路由表决定走本地发送还是内部调用。之所以 Nginx 要做一致性哈希而不是简单轮询关键是因为 WebSocket 连接是有状态的如果同一用户刷新页面后连接被轮到另一台机器旧连接还没断开新连接已经建立就会出现一个用户两个连接并存推送时很容易发到旧连接上造成消息“消失”的假象。Nginx 配置里可以用hash $remote_addr但更严谨的是用用户 ID 的哈希。让客户端在 URL 里带上 userIdNginx 从请求参数里取值做 hash。配置片段大概是stream { upstream ws_backend { hash $arg_uid consistent; server 10.0.0.1:9001; server 10.0.0.2:9002; } server { listen 8080; proxy_pass ws_backend; proxy_timeout 900s; proxy_connect_timeout 10s; } }如果客户端连接地址里没有 userId 这种参数就只能退而求其次用hash $remote_addr。这种方案的缺点是同一用户在不同网络环境下可能落到不同节点所以最好还是业务层把 userId 放到查询参数里。这种结构的优点是改动量最小。原来已经有单机 WebSocket 服务的团队把本地 map 换成“本地 map Redis 路由表”再给节点间加一条内部调用通道两天之内就能上线分布式版本。缺点是接入和业务耦合在一起节点发布时要同时处理连接断开和业务逻辑迁移后期如果连接规模继续膨胀每个节点的压力还是会比较平均地上涨。5.2 独立接入网关结构适合规模化如果项目已经明显是“连接数量大、业务消息少”的形态比如物联网设备接入、IM 长连接我建议把 WebSocket 接入拆成一个独立的网关层。这样做的核心逻辑是连接管理是一项与业务无关的基础能力它自己的扩缩容、生命周期、节点摘除逻辑应该和业务服务彻底解耦。网关层只做三件事接受客户端连接维护连接状态。心跳保活把会话路由信息写到 Redis。接收业务层通过内部 API 下发的推送指令从本地会话表找到连接并发送。业务服务不感知 WebSocket 连接在哪台机器上它只调用一个“推送服务”的接口由推送服务内部去查路由表并调用对应网关节点。网关层的节点可以独立扩容业务服务完全不受影响。这种结构在发布时也可以做到连接无损迁移网关滚动升级时新节点先加入集群旧节点把连接平滑迁移过去或者等客户端重连。独立网关层的实现语言我建议用 Netty 或者 Go 的 gorilla/websocket性能和并发模型都比 Spring Boot 自带的 Tomcat WebSocket 好很多。Spring Boot 里的 WebSocket 虽然方便但连接量过万之后上下文切换和内存开销就会比较明显Netty 在这方面的资源控制粒度细得多。5.3 不同技术栈方案选择的小结简单整理一下我个人的选型参考项目情况推荐路径Java 技术栈连接量不大团队人少Spring Boot Tomcat WebSocket加 Redis 路由表Java 技术栈连接量大需要细粒度控制Netty 实现 WebSocket 网关Node.js 技术栈socket.io 自带多节点适配用 Redis adapterGo 技术栈gorilla/websocket 单节点配合自研路由表物联网、设备接入为主直接考虑 EMQX 这类专业 MQTT Broker别用 WebSocket 硬扛6. 容器化部署、网关路由与稳定性细节架构方案只是第一步真正让分布式 WebSocket 在线上稳定运行靠的是各种细节。下面这几块是实际部署和压测过程中反复调整过的地方。6.1 容器网络与 Nginx 配置要注意的地方有了多节点之后Nginx 的配置就不再只是简单转发。除了前面说的一致性哈希还有几个参数必须根据场景调整。proxy_timeout这个值建议设置成大于客户端心跳间隔。如果客户端每 30 秒发一个 ping 帧proxy_timeout至少要大于 35 秒否则 Nginx 会先于应用层判定连接超时把连接踢掉。我见过一个项目客户端心跳是 60 秒一次Nginx 默认的proxy_timeout是 60s两边在极限值上拉锯客户端频繁掉线排查了三天才发现是这里的问题。Nginx 四层负载均衡需要用到 stream 模块默认编译不一定带。使用前先确认nginx -V输出里有没有--with-stream没有的话要重新编译或换 OpenResty。WebSocket 的握手升级是 HTTP 协议但后续的帧传输是 TCP 流四层转发最稳妥不要用普通的 HTTP 反向代理去处理升级请求。容器化部署时不要把 Nginx 和 WebSocket 服务放在同一个 Pod 里用 localhost 通信除非你有充分的理由。用 K8s Service 做负载均衡时externalTrafficPolicy建议设置成Local保留客户端来源 IP方便做哈希路由和日志分析。如果设置成Cluster来源 IP 变成节点 IP哈希策略就失效了而且跨节点转发会多一跳延迟和故障面都会变大。6.2 容器环境下的会话生命周期管理K8s 环境下做滚动发布时Pod 会被逐个替换。如果直接kubectl rollout restart旧 Pod 的 WebSocket 连接会瞬间全部切断客户端体验到的是“集体掉线”。正确做法是给 Pod 配置preStop钩子和优雅退出。顺序大概是下发滚动发布指令K8s 逐个终止旧 Pod。旧 Pod 收到 SIGTERM 前先执行preStop调用一个内部接口把自身状态标记为“排空”。“排空”状态下的节点不再接受新连接Redis 路由表里的会话 TTL 也会因心跳停止而迅速过期。等待存量连接自然断开或者给一个宽限期比如 30 秒之后强制关闭进程。这样新连接会打到健康节点旧连接在宽限期内会陆续被客户端重连到新节点整个发布过程不会出现“全部连接同时断开然后同时重连”的连接风暴。客户端侧也必须有对应的重连退避逻辑。假设一个群发通知触发 10000 个客户端同时重连如果每个客户端都在 1 秒内重试新节点瞬间就要扛住一万个握手请求。最基础的做法是随机退避重连间隔取0~3秒随机值。更平滑的做法是每次重连间隔指数增长最多加到 30 秒连接失败时重置。6.3 心跳、缓冲区与消息体约束应用层心跳是必须做的TCP keepalive 不能替代因为 TCP keepalive 默认探测周期太长而且只能在连接级别感知对端不可达不能感知业务层卡死。我做心跳保活的经验值是客户端每 30 秒发一个 ping 帧服务端收到后返回 pong服务端每 60 秒检查一次连接池把超过 90 秒没活动没收到 ping 也没发消息的连接主动关闭。这三个数值之间留出了足够的容错空间网络抖动不会误杀连接死连接也不会赖着不清理。WebSocket 数据帧的大小也要做限制。框架默认对消息体大小的限制各不相同有的很小有的不限制。我建议在网关层统一设置一个上限比如单条消息 64KB。超过这个值的消息要么压缩要么拆分成多条。这样做的好处是防止恶心客户端发超大数据把节点内存打爆。Netty 里对应的是WebSocketServerProtocolHandler的maxFramePayloadLength参数Spring Boot 里对应的是container.addCustomHandler里的相同逻辑。7. 常见问题与排查技巧实录分布式 WebSocket 的问题排查难点在于链路长涉及前端、网关、Nginx、Redis、MQ、业务服务。我把实际过程中遇到的高频问题整理成一张速查表排查时可以按图索骥。现象可能原因排查手段客户端频繁断开重连心跳超时、Nginx 超时过短、机器负载过高看 Nginx 日志、客户端重连时间、服务端活跃连接曲线单节点连接数不均衡哈希策略不生效、新节点上线没调整权重检查 Nginx hash 配置看连接数监控消息推送到用户但客户端没收到路由表指向旧节点、连接被顶号、消息丢失查 Redis 路由表查看目标节点本地会话是否存在广播消息延迟高MQ 消费慢、某个节点消费阻塞看 MQ 积压量检查是否有个别节点线程池打满Redis 抖动导致推送失败路由表查询超时给 Redis 配置连接池超时、重试会话表加本地缓存节点重启后客户端重连风暴没有平滑下线、没有退避配置优雅退出客户端加随机重连间隔顶号后旧连接收不到下线通知踢下线的内部调用没有确认机制新节点发送下线命令旧节点执行成功后返回结果失败重试这里单独展开两个我踩得最深的坑。第一个是本地缓存加不合理导致的会话表不一致。最初我为了减少 Redis 查询压力在本地做了一层路由缓存key 是 userIdvalue 是节点标识缓存时间设了 5 分钟。结果用户换设备登录后新节点写了新路由旧节点还按本地缓存认为用户在自己这里消息直接发到旧节点旧节点的本地会话已经没有这个用户消息就丢了。这个问题的教训是路由缓存只能用于“查询本地节点是否存在连接”的场景不能用于“这个用户的连接一定在哪”的强一致判断。最终我把缓存时间缩短到 30 秒并且每次推送前还是要查一次 Redis 确认才算稳定。第二个是广播消息时本地连接遍历效率太差。初版广播实现直接遍历本地会话 Map每条连接调用一次send结果 Map 里两万条连接一条广播消息就会占住线程两秒其他消息全部排队。后来改成批量预取连接列表再用连接组或者消息扇出的方式并行发送延迟从秒级降到百毫秒级。遍历时还要注意 session 为空、连接已关闭这些脏数据发送前做一次状态检查否则会出现ClosedChannelException填满日志的情况。还有一个小技巧值得分享在 WebSocket 网关进程里打印日志时一定要把 userId、连接 ID、节点标识放进去并且用异步日志或日志队列防止磁盘 IO 拖垮推送链路。我见过一个生产事故服务本身没有性能问题结果日志同步写盘太慢连接事件日志排队排了两百万条最后内存被日志撑爆进程直接崩溃。日志是好东西但别让它成为系统的阿喀琉斯之踵。8. 方案评价与演进思路最后从个人角度评价一下这一整套方案。用 Redis 做会话表、用 MQ 做广播、用 Nginx 做路由这套组合最大的优点是每个组件都是成熟稳定的没有自研黑科技出了问题好排查团队也容易上手。对于绝大多数中小型实时通信场景比如客服系统、消息中心、弹幕、协同工具这个架构的性价比是最高的。它也有天花板。连接量继续往上走达到百万级甚至更高时Redis 会话表本身的 QPS、存储容量和网络带宽都会成为瓶颈。这时候一般往两个方向演进一是引入专业的实时通信基础设施比如用 MQTT Broker 替代自研 WebSocket 网关设备接入和消息路由全部交给 Broker 体系处理业务只需要做消息的上层业务逻辑。二是做多级路由把路由表按用户维度分片比如按 userId 的 hash 再拆成多个路由集群每个路由集群只负责一部分用户的连接信息。查询压力被分散到不同的路由集群单个 Redis 的负载就降下来了。这两种方向都是架构工程问题不是普通业务团队在初期需要面对的。如果一个项目每天的活跃连接还没超过十万核心精力还是应该放在客户端体验和消息可靠性上不要为了分布式而分布式。分布式本质上是在为故障和流量买单你买的是稳定性和扩展性付出的代价是架构复杂度、运维成本、排查问题的时间。判断是否要做的标准很简单单机方案是否已经出现你无法接受的故障或者发布流程是否已经严重拖累了迭代速度。我在实际项目里最深的体会是WebSocket 分布式集群的难点从来不在 WebSocket 协议本身而在于如何把“连接状态”这个强有状态的东西用各种手段变成一个可以被分布式系统管理和调度的对象。会话表、消息通道、节点通信、生命周期钩子四件事环环相扣少了任何一个环节系统都可以运行但线上一定会出问题。按照这套思路把四件事都补齐剩下的就是抄配置和压测调优的活儿了。
返回列表