千万级全服弹幕群发必须用Redis Stream+消费者组,因map[*websocket.Conn]bool+RWMutex会导致内存超20GB、广播延迟毛刺超1.2s、write on closed connection panic;Redis Stream提供持久化、ACK、断点续推和横向扩容能力,是唯一可扩展底座。

千万级全服弹幕群发不能靠单机内存遍历连接广播,必须用 Redis Stream + 消费者组做消息分发中枢,配合 gorilla/websocket 连接池与心跳保活;否则 goroutine 泄露、写 panic、OOM、消息重复或丢失会集中爆发。
为什么不能用 map[*websocket.Conn]bool + RWMutex 全量广播
全服场景下,连接数轻松破千万,内存中维护每个 *websocket.Conn 指针会导致:① 单机内存占用超 20GB(按 2KB/goroutine 栈 + 指针 + 锁结构粗略估算);② RWMutex.RLock() 遍历时 GC 扫描压力陡增,实测 500 万连接下广播延迟毛刺超 1.2s;③ 连接断开时 delete 不及时,残留指针触发 write on closed connection panic。真实线上已淘汰该模式。
Redis Stream 是唯一可扩展的全服弹幕分发底座
Stream 天然支持追加、消费者组、pending 列表、ACK 机制和断点续推,是千万级弹幕的事实标准。关键约束有三点:
- 每条弹幕必须写入全局 stream,如
stream:all:messages,不能按房间分 key —— 全服广播意味着所有消费者都从同一源头读 - 必须用消费者组(
XGROUP CREATE)而非 Pub/Sub:Pub/Sub 无持久化、无 ACK、无法横向扩容消费者,一节点宕机即丢全量弹幕 - 消费者组名需带业务标识,如
consumer_group_all_v2,避免多服务实例误消费彼此数据
go-redis/v9 下 Stream 写入必须绕过 clusterClient.Pipeline()
在 Redis Cluster 环境中,clusterClient.Pipeline() 直接 panic,必须手动路由:
立即学习“go语言免费学习笔记(深入)”;
- 先调
clusterClient.KeySlot("stream:all:messages")获取 slot - 再用
clusterClient.NodeForKey(slot)拿到对应节点的*redis.Client - 在该节点 client 上调
node.XAdd(ctx, &redis.XAddArgs{...})—— 不要 pipeline,XAdd 本身是原子写,批量写反而增加失败面 - key 必须不含 hash tag(如不写
stream:{all}:messages),否则会被强制路由到固定 slot,失去集群负载均衡能力
消费者端必须用独立 goroutine + context 控制生命周期
每个消费者组实例需起一个长运行 goroutine 调用 XReadGroup,且必须满足:
- 每次
XReadGroup设置Count: 100和Block: 5000(5s),防空轮询耗 CPU - 必须用
ctx, cancel := context.WithTimeout(parentCtx, 30*time.Second)包裹整个处理循环,超时强制重连,避免卡死 - 收到消息后,解析出弹幕内容,再通过 channel 投递给 websocket 广播协程 —— 绝不能在
XReadGroup循环里直接调conn.WriteMessage(),否则一个慢连接拖垮整组消费 - ACK 必须在消息成功投递到广播队列后才调
XAck(),否则网络抖动导致重复消费
最容易被忽略的是:Stream 的 MAXLEN ~ 参数必须显式设置(如 MAXLEN 10000000),否则全服弹幕无限堆积,磁盘爆满后 Redis 直接只读。别信“自动清理”——它只在写入时触发,而全服弹幕写入频次可能低于清理阈值。


















