拓冰建站拓冰建站
首页 / 资讯中心 / 正文

WebSocket分布式集群化改造:从单机到十万连接的体育直播实践

对做体育直播的技术团队来说WebSocket几乎就是命门。比分推送、弹幕互动、房间状态、连麦信令……用户看到的每一帧画面背后都有一群WebSocket连接在疯狂收发消息。我带的平台早期就是一台裸机扛着所有连接后来用户量涨上来被迫一步步把单机架构拆成分布式集群期间踩的坑、趟出来的路今天完整整理出来。这篇文章适合已经有WebSocket基础、正准备做集群化改造或者已经在分布式路上挣扎的朋友。我会从单机架构的瓶颈讲起讲清楚为什么必须做分布式、连接路由怎么设计、状态同步怎么保证不丢不乱最后再附上实战中遇到的高频问题和排查思路。全程用体育直播平台的具体场景来讲但思路本身可以平移到在线教育、互动游戏、协同办公这类强实时交互的业务里。1. 单机时代一台服务器扛下所有的日子1.1 体育直播场景下WebSocket到底在传什么先回忆一下我们的业务场景。用户进直播间看比赛页面上除了视频流之外还有几样东西是实时变化的实时比分和比赛事件比如进球、红牌、暂停、节间休息弹幕和互动消息这部分消息量最大也最考验推送能力房间内的状态信息比如主播在线状态、礼物特效、连麦状态、竞猜入口开关业务主动下发的指令比如抽奖开始、违规提醒、服务端踢人。视频流走的是CDN加HTTP-FLV或者HLS这部分不需要WebSocket管。但后面四类数据如果靠客户端轮询延迟高不说还浪费大量HTTP请求。所以基本都会选择WebSocket长连接做实时推送。体育直播对延迟又特别敏感一个进球比分晚推两秒弹幕区早就炸了用户直接以为是平台卡顿。单机架构下服务端逻辑非常简单一台机器上跑一个WebSocket服务所有连接都在本机内存里维护。连接进来就是一个长连接对象存到内存Map里连接断开就从Map里删掉。一个直播间对应一个Room对象Room里面存着所有订阅这个房间的连接。消息广播就是在内存里遍历这个Map逐个写消息。这个阶段选型也不用纠结Java用NettyGo用Gorilla WebSocket或者gin框架自带的websocket封装写起来都很顺手。我们用Go写过一版核心代码大概长这样type Room struct { ID string conns map[*websocket.Conn]bool mu sync.RWMutex } func (r *Room) Broadcast(message []byte) { r.mu.RLock() defer r.mu.RUnlock() for conn : range r.conns { conn.WriteMessage(websocket.TextMessage, message) } }这段代码一度支撑了我们一个完整的直播季所有弹幕、比分、礼物通知都走它。但问题也随着用户量增长一点点暴露出来。1.2 单机架构的四个致命问题第一个问题是连接数的物理天花板。一台8核16G的云主机能稳定扛住的WebSocket长连接大概在5万到10万之间具体取决于消息频率和业务复杂度。连接数一旦超过这个量文件描述符先告警紧接着CPU和内存双双爆掉。系统对外的表现就是新用户连不上、老用户消息延迟、心跳超时被被动踢下线体验雪崩。第二个问题是热点房间问题。体育直播的流量分布极其不均匀热门比赛房可能占全站80%的流量。一个房间的广播是逐个连接写CPU消耗和连接数成正比。上万人的热门房间单靠一台机器广播推送延迟会从毫秒级劣化到秒级弹幕刷屏时尤其明显。第三个问题是发布部署问题。每次改业务逻辑要重启服务一重启所有连接全部断开。为了减少断连影响只能选凌晨这种低峰期发布。但体育赛事经常是晚上黄金档凌晨发布时比赛早结束了用户全跑光这个问题根本绕不开。第四个问题也是最终逼我们改造的问题就是单点故障。机器宕机上面所有连接全部断开用户集体掉线弹幕停摆、比分不更新。对体育直播来说比赛打到一半全体掉线等于事故级别的事件。提示在线人数不等于连接数。一个用户可能同时开网页端、App端、小程序端连接数通常是在线人数的2到3倍。做容量规划时这个系数一定要算进去。2. 分布式演进的第一个分水岭连接与业务解耦2.1 网关层的无状态化改造思路单机架构的问题本质上是连接和业务逻辑耦合在同一台机器上。要拆成集群第一步就是把连接层做成独立的网关服务。网关不关心弹幕怎么过滤、竞猜怎么结算只负责维护连接、处理心跳、收发消息。业务服务通过内部接口或消息队列跟网关通信。这一步的关键是无状态化。意思是任何一台网关机器都能处理任意一个连接不依赖本地存储的全量连接数据。真正做到这一点前提是把连接的路由信息从本地内存挪到外部存储。之后新加一台机器负载均衡把新连接分发过去网关之间不需要互相感知对方的连接情况整体架构横向扩展的能力就出来了。我们第一版改造的前两周非常痛苦因为很多代码默认连接在本地随手就是localConnMap[userId]这种调用。改造过程其实就是在跟这些本地假设做斗争把所有跟连接状态有关的东西都抽出来统一走外部存储。2.2 Redis接管连接路由信息每一条WebSocket连接进来我们生成一个全局唯一的ConnectionId然后记录三个关键映射用户ID到连接ID的映射连接ID到网关节点ID的映射房间ID到连接ID集合的映射这三组映射关系存到Redis里用Hash和Set结构维护。为什么选Redis因为连接建立、断开的频率非常高一场比赛下来一个用户可能断断续续重连几十次数据库扛不住这种高频读写而Redis的Hash和Set操作是原子的延迟又低刚好匹配这个场景。实际设计如下# 用户到连接 HSET ws:user:10001 connId conn-9f2a # 连接到节点 SET ws:conn:conn-9f2a node gateway-2 # 节点上连接的具体地址 HSET ws:node:gateway-2 conn-9f2a 192.168.1.10:8080 # 房间到连接集合 SADD ws:room:10086 conn-9f2a这里有个细节是连接ID必须全局唯一不能用简单的自增。分布式环境下我们直接用UUID前缀带上节点ID比如gw2-9f2a41这样排查问题时一眼就能看出这条连接挂在哪个节点上。Redis的Key也要设计好TTL连接正常关闭时主动删除异常断开时靠TTL兜底清理。2.3 消息投递路径的重新设计单机时代一条弹幕从用户A发出服务端直接在同一台机器上路由给房间内所有连接。集群化之后发送者和接收者可能不在同一个网关节点上消息投递路径必须重设计。核心思路是两层路由第一层是业务服务确认消息属于哪个房间然后发到消息总线。第二层是每个网关节点从消息总线订阅所有消息收到后查Redis确认本节点有没有这个房间的活跃连接。有就把消息写出去没有直接丢弃。这样每个节点只写自己本地的连接避免了跨节点的内存级连接调用。用一句话概括就是大家共享消息总线每个节点只处理自己的一亩三分地。响应时间肯定比单机时代本地直接Map查找慢一点但换来的是无限扩容的可能。这个取舍值得。3. 状态同步的核心从Redis Pub/Sub到Kafka3.1 消息总线选型为什么最终选了Kafka这里重点聊一下消息总线的选型过程。我们最早用Redis Pub/Sub因为配置简单Redis本来就有现成的发布订阅功能。但用了一个月就发现不对劲。Redis Pub/Sub有两大硬伤。第一消息不持久化消费者在消息发布的那一刻如果掉线或者正在GC消息就永远丢了。第二Pub/Sub没有消费者组的概念所有订阅者都会收到全量消息这会导致消息重复处理。后来我们把消息总线迁到Kafka。Kafka的几个特性在这个场景下非常合适消息持久化到磁盘消费者重启后可以从offset重新拉取消费者组机制天然支持多网关节点水平扩展一个分区只被组内一个消费者消费分区机制可以按房间ID取模保证同一个房间的消息有序到达。消息总线上流动的消息我们定义了统一的内部协议大致包含五个字段消息IDmsgId、房间IDroomId、事件类型eventType、负载payload、时间戳timestamp。每个网关节点启动一个消费者组订阅房间消息这个Topic按房间ID的分区去消费。3.2 状态一致性与分布式锁的正确姿势体育直播的状态同步里比普通广播更麻烦的是状态类数据。比如直播间的开播状态、比分数据、竞猜截止状态。这类数据如果直接靠消息流去改并发情况下容易互相覆盖。举个例子第4节刚开始时比分是88比90两个后台服务同时收到两个罚球命中事件一个把主队改成90一个把客队改成92如果两个写操作并发执行最终可能只更新了一个。这个时候会用到分布式锁。我们用的分布式锁就是基于Redis的SETNX加过期时间。但这里要重点提醒锁的粒度一定要控制好。不要做成全房间一把锁那样的话并发的弹幕消息全被锁卡住了性能直接崩塌。我们的原则是普通消息广播绝不加锁只有低频状态写入才用锁。分布式锁有个臭名昭著的问题就是锁过期但业务没执行完导致两个客户端同时拿到锁。我们解决方式是给锁加上续期机制用一个后台任务在锁快过期时自动续期业务执行完立刻释放。Redis官方给的Redisson库里有现成的看门狗功能但如果你不用Redisson就得自己写一个定时续期这块不能省。3.3 水平扩容与故障转移的落地细节集群做好之后扩容变得很简单新起一台网关节点接入同一个Kafka集群用同一个消费者组ID把节点信息注册进服务发现中心然后负载均衡层把新连接导流过去。我们用的是Nacos也接触过用Consul和Etcd的团队选型主要看团队已有基础设施没必要为了新项目专门引一套。Kafka本身也可以承担一部分服务发现的职责但不如独立的注册中心直观。故障转移我们分两个层面处理。节点级故障某一台网关宕机服务发现摘掉这个节点新连接不再分配过来。老连接会超时断开客户端触发重连后会被负载均衡设备分配到其他健康节点。这个流程依赖客户端重连机制所以客户端SDK的重连逻辑必须写好否则节点故障后用户就一直连不上。连接级故障单条连接断开客户端走断线重连逻辑重连成功后网关根据UserId从Redis恢复用户状态把新连接重新注册到对应房间然后触发状态补偿。这条流程涉及的细节比较多接下来用一节完整实操来演示。4. 完整实操一条比分消息的集群之旅4.1 从后台服务到网关节点的推送链路我们用一个具体场景来讲后台计分服务收到一场篮球比赛的比分变更通知要推送给正在观看这场比赛的所有用户。第一步后台服务把比分消息发到Kafka的room-message这个Topickey就用roomId这样同一个房间的所有消息都会进同一个分区顺序有保证。String eventJson JsonUtils.toJson(scoreEvent); kafkaTemplate.send(room-message, scoreEvent.getRoomId(), eventJson);第二步每个网关节点都有一个Kafka消费者消费线程拿到消息后先去Redis查这个房间在本节点有哪些连接然后只对本地连接进行写入。func consumeRoomMessage(msg *kafka.Message) { roomID : extractRoomID(msg.Key) connIDs, _ : redisClient.SMembers(ws:room: roomID).Result() for _, connID : range connIDs { // 只处理本节点的连接 nodeID, _ : redisClient.Get(ws:conn: connID).Result() if nodeID ! localNodeID { continue } conn, ok : localConnMap.Load(connID) if ok { conn.(*websocket.Conn).WriteMessage(websocket.TextMessage, msg.Value) } } }第三步客户端收到消息后更新比分面板。这条链路从消息产生到客户端渲染我们实测的延迟在100毫秒左右体感上是即时的。这个方案每个人都会写但有几个坑必须注意。查Redis的SMembers在高并发下可能会成为瓶颈所以我们加了本地缓存每个节点每隔几秒同步一次房间成员列表。如果房间成员变化不频繁这个缓存策略非常有效。4.2 房间状态标记的原子化更新直播间有一些状态是持久且权威的比如当前比分、进行到第几节、竞猜截止没截止。这类数据不能只靠广播一个新用户中途进来他需要马上看到当前状态而不是等下一次事件推送。所以状态的权威版本要单独存广播只是通知在线用户状态变了。我们的做法是每次状态变更后台服务先更新Redis里的权威状态再发广播消息。新用户建立连接后网关先查Redis里的状态快照回填给新连接然后接入实时消息流。这个过程我们内部叫状态回填。代码示例如下public void updateScore(ScoreUpdateEvent event) { String lockKey lock:room:score: event.getRoomId(); Boolean locked redis.setIfAbsent(lockKey, 1, Duration.ofSeconds(3)); if (!Boolean.TRUE.equals(locked)) { // 拿不到锁说明同一时刻已有别的更新在跑本次操作丢弃或者重试 return; } try { String redisKey state:room:score: event.getRoomId(); redis.set(redisKey, event.toJson()); kafkaTemplate.send(room-message, event.getRoomId(), buildEvent(event)); } finally { redis.delete(lockKey); } }这段代码的逻辑是先加锁再更新状态最后广播。锁的过期时间设3秒正常操作都在几十毫秒内完成续期任务会保证长操作不丢锁。状态更新和消息广播放在同一个事务里做不到严格一致但因为状态权威数据在Redis里广播消息只是通知偶尔丢一条广播用户下次刷新页面看到的状态还是正确的。4.3 断线重连与状态补偿机制断线重连是WebSocket方案里最影响体验的环节。体育直播场景下用户可能在地铁、电梯、地下车库网络随时断开客户端必须能自动恢复否则用户看一眼比赛回来发现比分不走了第一反应就是平台垃圾。我们的重连策略是这样设计的客户端检测到连接关闭立即尝试重连重连间隔逐步退避按1秒、2秒、4秒、8秒递增最大不超过30秒重连成功后客户端带上自己收到的最后一条msgId网关收到带上一条msgId的连接去Redis查该房间最近N条消息记录找到比这个msgId更新的消息按顺序补推给客户端。这里的最近N条消息我们直接用Redis的List结构维护每个房间一个Key消息推出去的同时塞进List尾部超过100条就从头部弹出旧消息。List用Trim来控制长度。# 发布新消息时同时记录到最近消息列表 LPUSH ws:recent:10086 {msgId:12345, ...} LTRIM ws:recent:10086 0 99用户重连时把断线前最后一条msgId发上来服务端从List里读取消息过滤出比这个msgId更新的按顺序补推。100条是我们的经验值补太多重连瞬间的网络开销大反而拖慢正常消息的接收补太少用户掉线时间一长就丢消息。提示心跳和重连要配合好。如果心跳正常但长时间没有业务消息连接也要保持住。有些客户端SDK会默认在空闲一段时间后主动断开这个行为一定要关掉。5. 高频问题排查与避坑实录5.1 onclose code 1006连接为什么异常断开1006是WebSocket里比较特殊的错误码它表示连接被异常关闭但客户端没有收到正常的关闭帧。我们一度隔三差五接到反馈说用户弹幕不刷新一查日志全是1006。排查下来最典型的几种原因网关机器内存不足进程被操作系统OOM Kill。这个看系统日志就能确认属于容量规划问题服务端发心跳的周期和客户端超时时间不匹配。比如服务端60秒才发一次心跳客户端30秒没收到数据就判定超时主动断开两边就对不上了网关前面挂的Nginx或者负载均衡设备有默认的空闲超时时间。很多LB默认60秒空闲就断连WebSocket如果长时间没有业务消息就会被LB从中间掐断表现就是1006。我们最终的解决方案是Nginx配置proxy_read_timeout 3600s同时服务端每30秒发一次ping客户端回pong。这样既绕过了LB的空闲超时也能及时感知死连接。切记心跳是双向的光服务端发不算完客户端不回pong的连接要及时清理不然会堆积大量僵尸连接。5.2 消息重复与乱序怎么处理Kafka的投递语义是至少一次at least once消费者在处理完消息之后才提交offset如果处理逻辑报错导致offset没提交上去重启之后会重新消费一批消息就产生了重复投递。我们的应对是全局消息ID加客户端去重。每条消息进入Kafka时都有一个唯一的msgId客户端在WebSocket层维护最近处理过的msgId集合重复的直接忽略。这个去重逻辑虽然土但非常有效也不依赖复杂的分布式事务框架。乱序问题主要靠Kafka分区来规避。同一个roomId的消息只要确保发消息时key传的是roomId就一定会进同一个分区分区内的消息顺序是有保证的。如果key传了别的字段比如userId那同一个房间的消息会分散到多个分区消费端拿到的顺序就是乱的这种错误一旦上线很难排查。5.3 重连风暴最容易忽视的集群事故重连风暴是个平时想不到、一出事就致命的坑。假设某台网关节点宕机这台机器上挂了3万个连接3万个客户端几乎同时检测到断开同时发起重连。负载均衡设备瞬间压力巨大其他健康节点也会被一股脑的连接请求打满严重时会导致整个集群雪崩式不可用。我们加了两个防护第一客户端重连延迟加入随机抖动。重连间隔不再是固定的1秒、2秒而是在基础间隔上叠加一个0到5秒的随机值让重连请求均匀分布。第二服务端对单个IP的并发连接数做限制。正常用户最多同时开两三个连接如果同一个IP突然冒出来几千个连接大概率是异常情况直接拒绝多余连接避免一台出口网关把整个集群打崩。这两个防护实测下来很有用。自从加上抖动后我们的网关节点重启对用户的影响几乎降到了零。5.4 心跳参数调优经验参考心跳参数不能从网上随便抄一份必须根据实际网络环境来调。我们的参数组合是服务端每30秒发一次ping客户端在60秒内没收到任何消息就主动断开服务端连续3次ping无pong就关闭连接。这个组合能保证在大多数移动网络下3分钟内检测出无效连接又不至于在弱网环境下频繁误杀。如果网络环境特别差比如很多用户在电梯和地下车库可以适当放宽超时时间。但放宽的代价是僵尸连接存活时间变长占着资源不干活需要根据实际场景权衡。我把常见问题整理成了一个排查速查表方便遇到问题时快速定位现象可能原因排查方向连接频繁1006LB空闲超时、心跳不匹配、进程OOM查LB配置、心跳日志、dmesg消息重复Kafka重复消费、客户端未去重查offset提交、接收端msgId去重逻辑消息乱序Kafka key传错、分区数变化查发送key、分区策略重连后状态缺失状态回填逻辑不完整查Redis权威状态、最近消息List特定节点连接数暴增负载均衡策略问题、重连风暴查LB算法、客户端抖动策略直播间广播延迟高房间连接数过大、Redis查询瓶颈查本地缓存、消息批量写优化5.5 本地开发调试的几点提醒多节点环境调试比单机麻烦很多我们踩过的坑包括本地电脑连测试环境的Kafka集群连不上后来发现是安全认证配置的问题IDEA里多进程服务启动后端口冲突导致节点注册失败最后是在启动参数里通过环境变量区分节点编号解决的。本地调试时我们可以这样模拟多节点# 在IDE里分别配置启动参数模拟两个不同节点 SERVER_NODEgateway-1 go run main.go SERVER_NODEgateway-2 go run main.go另一件容易被忽略的是WebSocket连接被安全网关拦截的问题。我们有个阶段H5页面在本地连WebSocket正常打包成App就连接失败排查后是移动端安全网关对Upgrade请求做了拦截需要在安全策略里放行WebSocket的握手路径。WebSocket连接走的是HTTP Upgrade机制任何中间层如果对HTTP协议做深度检查都可能误伤。结尾这套架构从单机改造到分布式集群前后折腾了两个多月最大体会是不要为了分布式而分布式。如果你的平台在线人数还不到一台机器能扛的量单机加优化是完全够用的。但当连接数要跨过十万这个量级集群化改造就必须提前规划尤其是状态同步这部分越早设计好后面扩容和运维越省心。我个人在改造过程中最推荐的做法是小范围灰度挑几个热门直播间跑新集群其他房间继续走老的单机逻辑观察一两周再全量切换。任何架构升级稳定压倒一切。还有一个小技巧消息总线的消费端一定要加监控告警消费Lag一旦堆起来意味着推送链路已经出问题了这时候再看日志往往已经晚了。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门