十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

WebSocket集群状态同步方案:从单机到分布式实战

WebSocket集群状态同步方案:从单机到分布式实战 做体育直播平台这几年我对 WebSocket 集群方案从单机到分布式的整个演进过程算是有比较深的体会。这类业务场景很典型用户量大、实时性要求高、峰值冲击凶猛一场焦点比赛开打前用户像潮水一样涌进来如果底层架构撑不住卡顿、掉线、消息丢失这些问题就会在瞬间集中爆发。这次把我在实际项目中经历和总结的“WebSocket集群状态同步方案”完整梳理一遍从单机版怎么起步到集群化之后遇到哪些坑再到分布式状态下怎么同步状态、怎么推送消息一步一步说清楚。这套方案的核心价值是解决一个现实问题当一台服务器扛不住海量长连接时多台服务器组成的集群如何保证“用户在任何一台节点上都能收到消息、且各节点之间的连接状态保持统一认知”。适合正在做直播、IM、实时互动类项目的后端同学参考尤其是处在“单机够用但担心扩展”、“已经多机部署但消息乱发”这个阶段的团队。1. 内容整体设计与思路拆解为什么单机方案先“能用”再去谈“好用”1.1 单机时代的业务背景与架构状态刚开始做体育直播的实时推送时体量没现在这么大。当时的架构非常朴素一台应用服务器内置 WebSocket 服务前端建立连接后服务端保存所有在线会话Session然后通过一个全局的房间管理器记录每个用户属于哪个房间。比赛进行中后端将比分、角球、红黄牌、判罚等事件写入消息队列再同步推送到对应房间的所有连接。单机方案的核心代码逻辑用一个 ConcurrentHashMap 维护连接和用户信息就够了大概长这样Component public class LiveSessionManager { // userId - WebSocketSession private final ConcurrentHashMapString, WebSocketSession userSessions new ConcurrentHashMap(); // roomId - SetuserId private final ConcurrentHashMapString, SetString roomUsers new ConcurrentHashMap(); public void register(String userId, String roomId, WebSocketSession session) { userSessions.put(userId, session); roomUsers.computeIfAbsent(roomId, k - ConcurrentHashMap.newKeySet()).add(userId); } public void unregister(String userId, String roomId) { userSessions.remove(userId); roomUsers.getOrDefault(roomId, ConcurrentHashMap.newKeySet()).remove(userId); } public void sendToRoom(String roomId, Object message) { SetString userIds roomUsers.getOrDefault(roomId, Collections.emptySet()); for (String userId : userIds) { WebSocketSession session userSessions.get(userId); if (session ! null session.isOpen()) { session.sendMessage(new TextMessage(JSON.toJSONString(message))); } } } }这套方案在每天在线人数几千、单个房间峰值几百人时跑得很稳。但问题在于它的扩展能力完全受限于单台服务器的最大连接数、内存大小和 CPU 性能。当一场热门赛事同时在线人数到了几万服务器 CPU 直接打满堆内存疯狂增长GC 频繁到令人绝望连接开始大量超时。1.2 “多开节点”带来的第一个问题消息发给了谁为了解决容量问题最直觉的思路是把服务多部署几个实例前面加一层负载均衡让用户的连接分散到不同机器上。这个方向没有错但紧接着就会撞上一个问题——用户连在节点 A消息却从节点 B 发出来了。比如直播间后台的“进球事件”分发模块如果它在节点 B 上运行而观看用户全部连接在节点 A那么按照单机版的逻辑节点 B 根本不知道这些连接存在消息发给谁没人。这就是典型的“连接与服务耦合在一起”导致的分布式难题。所以集群化之后我们必须先回答三个问题用户的连接登记在哪里如何让所有节点知道某个用户当前连在哪个节点上一条消息产生后怎样才能只推给持有目标连接的节点而不是让所有节点重复推送如果某个节点挂了挂在它上面的用户如何被重新分配状态如何恢复这就是状态同步要解决的核心问题。可以在单机版的思路上演进但绝不能继续把所有东西都放在进程内。2. 核心细节解析与实操要点路由层、注册层、同步层三管齐下2.1 负载均衡层的“Session 亲和性”不是银弹很多人一开始会想到用 Nginx 或者网关做 Session 粘滞Session Affinity也就是同一个用户的所有请求都转发到固定的一台服务器。这样做确实能解决部分重连问题——用户断开后再连还会回到同一台机器之前的连接信息还在。但这里有个容易被忽略的隐患Session 亲和性依赖 IP 或者 Cookie 做哈希用户切换了网络比如从 WiFi 切到 5G或者代理出口 IP 变化亲和性就失效请求会落到另一台节点而原来的节点根本没有这个用户的连接信息结果就是“连接断了却找不到新家”。另一个坑是如果某台节点宕机负载均衡会把原本调度到它的用户全部转发到其他节点而其他节点的内存里没有这些用户的 Session等于所有在途连接全部失效。用户感受到的现象就是“比赛最激烈的时候直播突然黑屏/转圈/断开”。所以 Session 亲和性只能作为辅助手段真正要解决的是“全局可见的连接状态”也就是把连接信息从单机内存搬到公共区域。2.2 全局路由表用 Redis 保存“用户当前在哪个节点”我的做法是在集群中引入一层 Redis作为全局连接路由表。每台 WebSocket 节点启动后生成一个唯一的节点标识NodeId例如ws-node-1、ws-node-2。当用户连接进来时节点在 Redis 里写入一条映射userId - NodeId并设置过期时间比如 60 秒同时节点在本地维护一份用户的 Session 对象用于真正收发消息。这样全局路由表回答了一个关键问题任何时候要知道某个用户当前连在哪台节点查 Redis 即可。这里要注意过期时间不能太短否则用户长时间不活跃比如静默观看、不发消息会导致映射被提前清理推消息时被认为不在线也不能太长否则用户断线后要等很久才能“摘除”带来幽灵连接。我实践中取过 30 秒、60 秒、90 秒综合考虑心跳间隔一般 30 秒一次和网络抖动最后定在 60 秒比较稳。有了路由表之后发消息的流程变成业务侧产生实时消息带上目标用户 ID 或者目标房间 ID。推送模块根据用户 ID 查 Redis找到 NodeId。通过 NodeId 找到对应节点的服务入口调用该节点的内部接口推送。节点拿到用户 ID 后从本地 Session 表找出连接进行发送。比较简单的实现是让推送模块通过 HTTP 调用目标节点暴露的内部接口public void pushToUser(String userId, LiveMessage message) { String nodeId redisTemplate.opsForValue().get(ws:route: userId); if (nodeId null) { // 用户不在线丢弃或者走离线消息 log.warn(user {} not online, userId); return; } // 内部 RPC/HTTP 调用目标节点推送 wsNodeClient.push(userId, message); }2.3 房间级状态同步不按人推按房间推体育直播场景有个特点绝大多数实时消息是广播给房间内所有人的比如进球、红牌、VAR 判罚。如果按用户逐个查路由表再逐个推送效率太低成千上万用户的房间会把这个流程压垮。更合理的做法是维护“房间维度”的节点分布信息每个房间的用户分布在哪些节点上这个信息单独记录在 Redis 里结构可以是roomId - MapNodeId, SetuserId。每当用户连接、断开时都增量更新这份结构。推送一条房间消息时先查这个房间分布在哪些节点然后把消息分发给这些节点的本地推送接口再由节点在自己的内存中找出该房间所属的本地连接进行广播。用伪代码表达就是public void sendToRoom(String roomId, LiveMessage message) { // 1. 查询房间分布在哪些节点 MapString, SetString nodeUsers redisTemplate.opsForHash().entries(ws:room: roomId); // 2. 按节点批量分发 for (Map.EntryString, SetString entry : nodeUsers.entrySet()) { String nodeId entry.getKey(); SetString userIds entry.getValue(); wsNodeClient.pushToUsers(nodeId, userIds, message); } }这个方案看起来顺理成章但实现时有个细节容易踩坑用户断线时Redis 里的路由信息和房间节点分布必须同步更新。如果只删了路由表忘了清理房间节点分布那么下一条房间消息还是会发到已经没有该用户的节点造成不必要的网络 IO。所以我把“注册”和“注销”都封装成原子操作同时操作 Redis 的多个 key并配合过期时间兜底。3. 实操过程与核心环节实现Redis Pub/Sub 状态同步的完整落地3.1 为什么最终选了 Redis Pub/Sub 做节点间状态同步前面说的路由表和节点分发解决的是“消息找到正确节点”的问题。但还有一个更深层的需求节点之间要互相通信做到状态变更的即时广播。比如某个用户进入房间其他节点不一定要马上知道但某个节点收到一条房间广播消息需要让其他节点也拿到并推送给它们本地的连接。这个场景下最直接的技术选型是 Redis Pub/Sub。理由有三个实现成本极低只要引入 Redis 客户端就能用 publish/subscribe 完成节点间的消息广播不需要额外部署一套 MQ。实时性高Pub/Sub 是推模式消息发出后订阅者几乎立即收到延时通常在毫秒级。机房内网络可靠Redis 作为长连接和所有 WebSocket 节点保持订阅关系避免了节点之间互相维护连接带来的复杂管理。相比用每对节点之间建立 RPC 连接的方式Redis Pub/Sub 天然支持“一对多广播”非常适合“一条消息推送全集群节点”的模型。3.2 Channel 规划和 JSON 消息协议设计我在项目中为 WebSocket 集群规划了三类 Channelws/node/events节点注册、心跳、下线通知。ws/room/{roomId}按房间粒度拆分的实时消息通道每个房间一个 Channel自己按房间 ID 动态订阅。ws/user/{userId}点对点消息比如单聊通知、系统通知。动态订阅房间 Channel 这个点很关键。假如一共有 1000 个直播间每个节点不可能全量订阅所有 Channel这样内存和连接开销都太大。更好的做法是“按需订阅”节点上的第一个用户进入某个房间时该节点去订阅这个房间的 Channel当最后一个用户离开房间时节点取消订阅。这样每个节点只订阅当前有用户观看的房间资源使用非常可控。消息体我定义成统一的 JSON{ type: ROOM_BROADCAST, roomId: live_10086, event: GOAL, data: { matchId: m20240101, teamId: home, playerName: 某球员, score: 1-0 }, timestamp: 1714552000000 }其中type字段用于区分是房间广播、点对点消息还是系统控制指令。event是最小业务事件类型解析后可以直接映射到前端的事件处理函数。3.3 订阅消息后的处理流程节点启动后订阅 Channelws/node/events同时根据自己当前承载的房间列表动态订阅ws/room/*。收到订阅消息后的处理逻辑如下RedisListener(channel ws/room/#{roomId}) public void onRoomMessage(String rawMessage) { LiveRoomMessage message JSON.parseObject(rawMessage, LiveRoomMessage.class); String roomId message.getRoomId(); // 只推给本节点上的房间内用户 SetString localUserIds localRoomManager.getLocalUsers(roomId); for (String userId : localUserIds) { WebSocketSession session localSessionManager.getSession(userId); sendIfOpen(session, message); } }这里有个非常重要的设计节点收到 Channel 消息后不会再把消息发回 Redis也不会再去查全局路由表而是直接查本地的“房间 - 用户 Session”映射。因为这个消息原本就是从某个源节点发出来的其他节点只需要在本地做过滤分发。这样既保证了一个房间内所有用户只收到一次消息也避免了消息在集群里反复转发形成风暴。3.4 连接注册、注销与心跳续约的完整流程为了让状态同步真正可用连接的生命周期管理必须做完整。我总结了一套标准流程用户连接建立 1. WebSocket 握手成功后解析 token 获取 userId 2. 本地注册 Session - localSessionManager.register(userId, session) 3. 本地房间成员加入 - localRoomManager.addUser(roomId, userId) 4. Redis 写入路由表 - SET ws:route:{userId} nodeId EX 60 5. Redis 更新房间节点分布 - HSET ws:room:{roomId} {nodeId} {userIds} 6. 如果本地是第一个进入该房间的连接动态订阅 ws/room/{roomId} 7. 返回握手成功消息给客户端 心跳续约 - 客户端每隔 30 秒发送 ping - 服务端收到后刷新 Redis 路由表的过期时间 - 同时刷新本地 Session 的最后活跃时间 用户断开连接 1. 从本地 Session 表移除 2. 从本地房间成员移除 3. 从 Redis 路由表删除或依赖过期时间自动删除 4. 更新 Redis 房间节点分布 5. 如果本地最后一个用户离开了该房间取消订阅 ws/room/{roomId}有一个容易被忽略的坑心跳续约不能只刷新 Redis还需要同时刷新“房间节点分布”的过期时间。我给房间节点分布设置了同样的过期时间如果用户长时间不活跃房间映射超时被清理下次推送时节点找不到房间分布信息消息就丢了。所以心跳回调里要统一做一次“路由续期 房间分布续期”。4. 常见问题与排查技巧实录从 1006 到消息风暴4.1 客户端莫名断开code 1006 的排查思路WebSocket 开发中onclose事件里遇到 code 1006 是最常见的现象之一。1006 是一个特殊的状态码表示连接被异常关闭通常不是正常握手后的 close 帧导致的而是底层连接直接断了。在集群场景下我踩过的坑主要有这几个没有周期性心跳服务端 Nginx 或者网关层的 idle timeout 会把长时间没有数据传输的连接断开。客户端必须周期性发送 ping比如每 30 秒一次。我在项目中用定时器发送 ping同时服务端实现PongMessage响应。负载均衡超时时间配置过小如果前面挂了 Nginxproxy_read_timeout默认 60 秒而你的心跳间隔是 30 秒看起来够用但碰上 GC 停顿或者网络抖动一次心跳晚到达就会被断开。我把心跳间隔降到 20 秒同时把 Nginx 的proxy_read_timeout设为 75 秒减少误杀。服务端节点重启节点发布重启时未完成优雅退出的话连接直接 RST客户端就会收到 1006。解决办法是监听 JVM 关闭钩子先向客户端广播一条“服务即将维护”的控制消息再关闭所有连接。客户端收到后延迟 1~2 秒自动重连。跨网络环境切换移动端从 WiFi 切到蜂窝网络TCP 连接在系统层面已经断了但应用层没有感知。客户端要监听网络状态变化主动触发reconnect。4.2 消息重复推送的陷阱Redis 重试机制带来的重复投递高并发下Redis Pub/Sub 一旦网络闪断客户端库重连后可能会重新订阅但如果业务代码在重连时没有做幂等处理节点会重复处理同一批房间消息用户就会看到重复弹幕、重复比分提示。我在实际项目中的处理方式分两层第一层消息带去重 ID。每条广播消息生成一个全局唯一的msgIdUUID 或雪花 ID节点内部的BloomFilter记录最近处理过的消息 ID重复消息直接丢弃。这个过滤放在业务分发之前能挡住很大一部分重复消息。第二层客户端幂等。前端拿到实时事件后根据event timestamp matchId生成事件哈希用 Set 保存最近 200 条重复事件不渲染。这样即使服务端偶发重复用户也感知不到。4.3 房间广播风暴如何避免大直播间拖垮全集群体育直播最极端的情况是一个超级热门房间同时在线 10 万人。如果不加控制一条进球消息会瞬间产生大量 Redis、网络和 CPU 开销。我的方案是做“节点本地聚合广播”。每个节点收到 Channel 消息后不再逐条发送 WebSocket 消息而是把同一房间的多个事件合并成一批在极短时间窗口内比如 50ms打包发送。用户在客户端感受到的效果是“几毫秒内连续几条消息一起到达”体验不降级但服务端把每秒钟几万次的发送调用降低到几百次。另外对超大房间做了分片策略如果房间在线人数超过 5000按userId % 10拆成 10 个子房间每个子房间独立维护成员和节点分布推送时并行发到 10 个 Channel。这样单 Channel 的订阅和推送压力就分开了。4.4 Redis 单点故障的兜底方案如果全局路由表和 Pub/Sub 都依赖的 Redis 挂了整个集群的状态同步就瘫痪了。所以在生产环境里我把 Redis 也部署成集群模式比如 Redis Cluster 或者主从 Sentinel保证高可用。这个环节其实踩过不少坑。有一次做故障演练kill 掉 Redis 主节点Sentinel 自动切换后WebSocket 节点上已建立的 Pub/Sub 订阅连接全部断开而且 Redis 客户端库没有自动重新订阅。导致的结果是节点还在但收不到任何房间消息用户看到的就是直播画面正常比分死活不更新。解决方法是监听 Redis 的连接状态事件一旦检测到连接断开并恢复重新执行一遍全量订阅。恢复订阅之前先做一次全量“掉线用户重挂”遍历本地 Session重新把路由关系写入 Redis确保切换期间漏掉的用户信息被补回来。这个“重挂”操作我用了一个简单的定时任务每分钟检查一次 Redis 里路由表数量和本地 Session 数量的偏差如果偏差超过 5%主动全量重挂一次。4.5 集群扩容与缩容时的平滑操作日常运维还有一个场景凌晨低峰期要缩容一台节点或者活动前要扩容两台节点。这种操作不能直接停机否则在线用户会瞬间断开。我用的策略是缩容先把节点标记为“draining”状态负载均衡不再给它分配新连接然后向该节点上的客户端推送一个“请切换节点”的控制消息客户端收到后主动断开并重连这时新一轮连接会均匀分布到其他节点。等该节点在线连接数降到 0再安全下线。扩容新节点加入后它的本地状态是空的需要从 Redis 拉取所有房间的当前状态比如房间当前比分、事件序列号之后才能开始接收订阅消息。所以我在节点启动流程里加了一步“初始化同步”从 Redis 读取每个房间的最新状态快照加载到本地缓存。这样新节点加入后用户连接过来才能立刻展示正确的比赛信息。5. 方案的技术选型对比与最终架构总结5.1 三种常见状态同步方案的对比做 WebSocket 集群业界常见的状态同步方案不止 Redis Pub/Sub 一种我梳理一下特点对比方案实现成本实时性持久化适用场景Redis Pub/Sub低毫秒级无直播间实时广播、节点间命令同步Kafka/RocketMQ高秒级到毫秒级有需要回溯、削峰、审计的异步事件自研 RPC 网格高毫秒级无超大规模集群节点间需要精细控制实际项目中我把这三种方案做了分层Redis Pub/Sub 负责实时状态同步和节点命令Kafka 负责削峰、异步处理和历史事件落库自研 RPC 网格只有在大规模集群里用于控制面消息平时不使用。5.2 最终架构的组件全景整理一下完整的 WebSocket 集群状态同步方案涉及的组件WebSocket 服务节点N 个管理客户端长连接、本地 Session、房间成员。Redis Cluster存放全局路由表、房间节点分布、Pub/Sub 消息通道。负载均衡层Nginx/云负载均衡做四层 TCP 转发TLS 卸载放在七层。业务消息入口比赛事件产生后经过 Kafka 异步写入消息队列再由分发服务发布到 Redis Channel。监控组件Prometheus 采集每个节点的连接数、消息吞吐量告警规则设置连接数突降和消息积压。最终发一条“进球事件”的链路是比赛数据源 - 业务后端 - Kafka - 事件分发服务 - Redis Pub/Sub(ws/room/{roomId}) - 各 WebSocket 节点 - 本地 Session - 客户端这条链路里每个环节的容量都比较容易扩展WebSocket 节点可以水平扩Redis 可以集群化Kafka 可以扩大分区数瓶颈点已经转移到了数据库和带宽上而不再是 WebSocket 服务本身。5.3 如果从零开始落地推荐的实施步骤如果你正在做一个新的体育直播或者实时互动项目我建议按照这个顺序逐步落地第一步确定 WebSocket 服务框架。Java 用 Netty 或者 Spring WebSocketGo 用 gorilla/websocket。第二步封装本地 Session 管理、房间成员管理保证单机模式可用。第三步引入负载均衡解决多节点部署后的连接路由用 Redis 做全局路由表。第四步引入 Redis Pub/Sub 做房间消息广播完成节点间状态同步。第五步加心跳、重连、幂等这些稳定性能力处理 1006 和重复消息问题。第六步做监控和告警确保节点扩容缩容操作顺畅。每一步验证通过后再进入下一步不要一上来就追求最复杂的架构。6. 经验沉淀这些坑踩过之后我对状态同步的理解从单机到分布式这套 WebSocket 集群方案最核心的收获是状态同步的本质不是把状态复制到所有节点而是让每个节点都清楚“我应该对哪些用户负责以及其他人需要我处理什么”。只要节点间的职责边界清晰状态同步的复杂度就能降下来。具体到实际的心得有几个点值得展开说说。关于 NodeId 的设计我遇到过一个大坑最开始用进程 IP 作为 NodeId后来发现同一台物理机部署多个实例时不同实例的 NodeId 会冲突导致路由表互相覆盖。后来改成IP:端口的组合并加了一个随机后缀保证了全局唯一性。另外NodeId 会出现在 Redis 的 key 和消息体中所以尽量短但也要可读方便排查问题。我用的是ws-i-192-168-1-10-8080-7f3a这种格式虽然长但一眼能看出机房、IP、端口和进程标识。再一个是消息序列的问题。体育直播里事件顺序非常重要不能出现“比分已经 2-0 了客户端才展示 1-0”。Redis Pub/Sub 在正常情况下保持发布者的发送顺序但如果发生了网络闪断重连或者使用了多生产者并发发布顺序是无法保证的。我在事件分发服务里引入了“每房间消息序号”的概念每个房间维护一个自增的序号客户端根据序号判断是否出现了跳跃。如果发现序号跳跃超过 1说明中间有消息丢失客户端触发一次“增量状态恢复”从 REST 接口拉取该房间的最新状态把漏掉的事件补回来。这个策略在直播场景下非常有用保证比分、事件展示的最终一致性。关于分布式锁在其中的应用也有一个值得分享的场景用户在多个设备上同时打开直播比如手机和平板如果同一个 userId 的后端连接超过一个踢掉旧连接时要加锁防止两端同时注册导致路由表互相覆盖。我用了 Redis 的SET lock_key userId NX EX 5做分布式锁保证同一时刻只有一个连接在处理注册逻辑。这个锁的过期时间设置成 5 秒但注册逻辑本身非常快通常是毫秒级所以不会出现误删锁的问题。但如果逻辑比较复杂建议用 Redisson 的看门狗机制自动续期避免锁过期。最后关于监控给一个实用建议WebSocket 集群最需要盯的指标不是 CPU 和内存而是“在线连接数”和“消息推送 QPS”的曲线。这两条曲线的异常往往意味着架构问题。比如在线连接数突然下跌 10%大概率是某个节点出现 OOM 或者网络分区消息推送 QPS 持续上涨但在线连接数不变可能是出现了重复推送风暴。我把这些指标接入了 Grafana每次比赛结束后复盘曲线能提前发现很多隐患。兜底方案里还有一个细节本地 Session 与 Redis 路由表的一致性校验。我写了一个后台任务每 30 秒遍历本地的连接检查每个用户的 Redis 路由表是否仍指向当前节点如果不是说明出现了异常比如用户被其他设备的连接顶下线或者重连后路由未更新此时本地主动关闭这条连接让客户端重连并重新路由。这个任务看起来简单但在长时间运行的集群里能清理掉大量“僵尸连接”和提前过期产生的错误路由显著提升推送准确率。回到最开始说的那个流动场景——用户点进直播间连接建立节点分配路由注册消息经过 Kafka 到 Redis再被正确的节点接收并推送到用户的手机整个过程环环相扣。架构演进到这一步就不再是“一台服务器硬扛所有压力”而是整个集群各自分工、协调配合。这个过程里踩过的每一个坑最后都变成了这套方案里的一块拼图补上去之后系统才真正稳下来。
返回列表