10万在线WebSocket弹幕系统:集群架构与调优实战
发布时间:2026/9/16 5:18:49 作者:尧图编辑部 阅读量:1,286

先说结论如果你正打算给直播、体育赛事、电商大促这类高并发实时互动场景做长连接服务这篇文章里的方案可以直接拿来当蓝本。我负责过的一套体育直播弹幕系统峰值阶段稳定扛过10万人在线弹幕从发出到所有人都看到体感上就是“秒达”。核心组件就三个词WebSocket、Redis、Kafka外加一层Nginx做接入路由。下面把整个设计思路、踩坑记录、调优参数全部摊开讲。1. 业务场景与架构选型1.1 10万人在线到底意味着什么10万人在线很多人第一反应是“不就并发十万吗”其实完全不是。这里说的是同时建立了10万条WebSocket长连接不是10万次HTTP请求。两者差着量级。长连接意味着服务端要长期维护10万个活跃socket每个连接都有内存占用、读写缓冲、心跳定时器。假设每个连接平均占用内存30KB到50KB光是连接本身就要吃掉3GB到5GB堆内存这还没算业务对象和消息缓冲。加上比赛进行中弹幕密集单场热门赛事高峰期每秒会产生大几千条弹幕服务端要把每条弹幕广播到在线客户端那就是“每来一条消息往10万个连接里各写一遍”算下来每秒的写操作量级是几千万次。所以这个场景的核心矛盾不是“接收弹幕”而是广播路径上的写放大。架构选型必须围绕这个矛盾来设计。1.2 为什么是WebSocket而不是轮询或者SSE最早的直播弹幕用过HTTP轮询但效果很差。HTTP协议每次请求都要带上完整的头部一次轮询可能只有几条弹幕但请求头动不动就几百字节数据利用率惨不忍睹。更重要的是轮询是“客户端主动拉”弹幕延迟取决于轮询间隔想做到一秒内到达客户端就得每秒请求一次10万人同时每秒请求一次接入层直接被打穿。SSEServer-Sent Events解决了“服务端主动推”的问题但它有两个硬伤一是单向的客户端只能收不能发弹幕发送还要另走HTTP接口链路被拆成两半二是SSE本质上还是跑在HTTP上虽然用的同一个连接但代理层、浏览器层面都有各种兼容问题。所以我们最终选了WebSocket双向通信、头部开销小、天然支持服务端推送一次握手之后就是纯粹的二进制帧或文本帧是实时互动场景最合适的载体。1.3 单机方案为什么扛不住早期系统是单机部署一台8核16G的云主机Nginx代理到Spring Boot应用应用里用原生的WebSocket注解处理器。5000人在线的时候一切正常一万人开始出现广播延迟两万人的时候CPU直接飙到90%以上弹幕明显卡顿出现“比赛都进球了弹幕还没刷出来”的尴尬情况。复盘下来有三个瓶颈一是连接数垒在单机上单进程能承载的连接上限受文件描述符和内存制约调高ulimit之后能到几万但堆内存和GC压力已经很恐怖了。二是广播是同步循环伪代码大概是for (session: sessions) session.sendMessage(msg)一个慢客户端写不进去整个循环就卡住后面的连接全部等待。这个“木桶效应”在单机场景下就会导致整体延迟飙升。三是扩容只能靠堆配置加到32G内存之后单台能扛四五万连接但不敢再加了因为一旦宕机就是四五万人同时断线重建连接的风暴能把下游服务打挂。到这一步向集群演进已经不是什么“技术追求”而是活下去的必然选择。1.4 集群整体架构的组成最终落地的集群架构包含四个层次接入层一组Nginx节点对外暴露WebSocket入口用ip_hash策略做负载均衡。为什么不用round-robin因为WebSocket是长连接连接一旦建立就不该在Nginx层再迁移跨节点迁移需要专门机制ip_hash可以把同一IP的客户端尽量固定到同一台后端节点减少不必要的跨节点转发。WebSocket服务层若干无状态应用节点跑Spring Boot内嵌Netty或直接用Netty编写。节点之间天然不共享本地会话所有连接信息注册到Redis节点通过订阅Redis频道感知其他节点的在线状态。协调与存储层Redis Cluster保存会话映射和在线状态Kafka作为弹幕消息的传输管道。Redis负责“谁在哪台机器上”的路由问题Kafka负责“弹幕怎么派发到所有节点”的广播问题。业务后端弹幕过滤、敏感词检测、用户等级、礼物逻辑等业务处理通过Kafka解耦不直接参与长连接读写。这个架构的核心思想是连接无状态化路由集中化广播异步化。连接无状态化让节点可以随意扩缩容路由集中化让任意节点都能找到任意连接的归属广播异步化让业务处理和连接分发互不阻塞。2. 服务端核心模块设计与实操2.1 连接接入与鉴权别把WebSocket当成免鉴权通道WebSocket握手阶段其实是HTTP Upgrade请求所以鉴权可以挂在握手前完成。常见做法是客户端在URL上带token比如ws://gw.example.com/live/1001?tokenxxx服务端在握手拦截器里校验token。我这里给出一个简化的Spring Boot接入示例用的是Spring WebSocket的HandshakeInterceptorpublic class AuthHandshakeInterceptor implements HandshakeInterceptor { Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, MapString, Object attributes) { String query request.getURI().getQuery(); String token extractParam(query, token); Long userId authService.checkToken(token); if (userId null) { response.setStatusCode(HttpStatus.UNAUTHORIZED); return false; } attributes.put(userId, userId); String roomId extractParam(query, roomId); attributes.put(roomId, roomId); return true; } }注意token放在query里虽然方便但有被日志记录导致泄露的风险。生产环境建议token放在Sec-WebSocket-Protocol头里传递或者在握手后用自定义协议帧做一次应用层鉴权两种方式可以根据安全等级选择。服务端连接建立之后要向Redis写入会话路由信息。键的格式我建议用ws:route:room:{roomId}值是一个Hashfield是userIdvalue是{nodeId}:{connectionId}。这样一个房间的所有在线用户分布在哪些节点上一目了然。同时还要维护一个全局维度按userId维度记录这个用户当前在哪个房间、哪个节点方便做点对点消息投递。2.2 连接路由与会话管理节点之间怎么知道“谁在哪儿”集群里最尴尬的问题就是用户A连的是节点1用户B连的是节点2A给B发私信节点1怎么知道B在哪这时候Redis就是全局会话表的载体。我设计的路由表结构如下ws:conn:{userId}-{nodeId}:{connId}:{roomId}String类型带TTL每次心跳续期。这个键用来做私信和踢人。ws:room:{roomId}- HashfielduserIdvalueconnectionId。这个键用来做房间维度的热数据查询比如后台要看一个直播间哪些人在线。这里有个细节如果只在Redis里存路由节点查询路由就要走一次网络IO高并发场景下会拖慢消息链路。所以每个节点本地还有一层Caffeine缓存存“最近查过的路由”命中率极高。因为直播场景里弹幕广播根本不需要查路由只有点对点消息才查私信频率远低于弹幕本地缓存基本能挡住大部分流量。当用户断开连接时节点需要清理自己的本地会话同时删除Redis里的路由键。删除这里要小心的坑是节点的连接其实有“换人”的可能。比如用户断线重连新连接连到了节点3但节点1还没来得及清理旧路由此时Redis里就存在两条路由旧路由会抢占新路由。我的解决办法是写入前先对比旧值如果旧值里的nodeId不等于当前节点的nodeId就主动通知旧节点把旧连接踢下线。这就是典型的“踢旧push新”流程保证同一个userId全局只有一条有效连接。2.3 弹幕广播链路Kafka在这里扮演的角色弹幕发送的完整链路是这样的用户通过WebSocket发送一条弹幕文本客户端收到的是文本帧。节点解析出弹幕内容、房间号、用户名把原始消息丢给Kafka的danmaku-inputtopic。业务消费者从topic里拉消息做敏感词过滤、等级校验、频率限制通过后写入danmaku-outputtopic。WebSocket服务层的每个节点都消费danmaku-output把经过校验的弹幕写到本节点维护的、属于该房间的所有连接上。为什么中间要插两层Kafka而不是节点直接处理完就发送原因有三解耦业务敏感词过滤、频率控制、黑名单这类业务逻辑不应该阻塞长连接线程。长连接服务最忌讳在IO线程里做耗时操作Kafka天然把生产者和消费者分开。削峰比赛突然进球弹幕量可能瞬间高出平时十倍。Kafka可以缓冲峰值消息消费者按速率拉取防止后端服务被瞬时洪峰击垮。一致性所有节点消费同一个topic广播路径就有了统一的“顺序约定”。虽然业务上不要求所有弹幕严格有序但至少同一条弹幕不会出现“节点1先发、节点2后发”的乱序感。Kafka topic的分区数设置也很关键。danmaku-output的分区数建议和WebSocket节点数一致或者设成节点数的倍数。这样每条消息只会分发到部分分区每个节点消费其中若干个分区避免所有节点都消费全部分区导致重复消息处理开销。每个节点消费到一条弹幕后遍历本节点的在线连接表找出属于目标房间的session逐一发送。这里有个性能优化点不要用for循环直接调用session.sendMessage()。Netty的写操作是异步的但大量并发写同一个Channel时如果超过高水位线默认64KB会触发Channel.isWritable()变成false此时如果你还硬写数据会堆积在Netty的outboundBuffer里内存暴涨延迟飙升。正确的做法是用ChannelGroup管理房间内的连接再用group.writeAndFlush(msg)做广播。ChannelGroup内部会帮你把消息逐条写到每个Channel并遵循Netty的写缓冲水位控制比手写循环安全得多。如果你用Spring WebSocket的SimpMessagingTemplate底层也封装了类似的逻辑但要注意它的广播接口在高并发下性能不如直接操作Netty的ChannelGroup。2.4 点对点消息投递私信、礼物、系统通知怎么找到人点对点消息私信、礼物通知、系统公告和广播不一样它必须精准投递给某个userId。投递路径需要经过三个步骤调用方拿到目标userId查Redis里的ws:conn:{userId}路由得到{nodeId}:{connId}:{roomId}。如果目标节点就是当前节点直接查本地连接发送如果目标节点在别处通过节点间共享的Redis频道向目标节点发送一条ws:notify:{nodeId}的消息内容包含目标connId和消息体。目标节点消费该频道消息找到连接并发送。这里有个很实用的经验节点间通信不要自己定义私有协议直接用Redis的Pub/Sub就够了。Redis Pub/Sub虽然不持久化但点对点投递的消息量级不大而且对实时性很敏感刚好匹配它的特点。广播类型的消息走Kafka点对点消息走Redis Pub/Sub两条管道互不干扰架构上非常清晰。3. 集群化改造的关键细节3.1 Nginx负载均衡策略长连接场景不能用默认轮询WebSocket是长连接一旦建立就维持很久如果Nginx用默认的round-robin策略每个新建连接的分布是均匀的但长连接不会频繁断开重建久而久之可能出现某些节点因为建连时间集中而负载偏高另一些节点却闲着。再加上同一用户断线重连时如果被分配到不同节点路由表会频繁变更增加跨节点通信开销。生产环境的Nginx配置我建议这样写upstream ws_cluster { ip_hash; server 10.0.1.11:8080 weight1 max_fails2 fail_timeout30s; server 10.0.1.12:8080 weight1 max_fails2 fail_timeout30s; server 10.0.1.13:8080 weight1 max_fails2 fail_timeout30s; keepalive 1024; } server { listen 80; server_name ws.example.com; location /live/ { proxy_pass http://ws_cluster; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }几个容易踩的坑proxy_read_timeout和proxy_send_timeout必须调大Nginx默认是60秒如果收到不到数据会主动断开连接。虽然WebSocket有心跳但心跳间隔如果大于这个超时时间连接照样会被掐断。我一般直接设成3600秒配合应用层心跳。keepalive 1024是Nginx到后端节点的空闲长连接缓存数。这个参数很多人忽略其实很关键。没有它Nginx每次转发WebSocket升级请求都要新建到后端的TCP连接握手性能会差很多。ip_hash有个著名的坑客户端如果通过移动网络上网IP可能会变同一个用户在不同时刻可能被哈希到不同节点。这个问题可以通过断线重连机制兜底所以ip_hash不是完美的但配合好重连逻辑问题不大。如果要求更均衡可以把ip_hash换成一致性哈希按userId计算哈希配置会稍微复杂一些但节点增减时的连接迁移范围会更小。3.2 连接迁移机制缩容和发布时怎么不掉线集群必然要面对发布和扩容缩容。如果直接停节点上面几万条WebSocket连接会瞬间断开客户端集体重连造成“重连风暴”。所以必须有一套连接迁移机制。我的做法是“优雅下线”三步走标记下线运维把某节点设置为draining状态节点不再接受新的连接分配但保留已有连接。通知迁移节点后台线程扫描本机所有连接向每个客户端发送一个“reconnect”指令帧。客户端收到后断开当前连接重新走一次建连流程Nginx会把它分发到其他健康节点。等待清空节点等待所有连接迁移完毕或者等待超时比如60秒后强制退出。因为所有连接的信息都实时同步在Redis里新连接建立后自动更新路由旧节点退出不会影响消息投递。这个方案比“先停服再重连”体验好得多用户体感是“闪断了一下”弹幕几乎无感知。但要注意迁移期间新连接流量会全部涌到其他节点要提前评估其他节点的余量。我的经验是单节点负载不超过60%时可以放心让一个节点迁移负载超过80%最好分批次迁移一次只迁一半。3.3 心跳与超时策略参数不是随便拍脑袋定的WebSocket心跳的作用有两个一是维持Nginx的代理连接不超时二是让服务端及时发现死连接。客户端主动发Ping帧服务端回Pong帧这是标准做法。心跳间隔怎么定我的参数是客户端每30秒发一次Ping服务端连续两轮60秒没收到Ping就判定连接死亡清理本地session和Redis路由。为什么是30秒而不是10秒或120秒这背后有个计算逻辑心跳包平均大小约20字节10万个连接每30秒一次心跳每秒产生的额外流量大约是100000 * 20 / 30 ≈ 66KB/s对带宽来说可以忽略。但如果把间隔压到10秒心跳流量变成三倍而且Nginx层处理的包数量激增纯粹增加无意义的CPU开销。反过来间隔超过120秒一旦网络闪断服务端要等两分钟才发现用户体验很差。30秒是一个成本与感知之间的折中值。服务端侧的逻辑是这样写的public class HeartbeatHandler extends ChannelInboundHandlerAdapter { private static final int READ_TIMEOUT_SECONDS 60; Override public void channelActive(ChannelHandlerContext ctx) { // 每隔30秒检查一次这个连接是否在最近60秒内有收包 ctx.executor().scheduleAtFixedRate(() - { long lastReadTime getLastReadTime(ctx.channel()); if (System.currentTimeMillis() - lastReadTime READ_TIMEOUT_SECONDS * 1000L) { ctx.close(); } }, 30, 30, TimeUnit.SECONDS); } }注意这里的心跳判断是基于“最近一次收到任意数据包”的时间而不是只依赖Pong帧。因为客户端在发弹幕、发互动消息时本质上也是一种活性证明按收包时间判断更不容易误杀正常连接。3.4 断线重连与消息补偿客户端也不省心服务端设计得再好网络环境永远是不可控的。客户端必须做好断线重连和消息补偿不然用户会发现闪断之后弹幕少了一段。客户端重连策略建议用指数退避随机抖动第一次断线立即重连第二次等2秒第三次等4秒最大不超过30秒。每次重连前加一个0到1000毫秒的随机抖动防止大量客户端同时重连形成“重连潮”尤其是一场大型比赛结束用户集体退出页面再进来时。消息补偿的思路是客户端在每次心跳时带上自己收到的最后一条弹幕序号服务端根据这个序号找到断线期间的增量消息通过单独的补偿通道发给客户端。这里有一个取舍如果只要求弹幕“尽量不丢”可以把补偿设计成只补最近30秒的消息如果要求严格有序就要在服务端为每个房间维护一个短暂的消息序号寻址窗口用Redis的List或者Kafka的offset来实现。我们的业务属于“丢一两条不影响体验”的类型所以补偿窗口只保留最近50条弹幕序号放在客户端本地Storage里。重连成功后客户端发一条sync请求服务端把大于本地序号的弹幕批量返回。这样做的好处是补偿逻辑简单不会占用过多内存。4. 性能调优与压测实录4.1 一次压测把连接压崩的复盘上线前做了一次压测目标5万在线结果跑到3万左右部分客户端开始掉线。排查发现后端节点的CPU不高内存也正常但GC日志里Full GC非常频繁每次停顿超过3秒。罪魁祸首是Netty的DefaultChannelGroup。它内部用了一个ConcurrentHashMap来存Channel但每次广播时要对所有Channel做一次遍历。10万连接广播一条弹幕遍历的代价本身就高更糟糕的是我们当时为了统计在线数在广播循环里做了channelGroup.size()之类的额外操作这些每次调用都是O(n)的。再加上GC要扫描大对象Old区很快就满了直接触发Full GC。解法有两个层面。一是缩短对象生命周期广播的消息对象尽量复用用ByteBuf而非创建新String二是把“全量广播”拆成“分房间广播”避免所有连接共用一个大Group。每个房间维护自己的小型ChannelGroup广播时只遍历本房间的连接规模从10万降到几千GC压力和CPU占用都大幅下降。4.2 Netty线程模型和背压处理Netty默认的EventLoop线程数是CPU核数的两倍。对于WebSocket服务这个默认值通常足够了不需要盲目调大。关键是理解一个Channel上的所有操作都由同一个EventLoop线程串行执行所以同一个连接上绝对不会出现并发写冲突。但是不同连接可能被分配到不同EventLoop广播一个房间的消息时会跨多个EventLoop需要小心并发写放大。背压处理我总结成三条原则写消息前先检查channel.isWritable()如果返回false说明对端消费不过来缓冲区已经高于高水位这时候不要再往这个channel写把这部分消息挂到pending队列。给pending队列设一个上限比如1000条。超过上限直接丢弃队列头部的旧消息保证新消息能尽快送达。对弹幕产品来说“丢弃旧的保留新的”比“全部堵住”体验更好。用Netty的ChannelFutureListener监听写完成事件如果某个连接连续出现多次写失败直接关闭该连接让客户端重连。我在代码里给每个WebSocket连接加了三个数字指标已发消息总数、写失败次数、pending队列长度。通过Prometheus暴露出来再配上Grafana看板业务高峰期可以实时看到哪些节点存在背压风险。有一次压测就是通过这个看板发现两台节点的pending队列持续堆积排查后发现是那两台机器所在的机房网络出口带宽不够加带宽后解决了。4.3 JVM和内核参数调优用了Netty的WebSocket服务JVM参数配置对稳定性影响很大。我最终稳定运行的一套参数是java -Xms8g -Xmx8g -Xmn4g \ -XX:UseG1GC \ -XX:MaxGCPauseMillis50 \ -XX:ParallelRefProcEnabled \ -XX:MaxDirectMemorySize2g \ -Xlog:gc*:gc.log:time,uptime,level:filecount10,filesize100m \ -jar ws-server.jar几个关键点说一下-Xms和-Xmx设置成相同值避免堆动态伸缩引起的性能抖动。新生代设成4GB保证大量短生命周期对象比如弹幕消息对象在Young GC阶段就被回收不会晋升到老年代。Netty使用堆外内存Direct Memory做网络读写缓冲堆外内存上限要用MaxDirectMemorySize单独控制不设的话默认等于堆上限遇到大流量可能触发OutOfMemoryError。内核层需要调的文件描述符上限和TCP参数# /etc/sysctl.conf net.core.somaxconn 65535 net.ipv4.tcp_max_syn_backlog 65535 net.ipv4.ip_local_port_range 1024 65535 net.ipv4.tcp_tw_reuse 1 net.ipv4.tcp_fin_timeout 15再配合ulimit -n 1024000提高进程文件描述符上限。调完之后单机连接数从原来的5万提升到16万才触顶而且是在内存足够的前提下。这里强调一下调优不是无脑加参数每个参数的变更都应该有压测数据作为依据。4.4 消费端慢导致弹幕延迟飙升Kafka消费链路也踩过一个大坑。由于danmaku-output的消息量很大分区分给了各个WebSocket节点但某个节点因为GC停顿或CPU抢占消费速度跟不上生产速度Lag越来越多导致这一节点上的用户弹幕延迟明显比别的节点高用户就会投诉“我看不到弹幕了”。解决这个问题有两个方向一是增加分区数让每个节点消费更多分区并行处理但分区数超过节点数之后单个节点的并行度并不会继续提高因为同一个分区只能被同一个消费者组里的一个实例消费。二是控制消费速率在消费端引入批量处理和线程池。Kafka消费者拉下来的消息先放进一个内存队列由业务线程池批量写入ChannelGroup。这样即使Kafka消费端短暂积压也不至于阻塞拉取线程。最终我采用的方案是后者。具体做法是消费者每批拉取500条消息然后拆成4个小批次交给4个线程并行广播。这样单节点内不同房间的广播可以并行推进而同一房间内的消息仍然按顺序发送既提高了吞吐又保证有序性。4.5 压测数据参考最后放一组我们压测环境的参考数据配置为3台8核16G的WebSocket节点2台Kafka节点3台Redis Cluster指标数值总连接数100,000弹幕发送峰值5,000条/秒单节点广播QPS约160万次写/秒弹幕端到端延迟P99 约320毫秒节点CPU峰值约70%堆内存峰值约6.2GB这个数据的含义是单条弹幕从用户发出到所有在线客户端收到最慢的10%也不会超过320毫秒大部分在200毫秒以内体感接近“秒达”。如果继续加节点水平扩容到5台、10台都能线性扩展瓶颈主要在Nginx层和Redis的访问延迟上服务本身已经没有单点问题。5. 常见问题速查表整理一下我在这个项目中遇到频率最高的问题以及排查结论做成一个速查表遇到类似症状可以先对照排查。问题现象根本原因解决方案连接建立后几秒内断开前端报1006Nginx代理层丢弃了WebSocket升级请求检查Nginx是否配置了Upgrade和Connection头proxy_read_timeout是否太小所有客户端每隔一段时间集体断开一次应用层心跳间隔大于Nginx或SLB的空闲超时应用心跳间隔必须小于代理层超时建议2/3原则某节点CPU高但连接数不高广播逻辑遍历了全量连接出现写放大按房间拆分ChannelGroup只遍历目标房间Redis路由被反复覆盖用户被踢下线用户断线重连后新旧连接同时在线写入路由前比对旧值通知旧节点踢掉旧连接Kafka消费Lag持续增长消费者业务逻辑阻塞拉取线程消费者拉取与处理分离用线程池异步批量发送突然出现大量TIME_WAIT客户端频繁建连断开调整TCP参数开启tcp_tw_reuse客户端采用指数退避重连单节点Full GC频繁大量短生命周期对象进入老年代调大新生代使用G1并设置目标停顿时间减少直连写大对象弹幕顺序错乱多个线程并发写同一个Channel或同一个房间同一房间的广播串行化或者用ChannelGroup确保写的有序性某房间弹幕一直无法送达该房间的ChannelGroup没有正确注册新连接检查连接建立时是否加入了房间对应的Group断开时是否移出WebSocket集群的难点从来不在“连接怎么建立”而在于连接建立之后的一整套状态管理、路由同步、广播分发、故障恢复机制。希望这套方案能给你一些参考特别是那些“不试不知道”的细节比如Nginx超时参数、ChannelGroup的坑、Kafka消费解耦、踢旧建新流程每一个都花过不少时间去填。如果你也在做类似的实时互动场景欢迎对照着压测一遍自己的系统。