ARTICLE DETAIL

资讯详情

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

WebSocket分布式集群演进:从单机到高可用实时推送方案

WebSocket分布式集群演进:从单机到高可用实时推送方案 凌晨三点我盯着监控屏上的错误告警整层楼只剩下机房的散热风扇在响。某个节点的连接数已经逼近上限而外面的世界杯比赛正踢到加时赛——如果再往下走今晚的比分推送就会大面积断线。那天之后我把整个WebSocket推送链路从单机架构重构成了分布式集群过程踩了不少坑也沉淀了一套相对完整的方案。这篇博文就把这段演进过程原原本本拆给大家看单机是怎么撑过前一万用户的、集群化之后连接和session怎么管理、跨节点消息怎么做到不重不漏以及高可用和故障迁移的完整操作路径。需要说明的是这套方案的背景是体育直播平台的实时比分推送前端要求同一场比赛的所有在线用户都能在1秒内收到事件单个节点峰值连接数目标是10万以上。如果你也在做直播弹幕、行情推送、IM消息这类强实时业务这套从单机到分布式的思路基本可以直接照搬。1. 起步单机WebSocket服务是怎么撑过最初一万人的1.1 最开始的服务模型早期业务量不大一台4核8G的云服务器就能跑得很舒服。技术选型也很朴素前端通过Nginx的HTTP Upgrade机制把连接升级成WebSocket长连接后端用Go写了一个无状态网关服务所有在线连接维护在一个全局Map里。连接建立之后结构大概是这样的type Client struct { ID string UserID int64 MatchID int64 Conn *websocket.Conn SendCh chan []byte }每个连接对应一个ClientMatchID表示这个用户当前订阅了哪场比赛。推送的时候拿到某场比赛的ID遍历全局Map找到所有订阅了这场比赛的连接把消息丢进各自的SendCh由每个连接对应的goroutine负责写回前端。这套模型最直白也是大多数团队起步时的做法。它不需要专门的消息中间件不需要考虑跨节点通信甚至不需要引入Redis来做session管理。只要把连接池维护好、消息队列的缓冲区设计得当单机撑住一万到三万的实时在线连接完全没问题。1.2 单机的三个隐性天花板能上线和能扛住爆发是两码事。我们当时做了压测发现问题并不是单纯连接数上不去而是三个隐性天花板一起逼过来的第一个是连接数上限。操作系统默认的ulimit只有1024哪怕调到65535一个Node进程的FD数、内存、goroutine调度能力都有物理上限。Go的goroutine很轻量但每个WebSocket连接背后还有读缓冲、写缓冲、业务状态维护8G内存的机器跑到5万连接就明显感觉到GC压力。第二个是广播放大效应。看球赛的用户会集中在同一场比赛里一场热门比赛的在线用户可能占全站总连接的60%。服务端要按比赛维度做广播消息体虽然就几百字节但乘以几万连接就是一次数百万次的内存拷贝和写操作CPU很容易先撑不住。第三个是单点故障。Nginx后面只有一台后端服务进程一挂全站在线用户全部掉线。掉线之后所有客户端同时发起重连又会产生雪崩式的重连风暴。当时的监控曲线我记得很清楚一场焦点战开赛前半小时连接数线性上涨到了第70分钟某个客场球队进球推送峰值达到每秒40万次写操作CPU毛刺直接冲到90%某个时刻甚至出现了客户端收不到推送的秒级延迟。也是从那次之后我们下定决心做分布式改造。单机方案的能跑和集群方案的能扛中间差了整整一个量级。2. 集群化改造先解决连得上再解决推得对2.1 网关层的连接亲和与负载均衡分布式改造第一步是加机器。最简单的做法多起几个WebSocket服务节点前面用Nginx做负载均衡客户端通过ws://域名连接。这里有个关键点WebSocket是长连接一旦建立后续客户端和服务端的连接关系就固定了。所以负载均衡策略不能选普通的round-robin轮询——一个客户端断开重连后可能被分到另一台机器如果后端session不做迁移用户就被下线了。常规做法是Nginx配置ip_hash让同一个IP的请求尽量打到同一台后端节点upstream ws_cluster { ip_hash; server 10.0.0.2:8080 max_fails2 fail_timeout10s; server 10.0.0.3:8080 max_fails2 fail_timeout10s; server 10.0.0.4:8080 max_fails2 fail_timeout10s; keepalive 1024; } server { listen 80; server_name push.example.com; location /ws { proxy_pass http://ws_cluster; 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_read_timeout 3600s; proxy_send_timeout 3600s; proxy_buffering off; } }ip_hash其实也只是权宜之计。客户端如果从Wi-Fi切到4GIP变了照样会连接到另一台机器。后面我们的门户和App端统一走了注册中心客户端每次连接前先从HTTP接口拿一个gateway地址列表再按最少连接数策略选节点。这样即便IP变了只要用户标识稳定网关层就能根据用户ID做更精细的亲和。2.2 Session如何进行分布式管理连接上来了session怎么管理成了第一道坎。单机时全局Map就能搞定分布式节点上每个节点只知道跟自己在线的连接不掌握全局状态。我们一开始尝试了全量共享的方案用Redis集中存储所有在线连接信息键是client_id值是节点ID、用户ID、订阅的比赛ID。但这个方案的代价很大——每次心跳都要写一次Redis在线用户几万的时候连接信息频繁更新Redis的写瓶颈比WebSocket本身更早暴露。后来换成了本地缓存路由表的混合方案每个节点把与自己建立连接的客户端信息放在本地内存中连接对象、鉴权信息、订阅关系另有一张user_route表记录的是全局颗粒度的映射user_id - node_id, client_id这张表放进Redis并设置TTL每次心跳时刷新TTL查询的时候优先查本地Map本地没有再去Redis找对应节点然后通过跨节点消息通道转发。这套设计的好处是连接级别的状态比如发送通道、连接对象不需要跨节点同步性能接近单机跨节点同步的只是轻量级的路由坐标数据量小Redis完全扛得住。关于鉴权我们的做法是在WebSocket握手阶段用Query参数携带一次性token服务端先校验token再允许升级连接。token可以从HTTP服务签发有效期很短这样即使用户的连接被劫持影响的也只是这一条连接而不是整站账号体系。Netty项目里做WebSocket鉴权时经常见到在HttpServerCodec和WebSocketServerProtocolHandler之间加一个HttpRequestHandler专门处理token校验思路完全一样。3. 状态同步从广播到精确投递的架构升级3.1 按房间/比赛维度建模消息通道集群化之后最核心的问题来了A节点上的用户订阅了某场比赛这场比赛的比分事件却是从B节点上的推送服务发出来的怎么把消息跨节点送过去很多方案喜欢一张大网把所有节点连起来任意节点之间都可以互相转发消息看起来灵活实际上路由复杂度极高。我个人的经验是体育直播这类场景天然有业务维度的隔离属性——比赛就是房间事件归属于具体比赛。按比赛维度建模消息通道会让跨节点转发的拓扑简单得多。具体来说我们把每场比赛抽象成一个房间。每个WebSocket节点按需订阅若干房间Redis里存一份match_id - node列表的关系表。推送服务产生一条比分事件后只需要查询这张表确定当前这场比赛正在被哪些节点的用户订阅然后将消息定向投递给这些节点由各节点做本地的最终扇出。这个模型有几个明显的好处消息不会全局广播只会发到有订阅者的节点流量骤减节点的订阅关系是动态的用户进入/退出比赛订阅时更新节点级订阅表不需要精确到每一条连接做全局路由单台节点宕机后只需要把该节点从所有match_id的订阅列表里摘除不会影响其他节点。房间维度的建模其实和直播弹幕的做法很接近。很多弹幕系统的路由表就是维护直播间ID到节点列表的映射消息按直播间分发而不是按单个用户分发因为弹幕场景天然就是房间内广播。3.2 消息序列号与位点补偿机制状态同步还有一个容易被忽视的点消息是有序性的保证。体育比赛事件必须严格按时间顺序投递——进球事件的推送如果比角球事件的推送先到客户端就会显示错乱。我们给每场比赛单独维护了一个自增序列号。事件产生时按match_id维度递增并存储到Redis消息体里带上这个序列号。客户端收到消息后根据序列号判断有没有漏消息如果发现有gap客户端会主动向服务端发起一次补拉请求从断点位置开始重新推送这批消息。这个机制是典型的推送补偿双通道设计。对于实时性要求特别高的场景比如进球瞬间的庆祝动效推送通道保证低延迟对于可靠性要求高的场景比如比分面板的数值一致性补偿通道做兜底。举个具体数字一次进球事件从产生到推送到用户手机我们会拆成两个动作。第一个动作是写入事件存储并追加序列号耗时极短第二个动作是通过消息通道实时推送给在线用户。如果用户恰好处在断网或弱网状态App端会在网络恢复后通过补偿接口拉取断点之前的消息。实测下来断线重连的恢复时间控制在500毫秒以内而比分面板的最终一致性能在重连后3秒内收敛。3.3 用Redis Pub/Sub还是Kafka跨节点消息通道是集群方案的主动脉选型上我们纠结了很久。网上讨论最多的两条路是Redis Pub/Sub和Kafka。我两边都试过结论是没有绝对的好坏只有场景匹配不匹配。先说我踩过的坑。第一版我们图省事直接用了Redis Pub/Sub做节点间广播。它的好处是轻量、延迟低而且周围同学对Redis都很熟接入成本几乎为零。但它的致命问题是消息不持久化且消费能力受单节点性能限制。一旦某个网关节点处理不过来导致订阅端积压消息就悄悄丢了而且没有任何补偿机制。后来把核心链路换成了Kafka结构变成了这样推送服务产生事件写入Kafka的match_event主题每个网关节点作为一个消费组去消费自己关心的分区消费端拿到事件后查询本地订阅关系完成本地扇出。Kafka带来的收益非常明显分区机制让我们可以按match_id做key划分同一场比赛的事件保证进入同一个分区顺序性天然得到保证消费组机制让节点增减时自动做rebalance最重要的是消息落盘即使某个节点处理到一半宕机重启后可以从上次提交的offset继续消费。选型时我做了一张对比表贴出来供参考维度Redis Pub/SubKafka延迟亚毫秒级毫秒级可接受持久化无有按配置保留多天顺序性受网络影响分区内严格有序积压处理无可保留历史消息重放运维复杂度低中等适用场景节点间广播通知核心事件可靠投递最终我们的架构是双通道并存Redis Pub/Sub只做轻量的节点心跳和路由变更通知Kafka负责核心的比分事件投递。这个分工现在看起来还是很合理的。4. 高可用节点宕机、断线重连、故障迁移的那些事4.1 心跳与健康检查分布式系统的稳定性一半靠设计一半靠运维护栏。护栏的第一道就是心跳和健康检查。WebSocket协议自带ping/pong机制前端按30秒间隔发ping服务端必须在超时时间内回pong。我们在后端也对每个客户端维护了心跳超时检测超过90秒没有收到客户端的任何消息就判定为失活主动断开连接并清理路由表。节点间的健康检查则用另一套体系。每个网关节点会向注册中心上报自己的CPU、内存、在线连接数和最近1分钟消息处理速率。Nginx的健康检查器每5秒探测一次后端节点的/healthz接口连续3次失败就把节点从upstream里摘除。有一次正值早场英超Nginx日志里突然出现大量502表面上看是后端进程假死。排查后发现问题出在goroutine泄漏——某个版本的消息重试逻辑在Kafka消费失败时会无限重试且不设上限导致goroutine堆积、内存飙升最终GC卡死。从那以后我们给所有消费逻辑加了重试次数上限和熔断降级开关超出10次直接写死信队列人工介入处理。4.2 断线重连与消息幂等WebSocket是长连接网络抖动、切换Wi-Fi、服务器重启都会导致连接断开。客户端的断线重连机制必须做对否则要么是重连风暴把服务打垮要么是消息重发导致数据重复。我们的客户端重连策略分三档首次断线后立即重连连续失败后退避至2秒、5秒、10秒最多到30秒封顶断线期间如果业务层需要状态通过HTTP补偿接口直接拉取。服务端这侧需要解决的是消息幂等。推送通道里客户端可能收到重复消息原因可能是Kafka消费端重平衡后重复下发也可能是客户端断线前已经收到数据但ack丢了。我们的做法是给每条消息加event_id客户端按客户端ID维护最近收到的500个event_id做去重。这个方案简单直接效果也很稳实测重复推送率能压到万分之一以内。4.3 故障迁移与流量调度节点宕机后的流量迁移是整套方案里操作风险最高的环节。我们设计过一个标准流程每次演练都会完整走一遍Nginx先从upstream摘除故障节点新连接不再分发过去故障节点上的活跃连接由心跳超时自然断开客户端开始走断线重连逻辑客户端重连到其他健康节点后服务端根据user_route表重建路由关系如果故障节点之前订阅了大量比赛的房间其他节点会收到订阅变更通知把自己在Redis里的订阅关系补充完整。这套流程里最容易出问题的是第3步。重连后用户请求打到新的节点但新节点的本地缓存里没有这个用户的订阅关系如果不同步历史订阅用户会假在线——连接是通的但收不到任何推送。后来我们的做法是客户端重连上来后必须上报当前正在观看的比赛ID列表服务端按这个列表重建本地订阅同时更新路由表。既做状态同步又解决冷启动问题。5. 性能调优从指标到参数的逐项排查实录5.1 内核与网关参数WebSocket是长连接和HTTP短连接相比性能瓶颈更多集中在操作系统层和网络层。如果你的服务要撑10万级的长连接下面几个参数几乎是必调的# 文件描述符上限 ulimit -n 1048576 # TCP连接复用与探活 sysctl -w net.ipv4.tcp_tw_reuse1 sysctl -w net.ipv4.tcp_fin_timeout30 sysctl -w net.ipv4.tcp_keepalive_time600 sysctl -w net.ipv4.tcp_keepalive_intvl30 sysctl -w net.ipv4.tcp_keepalive_probes3 # 最大半连接队列 sysctl -w net.core.somaxconn65535 sysctl -w net.ipv4.tcp_max_syn_backlog65535Nginx层需要特别注意的是proxy_buffering off。HTTP短连接开缓冲没问题但WebSocket是双向流式通信一旦开了缓冲服务端推送的消息要等缓冲区满了才会被转发到客户端肉眼感知就是推送延迟忽高忽低。应用层的连接读写缓冲区也要精心设置。我们的经验是每连接读缓冲32KB、写缓冲64KB这一个配置能把内存占用减少30%以上。当时压测数据是8C16G的节点撑8万在线连接内存稳定在3.2GB以内GC毛刺在100ms以内。5.2 合并推送与批量发布客户端收到的比分事件看起来是一条一条的实际上服务端推送时我们做了非常激进的合并策略。进球、红牌、换人这类事件在一个比赛瞬间可能连发好几条如果每来一条就发起一次WebSocket写操作高频场景下的系统调用开销很大GC压力也会上升。我们的方案是给每个连接加了一个聚合发送队列同一个客户端在100毫秒窗口内的多条消息会合并成一条batch消息一次性写出。这个优化做完推送吞吐量直接翻了一倍。另外还要注意Kafka消费端不要逐条处理。我们一开始用单线程逐条消费事件每条都去做一次Redis查询、一次本地广播吞吐量死活上不去CPU却很高。后来改成批量拉取每次拉取最多500条消息批量处理、批量提交offset消费者的吞吐从每秒2000条涨到了每秒3万条以上。5.3 本地联调用IDEA多开进程模拟分布式节点分布式环境搭好了本地开发却遇到一个很现实的问题代码改了一行如何在本地模拟多节点、验证消息是否跨节点转发正常我们当时用了一个很土但很管用的办法——用IDEA同时启动多个服务实例通过不同端口模拟多个节点。具体操作是给网关服务加上端口参数Spring Boot项目启动时指定--server.port8080、8081、8082三个实例。IDEA的Run Configuration里复制三份启动配置分别设置不同的VM参数和应用参数点击并行启动。这样本地就有一个三节点的伪集群路由表注册到同一个本地RedisKafka连接同一套本地集群完全能模拟出跨节点消息流转的效果。我用这个方式做过一次很有价值的联调A节点上建立连接B节点上产生一条比赛事件通过Kafka分发到A节点验证A节点上的客户端能不能正确收到消息。这个测试在单机模式下根本跑不出来但通过IDEA多开进程整个链路20分钟就能跑通验证。6. 踩坑实录排查解决过的典型问题6.1 WebSocket 1006 非正常断开的真相前端的WebSocket连接经常报onclose code: 1006这个状态码的含义是连接异常关闭没有收到正常的Close帧。我们的业务里出现过三次大规模1006原因各不相同第一次是Nginx的proxy_read_timeout设置太短默认60秒就断开了长连接。这个问题最好排查调成3600秒即可。第二次是健康检查器把节点从upstream摘除后已有连接没有做优雅关闭。前端收到的是裸TCP断开自然就是1006。后来我们专门写了优雅关停逻辑节点收到SIGTERM后先告知Nginx摘除自己再向所有在线客户端发送Close帧并等待3秒最后才真正退出进程。第三次最隐蔽是客户端从H5页面能正常连接打包成App就连接不上。排查了半天发现是App端走了HTTPS代理代理对WebSocket的Upgrade头做了特殊处理导致服务端没有正确识别升级请求。最终在App的WebView设置里禁用了代理、放行WebSocket域名才解决。6.2 消息延迟与积压的告警定位有一次比赛日中午监控突然报警推送延迟从平时的200ms涨到了5秒。第一反应是Kafka消费出了问题但查看消费组lag一切正常。后来发现延迟的根子在Redis——某场比赛的订阅关系表大量过期重写节点之间反复推送订阅变更Redis的CPU被打满了。从那之后我们定了一个规矩节点级订阅关系表不做过期用主动通知的方式变更。由于订阅关系始终存在查询压力下降了一个量级延迟恢复到200ms以内。另外我们还为消息延迟设了双层告警业务侧统计从事件产生到推送完成的端到端耗时超过1秒告警系统侧监控Kafka消费组的lag超过1万条告警。两层告警互相校验能快速定位到底是业务链路的问题还是消息队列的问题。6.3 分布式锁误用与路由表冲突分布式改造过程中我们也踩过分布式锁的坑。最初想在节点摘除和订阅变更时用Redis分布式锁保证一致性结果锁的续期机制没做好持有锁的节点处理超时后锁已过期另一个节点抢到锁又执行了一遍变更逻辑两次变更互相覆盖路由表出现了短暂的不一致。后来我们删掉了大部分分布式锁的使用场景改用版本号条件更新的方式每次路由变更携带一个自增版本号Redis在更新时比较版本号只接受更新的版本。这样既避免了锁的复杂性又保证了变更最终一致。分布式锁在集群方案里是典型的看着优雅、用着难受的技术能用数据结构保证的就不要依赖锁。7. 复盘与实用建议整个改造做完我们的WebSocket集群从最初的3个节点扩展到了12个节点在线连接数从单机5万提升到了40万推送延迟P99稳定在800ms以内故障恢复时间从原来的分钟级缩短到了10秒以内。回看整个过程我觉得有三件事是决定成败的关键。第一是消息通道的选型Redis Pub/Sub虽然便捷但核心链路一定要有持久化和重放能力这是可靠性的底线。第二是按业务维度建模路由体育比赛就是房间消息按房间分发比按用户逐条路由简单太多。第三是断线补偿机制只做推送不做补偿的方案是脆弱的移动网络环境下一旦弱网抖动用户体验会全线崩溃。最后分享一个我保留至今的排查习惯每次故障处理完我们都会把根因、现象、解决手段整理成一条速查记录沉淀进团队的运维手册里。比如1006断线的几种原因、消息延迟的排查路径、Kafka重平衡导致重复消费的处理办法这些记录在之后的多次大促、杯赛、焦点战中都实实在在救过场。分布式系统的演进没有终点但把这些现场经验留住了后面的人就能少走很多弯路。
返回列表