ARTICLE DETAIL

资讯详情

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

WebSocket弹幕集群架构:高并发低延迟实时广播实战

WebSocket弹幕集群架构:高并发低延迟实时广播实战 1. 项目概述为什么体育直播平台的弹幕不能“卡一下”“体育直播平台10万人同时在线弹幕一秒送达”——这句话不是营销话术而是用户对实时互动体验的底线要求。我做过三年体育类直播后台架构从单机WebSocket服务起步到支撑中超、CBA季后赛峰值12.7万并发连接的集群系统踩过太多坑。弹幕延迟超过800毫秒用户就会觉得“卡”超过1.5秒大量用户会反复点击发送、刷屏重发反而加剧后端压力而一旦出现“弹幕发不出去”或“别人发的我看不见”投诉率当天就能翻三倍。这不是功能问题是信任崩塌。核心关键词就三个WebSocket、集群、弹幕。但它们组合在一起就不再是简单地把WebSocket服务多起几台——它直面的是高并发写入低延迟广播状态强一致性瞬时流量脉冲四重压力。体育赛事有天然的节奏进球、红牌、加时赛哨响瞬间涌出数万条弹幕像海啸拍向服务器。这时候单点WebSocket服务扛不住消息队列积压Redis缓存击穿下游消费延迟整个链路就“糊”了。我们最终落地的方案不是堆机器而是用分层解耦精准路由异步削峰状态隔离的组合拳在不牺牲实时性的前提下把10万并发稳稳接住平均端到端延迟压在320ms以内含网络RTT99分位延迟680ms。这篇文章不讲理论模型只说我们怎么一步步把“一秒送达”从PPT变成线上真实跑着的系统包括每个关键决策背后的算账过程、实测数据和血泪教训。2. 整体架构设计为什么必须放弃“一个集群打天下”的幻想2.1 传统单集群模式的致命缺陷很多团队一开始会想“WebSocket不就是长连接吗上个Nginx做TCP负载均衡后端起10台WebSocket服务连上Redis共享在线用户列表再挂个Kafka收发弹幕搞定。”我试过上线三天就被打趴。问题出在三个地方第一连接状态无法真正共享。Nginx的IP Hash或最少连接数调度只能保证同一个客户端IP始终落到同一台后端但用户换WiFi、切4G/5G、APP后台被杀再拉起IP就变了。结果就是用户A在机器1上发弹幕用户B在机器2上机器2查Redis发现用户A“不在线”直接丢弃——弹幕根本没广播出去。我们监控里看到过单日23%的弹幕“已接收但未广播”根源就在这儿。第二广播风暴让Redis成为瓶颈。早期我们用Redis Pub/Sub所有机器都订阅同一个频道。一条弹幕进来Kafka消费者推给所有机器每台机器再遍历自己维护的本地连接列表去发。问题来了10万台机器每台平均维护1万连接一条弹幕就要触发10万次本地for循环10万次write系统调用。CPU在广播线程上直接飙到95%GC频繁连接开始超时断开。更糟的是Redis Pub/Sub本身不保证消息不丢网络抖动时某台机器掉线几秒再连上就永远收不到那几秒的弹幕。第三瞬时脉冲让Kafka积压不可控。进球瞬间1秒内涌入4.7万条弹幕实测数据。我们的Kafka Topic只有3个分区单分区吞吐上限约1.2万条/秒立刻积压。消费者来不及处理延迟从毫秒级跳到秒级用户看到的弹幕全是“30秒前的旧闻”。我们曾紧急扩容到12分区但分区数翻4倍消费者实例也得同步翻4倍而消费者实例又依赖JVM堆内存和GC扩容后反而因GC停顿导致整体吞吐下降。提示别迷信“消息队列能解决一切”。Kafka擅长高吞吐、可持久化但不擅长低延迟广播。把它当“弹幕中转站”用等于用货车拉快递——运得多但送到门口慢。2.2 我们采用的四层分治架构我们彻底重构为四层接入层 → 路由层 → 广播层 → 存储层。每一层只干一件事且彼此解耦。接入层Ingress用自研的LVSKeepalived做四层TCP负载不碰业务逻辑。所有WebSocket Upgrade请求按room_id % 1024哈希到1024个虚拟节点再映射到物理机器。这个哈希值全程透传后续所有环节都基于它路由。好处是同一个房间的用户100%落在同一组机器上连接状态天然局部化。路由层Router这是核心大脑。它不处理连接只维护一张“房间→机器组”的映射表存在etcd里TTL 30秒心跳续期。当用户A加入room_123接入层把请求打到某台RouterRouter查表发现room_123当前由ws-group-07负责就返回该组的VIP地址客户端重定向连接。这样连接建立前就完成了精准路由避免了“先连再查”的状态不一致。广播层Broadcaster这才是真正发弹幕的地方。每个ws-group内部有3台机器组成小集群用Raft协议选主。主节点负责接收本组所有弹幕从节点只做状态同步。广播时主节点把弹幕序列化成二进制通过零拷贝的sendfile()系统调用直接推给本组所有在线连接。不经过Redis不走Kafka纯内存操作。实测单主节点可稳定广播8000条/秒弹幕延迟15ms。存储层Storage只做两件事一是用MySQL分库分表存弹幕历史按room_id % 64分64库每库128表二是用Elasticsearch建倒排索引支持按用户ID、关键词、时间范围快速检索。它完全不参与实时链路只是事后归档。这个设计的关键在于把“连接管理”和“消息广播”彻底分离。接入层管“谁连在哪”路由层管“谁该连哪”广播层管“怎么发最快”存储层管“发完存哪”。四层之间用轻量级gRPC通信协议定义清晰任何一层挂了都不影响其他层。上线后单点故障率下降92%扩容只需增加Router节点或Broadcast组无需动其他层。2.3 为什么不用K8s Service做服务发现热搜词里一堆K8s集群搭建但我们生产环境没用K8s Service做WebSocket服务发现。原因很实在K8s Service的Endpoint更新有延迟默认10秒而体育直播的房间创建是动态的——导播切到新球场后台API立刻创建room_456用户3秒内就得连上。如果等K8s慢慢同步Endpoint用户会看到“正在连接…”转圈超过5秒流失率飙升。我们用etcd自研Router从房间创建到Router感知并写入映射表全程800ms配合客户端3次指数退避重试100ms/300ms/900ms99%用户在2秒内完成重定向连接。这800ms就是用户体验的生死线。3. 核心细节解析弹幕“一秒送达”的技术锚点3.1 连接保活与断线重连别让网络抖动毁掉实时性WebSocket不是“连上就万事大吉”。移动网络下用户进出电梯、地铁隧道、Wi-Fi切换连接中断是常态。我们统计过单个用户每小时平均经历2.3次短暂断连5秒。如果每次断连都重新走完整握手流程HTTP Upgrade SSL握手 认证光TLS握手就要耗掉300~600ms弹幕延迟直接破秒。我们的解法是双通道保活 Token续期。双通道保活客户端除了主WebSocket连接额外建立一个极简的HTTP长连接/ping响应头带Connection: keep-alive。这个连接只传心跳包无业务数据开销极小。当主连接异常断开时客户端立刻用这个HTTP连接向Router发起/reconnect?tokenxxxlast_seq12345请求。Router验证token有效JWT有效期2小时并检查last_seq是否在本组广播窗口内我们维护最近10秒的弹幕序列号缓存如果是直接返回当前广播组的VIP和新的临时token客户端秒级重连无需重新认证。Token续期初始连接时服务端下发的JWT token里exp字段设为2小时但nbfnot before设为当前时间30分钟。客户端在nbf到期前5分钟自动发起/refresh请求Router校验旧token签名后签发一个新tokennbf再延30分钟。这样token永远有30分钟缓冲期避免因时钟漂移导致的意外过期。注意别用setInterval固定间隔发心跳。我们实测发现iOS Safari在页面后台时会节流JS定时器心跳间隔可能拉长到30秒以上。改用requestIdleCallbacksetTimeout兜底确保前台/后台都能稳定保活。3.2 弹幕广播的零拷贝优化从320ms到180ms的跨越广播延迟的大头往往不在业务逻辑而在数据拷贝。早期版本一条弹幕要经历Kafka Consumer反序列化 → 构造Java对象 → 序列化成JSON字符串 → 写入Netty Channel → Netty ByteBuf复制 → OS内核Socket Buffer复制 → 网卡DMA发送。光内存拷贝就4次加上GC压力P99延迟卡在320ms。我们做了三步改造协议扁平化弃用JSON自研二进制协议。头部4字节魔数0xCAFEBABE 2字节版本号 4字节body长度 bodyUTF-8编码的弹幕文本用户ID时间戳。序列化/反序列化耗时从1.2ms降到0.08ms。内存池复用Netty的PooledByteBufAllocator配置为maxOrder11支持最大4MB缓冲区预分配1024个4KB的CompositeByteBuf。广播时直接从池里取一个bufwriteBytes()填入二进制数据发完release()归还。避免频繁申请释放堆外内存Full GC频率下降76%。零拷贝推送关键一步。NettyChannel.writeAndFlush()默认会把数据拷贝到内核Socket Buffer。我们改用FileRegion将弹幕数据先写入DirectByteBuffer再用channel.write(new DefaultFileRegion(buffer, 0, buffer.readableBytes()))。这样数据从JVM堆外内存经DMA控制器直接送入网卡绕过内核Socket Buffer拷贝。实测单条弹幕网络栈耗时从110ms降到35ms。这三步下来端到端P99延迟从320ms压到180ms提升近一倍。但代价是协议不兼容老客户端我们用了灰度发布——新协议加X-Proto: v2header老客户端走JSON路径新客户端走二进制路径两周平稳过渡。3.3 房间状态同步如何让10万人知道“谁在说话”弹幕要显示用户名、头像、等级就得实时知道用户状态。如果每条弹幕都去查MySQLQPS轻松破10万DB直接跪。我们用两级缓存 增量同步L1本地Caffeine缓存。每个Broadcast节点内存里用Caffeine Cache存Mapuser_id, UserInfo最大容量10万expireAfterWrite 10分钟。UserInfo对象精简到只有nick_name、avatar_url、level三个字段序列化后200字节。L2Redis Cluster缓存。用Redis的Hash结构HSET user:12345 nick_name 张三 avatar_url https://a.com/1.jpg level 5。TTL设为30分钟比L1长避免L1失效时全量打Redis。增量同步用户资料变更如改名、升级前端调/api/user/update后端更新MySQL后发一条Kafka消息user_update:{user_id:12345, field:nick_name, value:李四}。所有Broadcast节点订阅此Topic收到后只更新本地Caffeine缓存对应字段不查DB。这样状态变更1秒内全网可见。我们压测过10万用户同时在线每秒5000次资料查询L1命中率92.7%Redis QPS稳定在400左右MySQL零压力。缓存穿透不存在的——用户ID是数字我们用布隆过滤器RedisBloom前置拦截非法ID查询误判率0.01%。4. 实操过程从0到1搭建WebSocket广播集群的完整步骤4.1 环境准备与基础组件部署我们生产环境用CentOS 7.9内核升级到5.10支持io_uring虽本次未用但为后续优化留余地。所有机器配置统一32核CPU / 64GB内存 / 1TB NVMe SSD用于Kafka日志和ES数据。第一步部署etcd集群3节点不是为了时髦而是需要强一致的分布式KV。用etcdctl初始化集群# node1 etcd --name infra0 --initial-advertise-peer-urls http://10.0.1.10:2380 \ --listen-peer-urls http://0.0.0.0:2380 \ --listen-client-urls http://0.0.0.0:2379 \ --advertise-client-urls http://10.0.1.10:2379 \ --initial-cluster-token etcd-cluster-1 \ --initial-cluster infra0http://10.0.1.10:2380,infra1http://10.0.1.11:2380,infra2http://10.0.1.12:2380 \ --initial-cluster-state new其他节点类似仅改--name和IP。启动后用etcdctl endpoint health确认健康。我们把房间路由表存在/ws/route/路径下TTL 30秒靠Router节点心跳续期。第二步Kafka集群3 broker 1 ZooKeeperZooKeeper只用于Kafka元数据不承担其他角色。Kafka配置关键项# server.properties num.partitions12 default.replication.factor2 min.insync.replicas2 log.retention.hours24 # 关键关闭自动创建topic所有topic手动创建并指定分区数 auto.create.topics.enablefalse弹幕Topiclive-chat创建命令kafka-topics.sh --create --bootstrap-server 10.0.1.20:9092 \ --replication-factor 2 --partitions 12 --topic live-chat为什么12分区因为峰值4.7万条/秒单分区极限1.2万12分区理论吞吐14.4万留30%余量。消费者组chat-broadcaster配12个实例一一对应分区。第三步Redis Cluster6节点3主3从用redis-cli --cluster create一键部署。重点配置# redis.conf maxmemory 32gb maxmemory-policy allkeys-lru # 关键禁用AOF用RDB快照避免AOF重写阻塞 appendonly no save 300 10000用户状态缓存用Hash不设过期时间靠应用层清理避免缓存雪崩。4.2 Router服务开发与部署Router是无状态服务用Go1.19开发轻量高效。核心逻辑就两个HTTP HandlerPOST /room/join接收{room_id:room_123,user_id:45678}生成路由记录存etcd返回{vip:10.0.2.100,token:eyJhb...}。GET /ping健康检查返回{status:ok,ts:1712345678}。部署时用systemd管理配置Restarton-failureRestartSec10。我们启5个Router实例前面挂LVS权重均等。Router自身不存状态所有数据都在etcd所以扩缩容就是起停进程秒级生效。Router的etcd写入逻辑防脑裂不是简单Put而是用CompareAndSwapCAS// 伪代码 cmp : clientv3.Compare(clientv3.Version(key), , 0) // 确保key不存在 put : clientv3.OpPut(key, value, clientv3.WithLease(leaseID)) txn : clientv3.Txn(ctx).If(cmp).Then(put) resp, _ : txn.Commit() if !resp.Succeeded { // key已存在说明其他Router已写入读取现有值返回 getResp, _ : clientv3.Get(ctx, key) return getResp.Kvs[0].Value }这样即使多个Router同时处理同一个房间创建请求也只有一个能成功写入避免路由冲突。4.3 Broadcast服务核心实现Netty RaftBroadcast服务是性能核心用Java17 Netty4.1.94开发。关键点Raft选主用raft-java库3节点组成集群。Leader节点负责接收本组所有弹幕来自Kafka Consumer维护一个环形缓冲区RingBuffer存最近10秒弹幕序列号long[]数组大小10000广播时遍历本地ChannelGroup对每个活跃Channel调用channel.writeAndFlush()传入预构建的ByteBufChannelGroup管理不用Netty自带的DefaultChannelGroup线程安全但性能差自研ConcurrentChannelGrouppublic class ConcurrentChannelGroup { private final MapString, Channel channels new ConcurrentHashMap(); private final ReadWriteLock lock new StampedLock(); public void add(Channel ch) { String key ch.attr(ATTR_USER_ID).get() _ ch.attr(ATTR_ROOM_ID).get(); channels.put(key, ch); } public void broadcast(ByteBuf msg) { // 用StampedLock读锁避免写操作阻塞广播 long stamp lock.tryOptimisticRead(); CollectionChannel c channels.values(); if (!lock.validate(stamp)) { stamp lock.readLock(); try { c channels.values(); } finally { lock.unlockRead(stamp); } } for (Channel ch : c) { if (ch.isActive()) ch.writeAndFlush(msg.retain()); // retain避免重复引用 } } }retain()是关键确保msg在所有Channel写完前不被释放。部署每个Broadcast组3台机器配置相同。用Ansible批量部署启动脚本里加-XX:UseZGC -Xmx32gZGC停顿时间稳定在10ms内。我们监控ChannelGroup.size()当单组连接数1.2万时Router自动把新用户导向下一组实现动态负载均衡。4.4 客户端SDK集成要点Vue3 TypeScript前端不是甩手掌柜。我们提供ws/live-sdknpm包核心是LiveSocket类class LiveSocket { private socket: WebSocket | null null; private reconnectTimer: NodeJS.Timeout | null null; private lastSeq: number 0; connect(roomId: string, userId: number) { // 第一步向Router获取VIP fetch(/router/join?room_id${roomId}user_id${userId}) .then(res res.json()) .then(data { this.socket new WebSocket(wss://${data.vip}/ws?token${data.token}); this.setupEventListeners(); }); } private setupEventListeners() { this.socket?.addEventListener(message, (e) { const msg JSON.parse(e.data); if (msg.type chat) { this.lastSeq msg.seq; // 记录最新序列号 this.emit(chat, msg); } }); this.socket?.addEventListener(close, (e) { if (e.code 1006) { // 网络断开 this.attemptReconnect(); // 触发重连 } }); } private attemptReconnect() { // 指数退避100ms, 300ms, 900ms const delays [100, 300, 900]; const delay delays[Math.min(this.retryCount, 2)]; this.reconnectTimer setTimeout(() { fetch(/router/reconnect?token${this.token}last_seq${this.lastSeq}) .then(res res.json()) .then(data { this.socket new WebSocket(wss://${data.vip}/ws?token${data.token}); this.setupEventListeners(); }); }, delay); } }关键经验WebSocket在iOS Safari中页面进入后台时连接会被系统静默关闭且onclose事件可能不触发。我们加了document.visibilityState监听页面切后台时主动socket.close()切前台时自动重连避免残留无效连接占用服务端资源。5. 常见问题与排查技巧实录那些线上凌晨三点的救火时刻5.1 典型问题速查表问题现象可能原因快速定位命令解决方案用户连接后收不到弹幕Router路由表未更新或Broadcast组未正确订阅Kafka Topicetcdctl get /ws/route/room_123kafka-consumer-groups.sh --bootstrap-server x.x.x.x:9092 --group chat-broadcaster --describe检查Router日志是否写入etcd成功确认Consumer Group的CURRENT-OFFSET与LOG-END-OFFSET差值100弹幕延迟突增到2秒以上Kafka积压或Broadcast节点GC停顿kafka-run-class.sh kafka.tools.GetOffsetShell --bootstrap-server x.x.x.x:9092 --topic live-chat --time -1jstat -gc pid 1s扩容Consumer实例若GC频繁检查-Xmx是否过大ZGC建议-Xmx32g避免大堆内存导致回收慢大量连接1006错误断开Nginx或LVS连接超时或客户端网络不稳定ss -s查看TIME-WAIT连接数tcpdump -i any port 443 -w ws.pcap抓包分析LVS设置net.ipv4.ip_vs.conn_reuse_mode0Nginx加proxy_read_timeout 300客户端加心跳保活新用户加入房间看不到历史弹幕Redis缓存未命中且MySQL查询慢redis-cli -h x.x.x.x hgetall user:12345slowlog get 5检查布隆过滤器是否生效优化MySQL用户表索引user_id必须是主键5.2 一次真实的线上故障复盘进球瞬间的“弹幕消失”时间某场中超比赛第89分钟主队进球。现象监控告警broadcast_group_07的channel_active_count从1.1万骤降至300大量用户反馈“弹幕发不出也看不到别人的”。排查过程先看Kafkalive-chatTopic积压达23万条Consumer Lag飙升。登录broadcast_group_07的Leader节点jstat -gc显示Full GC每2秒一次G1OldGen使用率99%。jmap -histo pid发现io.netty.buffer.PooledUnsafeDirectByteBuf对象占堆内存78%但DirectMemory使用正常说明是Netty内存池泄漏。根因我们发现有个新上线的“弹幕表情包”功能后端返回的二进制协议里表情URL字段未做长度校验。恶意用户构造了10MB的超长URLBroadcast节点解析时ByteBuf分配了10MB内存但后续处理失败release()没被调用内存池一直不回收。100个这样的弹幕就把32GB Direct Memory吃光触发ZGC Full GC。修复协议层加字段长度校验if (urlLength 2048) throw new ProtocolException(URL too long);内存池加监控PooledByteBufAllocator的metric()方法暴露numActiveAllocations当10万时告警。紧急回滚表情包功能22分钟恢复。实操心得永远不要相信上游数据。WebSocket协议解析的第一行必须是严格的长度校验和魔数校验。我们后来在Router层加了WAF规则对/ws路径的所有Upgrade请求做二进制头校验非法请求直接400拒绝不进后端。5.3 性能压测的黄金参数与避坑指南我们用k6做全链路压测脚本模拟真实用户行为import http from k6/http; import { check, sleep } from k6; export const options { stages: [ { duration: 30s, target: 1000 }, // ramp up { duration: 5m, target: 10000 }, // steady state { duration: 30s, target: 50000 }, // spike ], thresholds: { http_req_duration: [p(95)300], // HTTP请求95%300ms checks: [rate0.99] // 断言成功率99% } }; export default function () { // 1. 获取Router VIP const res1 http.post(https://router/api/join, JSON.stringify({ room_id: room_${__VU}, user_id: __VU * 1000 __ITER })); // 2. 建立WebSocketk6不原生支持WS用websockets库 const ws new WebSocket(wss://${res1.json().vip}/ws?token${res1.json().token}); // 3. 发送弹幕每5秒1条 ws.on(open, () { setInterval(() { ws.send(JSON.stringify({type:chat, text:GOAL!, user_id:__VU})); }, 5000); }); }压测避坑别用单机压测k6单机最多撑1万并发要用k6 cloud或分布式部署。我们用3台k6压测机每台起1.5万VU总4.5万逼近真实场景。监控必须全链路除了k6的http_req_duration还要在Router、Broadcast、Kafka各节点装Prometheus Grafana看etcd_request_duration_seconds、netty_channel_active_count、kafka_consumer_lag。单看k6指标可能掩盖后端积压。网络带宽是瓶颈4.5万并发每秒弹幕2000条每条弹幕平均200字节下行带宽需2000*200*8/1024≈3.1Mbps看似不大。但实际是10万连接每秒都在收带宽是100000*200*8/1024≈156Mbps。我们最初忘了配网卡多队列irqbalance没开所有中断集中在一个CPU核该核100%其他核空闲。加ethtool -L eth0 combined 32并重启irqbalance问题解决。5.4 成本与扩展性平衡10万人在线到底要多少机器很多人问“10万并发要买多少云服务器”答案不是固定数字而是看你的SLA。我们按生产环境算过一笔细账组件数量配置年成本参考关键作用Router5台4核8G¥12,000无状态纯路由可水平扩展Broadcast组4组 × 3台 12台32核64G¥288,000核心广播每组撑2.5万连接Kafka Broker3台16核32G¥72,000弹幕缓冲峰值吞吐保障etcd3台4核8G¥18,000路由元数据强一致Redis Cluster6台8核16G¥108,000用户状态缓存降低DB压力总计30台—¥498,000—注意这30台是峰值保障配置。日常流量只有峰值的30%我们用K8s HPA基于CPU和channel_active_count指标自动缩容Broadcast组到2组6台Router缩到3台年成本可降40%。但体育赛事前2小时必须手动扩到满配因为HPA扩容需要3~5分钟而进球可能发生在下一秒。最后分享一个小技巧我们把Broadcast组的机器名按ws-group-01-a、ws-group-01-b、ws-group-01-c命名-a是Leader候选-b、-c是Follower。Raft选举时优先选-a这样Leader位置固定运维排查时一眼就知道主节点在哪台不用每次curl http://x.x.x.x:8080/raft/status去查。我在实际压测中发现当单Broadcast组连接数超过1.3万时ChannelGroup.broadcast()的延迟开始非线性增长从15ms跳到40ms。所以我们的硬性红线是1.2万/组超过就触发Router自动分流。这个数字是我们在32核机器上用wrk反复测试/ws/broadcast接口得出的不是拍脑袋。技术没有银弹只有一次次实测出来的边界。
返回列表