从 MVP 到百万并发:直播弹幕系统完整设计
本文参考并延伸了王帅真《直播弹幕系统设计》中的演进思路,在其“轮询 MVP → 长连接推送 → MQ 削峰 → 热门直播间优化”的基础上,补充容量估算、协议设计、房间路由、顺序一致性、可靠性、慢消费者治理、内容安全、容灾和运维方案。
本文不是对某一家直播平台内部实现的复刻,而是一套可以用于系统设计面试、技术方案评审和真实项目落地的通用架构。
一、问题定义:直播弹幕系统到底在解决什么
从表面看,直播弹幕只是用户发送一段文本,然后让直播间里的其他用户看到。
但当一个直播间从几十人增长到几十万人时,问题就不再是“保存一条消息”,而是:
- 如何维持海量客户端长连接;
- 如何让一条弹幕尽快到达同一房间内的用户;
- 如何避免一条消息被直接复制几十万次,形成巨大的写扩散;
- 如何处理热门直播间的瞬时流量;
- 如何保证同一直播间内消息大体有序;
- 如何在断线重连后补齐少量遗漏消息;
- 如何过滤违规内容、广告、刷屏和恶意请求;
- 如何在系统过载时优先保住直播,而不是让弹幕拖垮整个平台;
- 是否需要保存历史弹幕,以及如何支持直播回放;
- 如何控制带宽、CPU、内存和存储成本。
直播弹幕系统与 IM 系统很像,本质都是在一个逻辑空间内收发消息。但二者的业务优先级并不完全相同:
| 对比项 | 即时通讯 | 直播弹幕 |
|---|---|---|
| 消息对象 | 单聊、群聊 | 直播间 |
| 消息价值 | 通常需要长期保留 | 普通弹幕时效性强,过期价值快速下降 |
| 到达要求 | 一般要求可靠到达 | 普通弹幕可允许少量丢失,礼物、付费消息不能随意丢 |
| 顺序要求 | 会话内顺序较重要 | 通常只要求房间内近似有序或局部有序 |
| 在线关系 | 联系人、群成员相对稳定 | 用户随时进入、退出直播间 |
| 流量形态 | 相对均匀 | 明显热点,容易在赛事、抽奖、带货节点瞬间爆发 |
| 扇出规模 | 一个群数十到数万人 | 热门房间可能同时在线数十万甚至更多 |
| 展示能力 | 客户端通常全部展示 | 屏幕每秒只能展示有限数量,天然需要采样和降噪 |
因此,直播弹幕系统的核心矛盾不是“有没有能力接收所有消息”,而是:
如何在有限的连接、计算和带宽资源下,把最有价值的消息,以足够低的延迟,稳定地送到尽可能多的用户面前。
二、先定义目标,不要一上来堆中间件
系统设计最容易犯的错误,是没有规模假设就直接画 Kafka、Redis、WebSocket 和几十个微服务。
下面以一个中大型直播平台为例,定义一组用于设计的假设。实际项目应使用真实业务数据替换。
2.1 规模假设
| 指标 | 假设值 |
|---|---|
| 日活用户 | 2000 万 |
| 峰值在线用户 | 200 万 |
| 同时开播房间 | 5 万 |
| 单个热门房间峰值在线 | 30 万 |
| 全站峰值弹幕发送量 | 10 万条/秒 |
| 热门房间峰值发送量 | 5000 条/秒 |
| 单条弹幕编码后平均大小 | 约 220 字节 |
| 普通用户端实际展示速率 | 20~50 条/秒 |
| WebSocket 推送 P99 延迟目标 | 小于 500 毫秒 |
| 发送接口 P99 响应时间 | 小于 100 毫秒 |
| 服务可用性目标 | 99.99% |
这些数值不是行业标准,只是容量设计的起点。
2.2 入口流量并不可怕,广播流量才可怕
全站每秒写入 10 万条消息,假设每条 220 字节:
1 | |
22 MB/s 的入口流量对一个经过扩容的消息集群并不夸张。
真正危险的是热门房间的广播。假设一个房间每秒产生 5000 条弹幕,在线用户 30 万,每条消息 220 字节,如果对每个用户逐条复制:
1 | |
这还没有计算 TCP、TLS、WebSocket 帧和公网传输开销。
因此,大型弹幕系统绝不能简单实现成:
1 | |
这段代码在功能测试里可能很优雅,在生产环境里则像一个精心编写的自毁按钮。
2.3 系统目标分级
不同类型的消息,可靠性和优先级不应相同。
| 消息类型 | 示例 | 可靠性 | 优先级 |
|---|---|---|---|
| 普通弹幕 | “主播好”“666” | 允许极少量丢失 | 低 |
| 系统通知 | 开播、封禁、房间状态变化 | 应可靠到达 | 高 |
| 礼物消息 | 用户赠送礼物 | 必须可追踪、可补偿 | 高 |
| 付费弹幕 | 醒目留言、超级留言 | 不应丢失 | 最高 |
| 互动事件 | 点赞聚合、投票变化 | 可合并、可采样 | 中 |
| 风控指令 | 禁言、踢出、关闭房间 | 必须及时执行 | 最高 |
这意味着系统不应只有一个统一的“消息通道”。至少要在协议、队列或调度层面区分优先级。
三、核心业务流程
一条弹幕从用户 A 到达同房间用户 B,至少经过以下阶段:
- 客户端建立连接并完成鉴权;
- 客户端加入直播间;
- 用户发送弹幕;
- 接入层校验参数、身份、频率和房间状态;
- 内容安全系统进行快速审核;
- 消息写入可靠消息通道;
- 房间处理器分配顺序标识;
- 路由系统找到当前订阅该房间的网关节点;
- 网关在本机对连接进行广播;
- 客户端去重、排序、限速并渲染;
- 历史写入服务异步落库;
- 实时分析系统消费消息,更新热度和运营指标。
flowchart LR
A[发送方客户端] --> B[接入层 / API Gateway]
B --> C[鉴权与限流]
C --> D[内容安全快速审核]
D --> E[消息队列]
E --> F[房间消息处理器]
F --> G[房间路由 / Fanout Router]
G --> H1[连接网关 1]
G --> H2[连接网关 2]
G --> H3[连接网关 N]
H1 --> I1[本机连接集合]
H2 --> I2[本机连接集合]
H3 --> I3[本机连接集合]
E --> J[历史存储消费者]
E --> K[实时计算与监控]
这里最重要的优化是:
消息不是从中心服务直接复制给每一个用户,而是先按“订阅该房间的网关节点”发送一份,再由网关在本机完成局部广播。
假设一个热门房间的 30 万用户分布在 30 台连接网关上,中心路由层每条消息只需要发送约 30 份,而不是 30 万份。
四、消息数据模型
直播弹幕看似只有 roomId、userId 和 content,但真正上线后会迅速出现去重、顺序、重试、审核、样式和回放问题。
推荐的数据结构如下:
1 | |
4.1 字段说明
| 字段 | 作用 |
|---|---|
messageId |
服务端生成的全局唯一消息 ID |
clientMsgId |
客户端生成,用于请求重试时幂等去重 |
roomId |
直播间标识,也是消息分区和路由的重要依据 |
messageType |
普通弹幕、系统消息、礼物、禁言指令等 |
serverTime |
服务端接收或确认时间,避免依赖不可信的客户端时间 |
roomSeq |
房间内顺序号,用于排序、断线补偿和缺口检测 |
priority |
推送和丢弃策略的依据 |
moderation.status |
内容审核状态 |
extension |
可扩展业务属性,避免频繁修改主协议 |
traceId |
链路追踪标识 |
4.2 直播是否需要视频时间偏移
纯实时直播中,弹幕通常按服务器时间到达即可,并不一定需要点播视频那样的固定偏移量。
但以下场景仍然需要视频时间轴字段:
- 直播回放;
- DVR 时移播放;
- 用户主动拖动直播进度;
- 主播端和观众端存在不同播放缓冲;
- 多 CDN 节点造成播放延迟差异;
- 需要把弹幕重新挂载到录播文件。
因此可以增加可选字段:
1 | |
其中 videoPts 表示消息对应的媒体时间位置,streamEpoch 用于区分断流重开后的不同直播时间轴。
五、第一阶段:轮询模式的 MVP
早期业务规模较小时,不需要立刻建设复杂的长连接集群。
最小可行版本可以采用:
- HTTP 短连接;
- 客户端每 1~2 秒轮询一次;
- Redis 保存最近一段时间的弹幕;
- 数据库异步保存需要回放或审计的消息;
- 服务端返回自上次游标之后的增量消息。
flowchart LR
A[客户端发送弹幕] --> B[弹幕 API]
B --> C[校验/限流/审核]
C --> D[(Redis 最近消息)]
C --> E[消息队列]
E --> F[(历史数据库)]
G[客户端定时轮询] --> H[查询 API]
H --> I[本地缓存]
I -->|未命中| D
H --> G
5.1 使用 Redis ZSet 保存最近消息
可以按直播间建立 ZSet:
1 | |
写入:
1 | |
查询:
1 | |
Redis Sorted Set 允许不同成员使用相同分值,但成员本身必须唯一。同分值成员会进一步按字典序排列。因此,不能只把弹幕正文作为 member,否则相同内容会互相覆盖;应把 messageId 或 roomSeq 编入 member。
5.2 时间戳游标的缺陷
如果客户端只记录 lastTimestamp,同一毫秒内存在多条消息时容易出现:
- 重复拉取;
- 边界消息遗漏;
- 同分值顺序不稳定;
- 客户端重试时难以精确恢复。
更稳妥的方式是返回复合游标:
1 | |
服务端按照 (serverTime, roomSeq) 做稳定排序。
如果系统已经具备连续房间序号,也可以直接使用 roomSeq 作为增量游标。
5.3 清理旧消息
Redis 只保留最近 30 秒或最近 5000 条消息:
1 | |
或者:
1 | |
清理动作可以由写入 Lua 脚本、定时任务或后台消费者完成,避免每次读请求都执行昂贵清理。
5.4 本地缓存降低重复读取
同一直播间的客户端会在相近时间请求几乎相同的数据。如果每次都回源 Redis,会造成大量重复读取。
查询服务可以缓存最近 1~5 秒的房间消息:
1 | |
优化原则:
- 只缓存活跃房间;
- 设置最大房间数量和最大内存;
- 使用 Caffeine 等带权重淘汰的本地缓存;
- 不要因为房间多就无限创建缓存项;
- 对热门房间采用请求合并,避免同一时刻大量回源;
- 使用一致性哈希或房间路由,使同一房间的读请求尽量命中相同节点。
5.5 轮询模式的问题
轮询方案实现简单,但规模上升后会暴露明显问题:
- 大量请求没有新消息,形成空轮询;
- HTTP 头和连接管理开销较高;
- 实时性受轮询周期限制;
- 客户端越多,Redis 重复查询越严重;
- 热门直播间的瞬时消息仍会冲击单个 Redis Key;
- 客户端轮询时间不一致,用户看到的弹幕不同步;
- 服务端难以主动下发封禁、房间关闭等控制指令。
因此,轮询适合作为早期 MVP 和后续降级通道,而不适合作为大规模系统的唯一方案。
六、传输协议选型:轮询、SSE 还是 WebSocket
| 方案 | 通信方向 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 短轮询 | 客户端主动请求 | 简单、兼容性好 | 空请求多、延迟高 | MVP、降级 |
| 长轮询 | 服务器延迟响应 | 比短轮询实时 | 每次消息后仍需重新发起请求 | 中小规模 |
| SSE | 服务端单向推送 | 基于 HTTP、自动重连、协议简单 | 浏览器侧主要是单向通道 | 只需要服务端推送 |
| WebSocket | 全双工 | 双向、长连接、帧开销低 | 连接治理和运维复杂 | 直播、IM、协同编辑 |
| WebTransport | 多流、基于 QUIC | 更灵活,弱网潜力更好 | 生态和基础设施成熟度要求高 | 特定新项目 |
直播弹幕既要客户端发送消息,也要服务端主动推送消息,因此生产系统通常选择 WebSocket。
WebSocket 在 RFC 6455 中定义为基于 TCP 的双向消息协议,典型方式是先通过 HTTP/1.1 Upgrade 完成握手。RFC 8441 又定义了在 HTTP/2 单个流上引导 WebSocket 的机制。因此,不能简单把 WebSocket 等同于“HTTP/2 自带的全双工通信”。
七、生产架构:WebSocket 推送模式
随着并发用户增长,系统需要将“业务处理”和“连接维护”拆开。
7.1 总体架构
flowchart TB
subgraph Client[客户端]
C1[Web / App / TV]
end
subgraph Edge[边缘接入层]
LB[四层/七层负载均衡]
GW1[Connection Gateway 1]
GW2[Connection Gateway 2]
GWN[Connection Gateway N]
end
subgraph Core[弹幕核心服务]
API[Danmaku Ingress]
AUTH[鉴权服务]
LIMIT[限流与风控]
SAFE[内容安全]
MQ[(消息队列)]
RP[Room Processor]
FR[Fanout Router]
ROUTE[(Room-Gateway 路由表)]
end
subgraph Data[数据层]
REDIS[(Redis 热数据)]
DB[(历史消息数据库)]
OBJ[(对象存储/离线数仓)]
end
C1 --> LB
LB --> GW1
LB --> GW2
LB --> GWN
GW1 --> API
GW2 --> API
GWN --> API
API --> AUTH
API --> LIMIT
API --> SAFE
SAFE --> MQ
MQ --> RP
RP --> FR
FR --> ROUTE
FR --> GW1
FR --> GW2
FR --> GWN
MQ --> DB
DB --> OBJ
RP --> REDIS
7.2 组件职责
Connection Gateway
只负责连接和协议,不承担复杂直播业务:
- 建立和关闭 WebSocket;
- 心跳检测;
- 解析二进制或 JSON 协议;
- 维护本机连接表;
- 维护本机
roomId -> connections索引; - 将用户上行消息转发给弹幕接入服务;
- 接收 Fanout Router 的房间消息;
- 在本机批量广播;
- 管理慢消费者和发送缓冲区;
- 连接优雅迁移与关闭。
它是有连接状态的服务,但业务状态应尽量外置。节点故障后,客户端可以重连到任意其他节点并恢复房间订阅。
Danmaku Ingress
负责写入链路:
- 验证用户身份和直播间状态;
- 校验消息长度、类型和协议版本;
- 用户级、设备级、IP 级和房间级限流;
- 幂等校验;
- 快速内容审核;
- 生成服务端消息 ID;
- 将消息写入消息队列;
- 返回发送确认。
Room Processor
负责房间维度的消息处理:
- 按房间串行消费;
- 分配
roomSeq; - 应用房间规则;
- 处理禁言、管理员、付费优先级;
- 聚合点赞等高频事件;
- 生成适合下游推送的统一事件。
Fanout Router
负责把一条房间消息发送到真正有该房间观众的网关节点:
- 查询
roomId -> gatewayId 集合; - 每个目标网关只发送一份;
- 批量发送多条消息;
- 对热门房间进行分层广播;
- 根据节点负载实施降级和丢弃策略;
- 收集推送结果和拥塞信息。
History Writer
异步保存需要回放、审计和分析的消息,不参与普通弹幕的主推送路径,避免数据库抖动直接影响实时链路。
八、WebSocket 连接生命周期
8.1 建立连接
浏览器原生 WebSocket API 不方便自定义任意 Authorization 请求头。可以采用以下方式:
- 客户端先通过 HTTPS 使用登录凭证换取一个短期、一次性的 WebSocket Ticket;
- 客户端使用 Ticket 建立连接;
- 网关校验 Ticket,绑定用户、设备和会话;
- Ticket 使用后立即失效。
1 | |
1 | |
不要把长期 Access Token 放在 URL 中,因为 URL 可能出现在代理日志、浏览器历史或监控系统中。
8.2 连接时序
sequenceDiagram
participant C as Client
participant A as Auth API
participant L as Load Balancer
participant G as Connection Gateway
participant R as Route Registry
C->>A: 申请一次性 WebSocket Ticket
A-->>C: Ticket,30 秒过期
C->>L: WebSocket Handshake + Ticket
L->>G: 转发握手
G->>A: 校验 Ticket
A-->>G: userId / deviceId / permissions
G-->>C: CONNECTED(sessionId, heartbeatInterval)
C->>G: JOIN_ROOM(roomId, lastSeq)
G->>R: 首个本地用户加入时注册 roomId -> gatewayId
G-->>C: JOINED(roomId, currentSeq)
8.3 负载均衡是否必须使用粘性会话
WebSocket 连接建立后,TCP 连接本身已经固定在某台网关上,因此连接存续期间不需要再做请求级粘性路由。
但需要注意:
- 重连可能进入另一台网关;
- 网关不能把关键业务状态只放在本机;
- 订阅关系需要重新恢复;
- 客户端应保存最近房间和
lastSeq; - 负载均衡器、Ingress 和反向代理的空闲超时必须大于心跳间隔;
- 滚动发布时需要先停止接收新连接,再排空旧连接。
8.4 心跳
推荐使用协议级 Ping/Pong 或应用级心跳:
1 | |
服务端响应:
1 | |
心跳间隔不能过短,否则百万连接会产生大量无意义请求。例如 200 万连接每 5 秒一次心跳,相当于每秒 40 万次心跳。可以设置为 20~60 秒,并配合 TCP Keepalive、移动端前后台策略和随机抖动。
九、协议设计
9.1 帧类型
| 帧类型 | 方向 | 作用 |
|---|---|---|
CONNECTED |
服务端 → 客户端 | 连接建立成功 |
JOIN_ROOM |
客户端 → 服务端 | 加入房间 |
LEAVE_ROOM |
客户端 → 服务端 | 离开房间 |
SEND_DANMAKU |
客户端 → 服务端 | 发送弹幕 |
SEND_ACK |
服务端 → 客户端 | 确认发送结果 |
ROOM_MESSAGE |
服务端 → 客户端 | 房间消息 |
BATCH_MESSAGES |
服务端 → 客户端 | 批量房间消息 |
RESUME |
客户端 → 服务端 | 断线恢复 |
PING/PONG |
双向 | 心跳 |
KICK_OUT |
服务端 → 客户端 | 强制下线 |
FLOW_CONTROL |
服务端 → 客户端 | 通知客户端降速 |
ERROR |
服务端 → 客户端 | 协议或业务错误 |
9.2 发送弹幕
客户端发送:
1 | |
服务端确认:
1 | |
客户端可以先本地乐观展示,再根据 ACK 将临时消息替换为正式消息。如果服务端拒绝,则标记发送失败,而不是悄悄消失。
9.3 批量推送
逐条发送 WebSocket 帧会增加系统调用、对象创建和协议开销。网关可以按 10~50 毫秒窗口聚合:
1 | |
批量窗口过大,会增加延迟;过小,则节省不了多少开销。需要结合消息速率动态调整。
9.4 JSON 还是 Protobuf
| 方案 | 优点 | 缺点 |
|---|---|---|
| JSON | 易调试、前端接入简单 | 体积较大、解析成本较高 |
| Protobuf | 体积小、解析快、Schema 明确 | 调试和版本治理更复杂 |
| MessagePack | 比 JSON 紧凑 | 生态和 Schema 约束弱于 Protobuf |
早期可以使用 JSON,规模上升后切换 Protobuf,并保留协议版本:
1 | |
十、消息写入链路
10.1 完整写入时序
sequenceDiagram
participant C as Client
participant G as Connection Gateway
participant I as Danmaku Ingress
participant S as Safety Service
participant M as MQ
participant P as Room Processor
participant F as Fanout Router
participant T as Target Gateways
participant D as History Writer
C->>G: SEND_DANMAKU(clientMsgId)
G->>I: 转发消息
I->>I: 鉴权、参数校验、限流、幂等
I->>S: 快速内容审核
S-->>I: PASS / REJECT / REVIEW
I->>M: 按 roomId 写入
M-->>I: Broker ACK
I-->>G: SEND_ACK(ACCEPTED)
G-->>C: 发送确认
M->>P: 顺序消费房间事件
P->>P: 分配 roomSeq
P->>F: 发布可推送事件
F->>T: 每个有订阅者的网关发送一份
T-->>C: ROOM_MESSAGE / BATCH_MESSAGES
M->>D: 异步落历史存储
10.2 ACK 到底代表什么
必须明确 ACK 语义:
| ACK 状态 | 含义 |
|---|---|
RECEIVED |
网关收到,但尚未进入可靠队列 |
ACCEPTED |
已通过基础校验并写入可靠消息队列 |
PUBLISHED |
已进入房间广播链路 |
REJECTED |
被审核、限流或业务规则拒绝 |
REVIEWING |
进入人工或异步审核,不立即广播 |
普通弹幕通常在 MQ 确认写入后返回 ACCEPTED 即可,不必等待所有观众收到。否则一个慢客户端就可能拖慢发送方响应。
10.3 为什么需要 clientMsgId
客户端超时后可能重试。如果没有幂等键,同一条弹幕会重复发送。
服务端可以在短时间窗口内记录:
1 | |
重复请求直接返回原 ACK。
幂等记录可以放在 Redis,并设置几分钟过期。对于付费消息,还应使用数据库唯一约束或业务流水号,不能只依赖短期缓存。
十一、消息队列与房间顺序
11.1 按 roomId 分区
消息队列的分区键应使用 roomId:
1 | |
这样同一房间的消息进入同一个分区,并由一个房间处理器顺序消费。
Kafka 官方文档明确说明:同一个 Key 的事件会进入同一分区,消费者读取某个 Topic-Partition 时,会按照写入顺序读取事件。
flowchart LR
A1[Room A 消息] --> P1[Partition 1]
A2[Room A 消息] --> P1
B1[Room B 消息] --> P2[Partition 2]
C1[Room C 消息] --> P3[Partition 3]
P1 --> C01[Consumer 1]
P2 --> C02[Consumer 2]
P3 --> C03[Consumer 3]
11.2 全局顺序没有必要
直播弹幕不需要全站统一顺序,只需要同一房间内具备可解释的顺序。
强行追求全局顺序会导致:
- 单分区瓶颈;
- 跨机房协调;
- 吞吐下降;
- 故障恢复复杂;
- 系统成本显著增加。
11.3 roomSeq 如何生成
可选方案:
方案一:Redis INCR
1 | |
优点是简单;缺点是热门房间会形成单 Key 热点,Redis 故障和迁移也需要处理序号连续性。
方案二:房间处理器内存计数
同一房间始终由同一个顺序消费者处理,在内存中递增,并周期性持久化检查点。
优点:
- 无需每条消息访问 Redis;
- 延迟低;
- 与分区消费模型自然结合。
缺点:
- 消费者切换时需要恢复;
- 检查点落后可能产生重复或跳号;
- 需要使用 Epoch 区分不同处理实例。
可使用:
1 | |
更准确地说,可以把 epoch 放在高位,把本地计数器放在低位,避免旧消费者恢复后与新消费者产生序号冲突。
方案三:直接使用消息队列 Offset
队列 Offset 在分区内单调递增,但一个分区通常承载多个房间,因此对某个房间来说会出现大量跳号。
它可以用于排序,但不适合直接作为“连续房间序号”进行缺口判断。
11.4 是否需要 Exactly Once
普通弹幕没有必要追求端到端 Exactly Once。
更现实的目标是:
- MQ 至少一次投递;
- 消息处理器允许重复消费;
- 使用
messageId幂等; - 网关和客户端去重;
- 顺序以
roomSeq为准; - 重要消息使用业务流水和持久化补偿。
Exactly Once 往往只是把复杂性转移到了事务、状态和边界定义中。对于普通弹幕,少量重复比大幅降低吞吐更容易处理。
十二、房间与网关路由
12.1 不要在中心存储每个用户的连接
一种直觉设计是:
1 | |
当房间有几十万用户时,每条消息都遍历用户列表,中心路由系统会成为瓶颈。
更好的模型是:
1 | |
全局只记录“哪些网关有这个房间的用户”,具体用户连接由网关本地维护。
flowchart TB
R[Room 9527]
REG[(全局路由表)]
G1[Gateway A<br/>本地 12000 用户]
G2[Gateway B<br/>本地 9500 用户]
G3[Gateway C<br/>本地 11300 用户]
R --> REG
REG --> G1
REG --> G2
REG --> G3
G1 --> U1[本机连接集合]
G2 --> U2[本机连接集合]
G3 --> U3[本机连接集合]
12.2 注册规则
- 某网关第一个本地用户加入房间时,注册
(roomId, gatewayId); - 后续用户加入只修改本地计数;
- 最后一个本地用户离开时,注销;
- 网关定期发送心跳;
- 路由记录设置 TTL;
- 网关宕机后,由 TTL 自动清除脏路由;
- Fanout Router 缓存热点房间路由,但必须支持版本或短 TTL 更新。
路由记录示例:
1 | |
12.3 路由存储选型
可选方案包括:
- Redis Set / Hash;
- etcd;
- 服务注册中心;
- 专用路由服务;
- 网关主动订阅 Fanout Router;
- 基于 Gossip 的节点状态同步。
需要注意:etcd 更适合低频配置和服务发现,不适合每个用户进出房间都直接写。通过“首个用户注册、最后一个用户注销”可以把写入频率从用户级降低到网关级。
十三、避免广播风暴:两级与三级扇出
13.1 两级扇出
1 | |
中心层按网关发送,网关按本地连接发送。
13.2 三级扇出
跨地域或超热门直播间可以增加区域中继:
1 | |
flowchart TB
F[Global Fanout Router]
R1[Singapore Relay]
R2[Tokyo Relay]
R3[Frankfurt Relay]
F --> R1
F --> R2
F --> R3
R1 --> G11[SG Gateway 1]
R1 --> G12[SG Gateway 2]
R2 --> G21[Tokyo Gateway 1]
R2 --> G22[Tokyo Gateway 2]
R3 --> G31[Frankfurt Gateway 1]
G11 --> C11[Clients]
G12 --> C12[Clients]
G21 --> C21[Clients]
G22 --> C22[Clients]
G31 --> C31[Clients]
这样跨地域只发送一份消息,由区域内完成二次扩散,减少跨地域带宽。
13.3 为什么不能让所有网关直接消费同一个 Kafka Topic
如果所有网关属于同一个 Consumer Group,一条消息只会被其中一个网关消费,不满足广播需求。
如果每个网关使用独立 Consumer Group,那么每个网关都会消费全站消息:
- 消费组数量巨大;
- 网关收到大量自己没有订阅者的房间消息;
- Broker 出口带宽快速增长;
- 网关需要过滤几乎全部无关消息;
- 节点扩缩容会带来频繁重平衡。
因此,更合理的方式是:
- Kafka 负责可靠顺序日志;
- Room Processor / Fanout Router 消费一次;
- 路由层只把消息发给真正订阅了该房间的网关。
十四、热门直播间设计
普通直播间和热门直播间不应完全使用同一套参数。
14.1 热门房间识别
可以综合以下指标:
- 当前在线人数;
- 最近 1 分钟进入速率;
- 最近 10 秒弹幕写入 QPS;
- 推送带宽;
- 房间增长率;
- 是否为赛事、发布会、抽奖或平台活动;
- 主播等级和历史峰值;
- 网关分布数量;
- 消息积压和客户端丢弃率。
不要只使用粉丝数判断。一个粉丝很多但未开播的主播并不是热点;一个临时赛事房间可能在几十秒内从普通房间变成超级热点。
14.2 热门房间分级
| 等级 | 在线人数示例 | 策略 |
|---|---|---|
| Normal | 小于 1 万 | 公共集群 |
| Hot | 1 万~10 万 | 提高缓存、独立限流、批量推送 |
| Super Hot | 10 万~50 万 | 专用分区、区域中继、动态采样 |
| Mega Event | 50 万以上 | 独立集群、预热、专项容量保障 |
14.3 房间消息分片
单个房间固定进入一个 MQ 分区可以保证顺序,但超级热门房间可能把该分区打满。
有三种处理思路:
方案一:继续单分区,但为热点房间独占分区
适合消息写入量仍可由一个分区承载的情况。实现最简单,顺序最好。
方案二:放松严格顺序,拆分多个子分区
1 | |
客户端按到达顺序展示,或在 50~100 毫秒的小窗口内按服务端时间重排。
弹幕对绝对顺序并不敏感,通常比强行保持全序更划算。
方案三:先统一排序,再并行广播
使用一个轻量 Sequencer 先分配 roomSeq,随后按消息 ID 拆到多个广播分区。网关收到后进行小窗口合并。
flowchart LR
I[热门房间入口] --> S[Room Sequencer]
S --> P1[Broadcast Lane 1]
S --> P2[Broadcast Lane 2]
S --> P3[Broadcast Lane N]
P1 --> M[Gateway Merge Buffer]
P2 --> M
P3 --> M
M --> C[Clients]
该方案吞吐高,但实现复杂,需要处理:
- 某条消息延迟导致重排等待;
- 缺口超时;
- 重复消息;
- Lane 扩缩容;
- Epoch 切换;
- 客户端最终显示顺序。
14.4 自适应采样
热门房间每秒产生 5000 条弹幕,但客户端屏幕可能只展示 30 条。把 5000 条全部推到客户端既浪费带宽,也会导致客户端主线程和渲染卡顿。
服务端可以做:
- 按用户等级加权;
- 付费和礼物消息优先;
- 同文案聚合为“× 128”;
- 重复字符压缩;
- 低质量消息随机采样;
- 按区域或用户分组采样;
- 按客户端性能下发不同速率;
- 对后台播放、弱网用户降低速率;
- 对系统消息保留独立高优先级通道。
采样不一定要求所有用户看到完全相同的普通弹幕。直播产品更关心整体氛围,而不是普通弹幕的强一致广播。
十五、流控、背压与慢消费者
大型长连接系统一定会遇到慢客户端:
- 用户网络弱;
- 手机进入后台;
- 客户端卡顿;
- TCP 发送窗口变小;
- 网关发送队列持续增长;
- 某些连接长时间无法写出。
如果网关无限缓存,最终会耗尽内存。
15.1 每连接发送队列
每个连接设置:
- 最大消息数量;
- 最大字节数;
- 最大等待时间;
- 不同优先级队列;
- 丢弃和断开策略。
示例:
1 | |
15.2 丢弃策略
当连接拥塞时:
- 优先丢弃过期普通弹幕;
- 丢弃重复、低权重弹幕;
- 合并点赞和同文案消息;
- 降低普通弹幕推送频率;
- 保留系统、礼物和付费消息;
- 队列持续增长时断开连接,让客户端指数退避后重连。
15.3 背压传播
flowchart RL
C[慢客户端] --> G[Gateway 发送队列上涨]
G --> R[区域 Relay]
R --> F[Fanout Router]
F --> P[Room Processor]
G -.拥塞指标.-> R
R -.区域负载.-> F
F -.动态采样/降速.-> P
网关应把以下指标反馈给上游:
- 当前连接数;
- 发送队列 P95/P99;
- 每秒发送字节;
- 丢弃消息数;
- 慢连接比例;
- Event Loop 延迟;
- Direct Memory 使用量。
上游根据这些指标调整采样率,而不是等到节点 OOM 后再处理。
十六、可靠性与断线重连
16.1 普通弹幕是否可以丢
从产品体验看,普通弹幕允许极少量丢失,但系统不能以此为借口完全不做可靠性。
建议分层:
| 等级 | 策略 |
|---|---|
| 普通弹幕 | MQ 持久化、至少一次、客户端去重,可在过载时采样 |
| 系统消息 | 重试、补拉、优先队列 |
| 礼物/付费消息 | 业务数据库流水、Outbox、幂等、补偿、审计 |
| 风控命令 | 高优先级、ACK、超时重试、状态兜底 |
16.2 断线恢复
客户端保存:
1 | |
重连后发送:
1 | |
服务端处理:
- 检查最近消息缓存是否仍包含
lastSeq之后的数据; - 如果缺口较小,返回增量消息;
- 如果缺口太大,只返回最新窗口;
- 重要状态通过房间快照重新同步;
- 客户端按
messageId去重; - 如果
roomSeq发生 Epoch 变化,客户端清空旧重排窗口。
sequenceDiagram
participant C as Client
participant G as Gateway
participant H as Recent History Cache
participant R as Room State Service
C-xG: 网络中断
C->>G: 重连并发送 RESUME(lastSeq=928381)
G->>H: 查询 seq > 928381
alt 最近缓存仍存在
H-->>G: 返回缺失消息
G-->>C: RESUME_OK + missingMessages
else 缺口超出缓存窗口
G->>R: 获取房间当前快照
R-->>G: currentSeq / roomState
G-->>C: RESUME_RESET + latestWindow
end
16.3 重连风暴
节点故障可能导致几十万客户端同时重连,形成二次事故。
客户端必须使用:
- 指数退避;
- 随机抖动;
- 服务端下发建议重连时间;
- DNS 或负载均衡多入口;
- 最大重试频率;
- 前后台差异化重连。
例如:
1 | |
十七、存储设计
17.1 分层存储
flowchart TB
M[实时弹幕流] --> R[(Redis 最近消息<br/>秒级到分钟级)]
M --> MQ[(Kafka / Pulsar<br/>小时到天)]
MQ --> H[(历史消息存储<br/>天到月)]
H --> O[(对象存储 / 数据湖<br/>长期归档)]
H --> S[搜索与审核索引]
| 层级 | 用途 | 数据保留 |
|---|---|---|
| 网关内存 | 即时批量与重发 | 毫秒到秒 |
| Redis | 断线补拉、最近消息 | 30 秒到几分钟 |
| MQ | 解耦、重放、削峰 | 数小时到数天 |
| 历史数据库 | 回放、审核、查询 | 数天到数月 |
| 对象存储 | 离线分析、低成本归档 | 长期 |
17.2 历史表设计
可以按 roomId + 时间桶 分区:
1 | |
单库 MySQL 适合中小规模或重要消息索引;大规模追加写可以考虑:
- 分库分表;
- Cassandra / ScyllaDB;
- HBase;
- ClickHouse 用于分析;
- 对象存储保存压缩文件;
- Elasticsearch/OpenSearch 只保存需要搜索和审核的索引,不要承担全量主存储。
17.3 Redis Streams 能否替代 MQ
Redis Streams 可以保存有序消息并支持消费者组,适合:
- 中小规模;
- 较短保留周期;
- Redis 已是核心基础设施;
- 对跨地域、超大日志保留要求不高。
但需要理解:
- Consumer Group 更适合组内分摊消费,不等于向所有网关广播;
- 每房间一个 Stream 可能产生大量 Key;
- 超热门房间仍会形成热点;
- 长期保留、磁盘容量、故障恢复和跨机房能力需要单独评估。
因此,Redis Streams 可以是 MVP 到中等规模的方案,但不能因为它名字里有“Stream”,就自动获得所有大型消息平台能力。
17.4 Redis Pub/Sub 的定位
Redis Pub/Sub 适合低延迟、临时广播,但它不应被当作可靠历史日志:
- 订阅者离线期间无法自动补回普通 Pub/Sub 消息;
- 不适合承担重要付费消息的唯一传输;
- 缺少持久重放语义;
- 集群广播范围和带宽需要评估。
可以把它用于:
- 网关控制通知;
- 非关键实时事件;
- 小规模房间广播;
- 可靠 MQ 之后的低延迟加速通道。
十八、内容安全与反作弊
直播弹幕是公开内容,不能在消息广播后再慢慢审核。
18.1 多级审核链路
flowchart LR
A[原始弹幕] --> B[本地规则<br/>长度/字符/黑名单]
B --> C[用户与设备风险]
C --> D[快速内容模型]
D -->|低风险| E[正常广播]
D -->|中风险| F[降权/仅自己可见/抽样审核]
D -->|高风险| G[拒绝/禁言/人工审核]
E --> H[异步深度审核]
H -->|发现违规| I[撤回/处罚/更新模型]
18.2 快速路径与深度路径
主链路审核必须低延迟,可以包括:
- 敏感词;
- URL、联系方式、广告模式;
- 重复刷屏;
- Unicode 混淆;
- 特殊字符轰炸;
- 用户风险分;
- 轻量文本模型。
异步深度审核可以做:
- 上下文语义;
- 房间主题;
- 多轮行为;
- 群体攻击;
- 账号关联;
- 人工复核。
18.3 仅自己可见
对于疑似垃圾消息,可以采用“Shadow Ban”:
- 发送者本地看到;
- 其他用户不收到;
- 风控系统继续收集行为;
- 避免攻击者立即感知策略。
但该策略涉及产品和合规判断,不能仅由技术团队擅自决定。
18.4 安全注意事项
- 必须使用
wss://; - 校验 WebSocket Origin;
- 限制单帧和单消息大小;
- 不允许客户端提交可执行 HTML;
- 前端对内容做转义,防止 XSS;
- 过滤控制字符和异常 Unicode;
- 防止压缩炸弹;
- 限制连接建立频率;
- 防止 Ticket 重放;
- 用户、设备、IP、ASN 多维限流;
- 对管理指令做独立鉴权;
- 所有付费消息关联支付流水。
十九、限流设计
19.1 多维限流
| 维度 | 示例 |
|---|---|
| 用户 | 每 2 秒最多 1 条普通弹幕 |
| 设备 | 单设备每分钟最多 60 条 |
| IP | 单 IP 建连速率和发送速率 |
| 房间 | 每秒最多接收 N 条,超出进入采样 |
| 主播 | 高风险主播房间使用更严格规则 |
| 消息类型 | 普通弹幕、礼物、系统消息分别限流 |
| 网关 | 保护单节点 CPU、内存和网络 |
| 区域 | 防止单区域流量打满跨区链路 |
19.2 本地限流与全局限流
- 本地令牌桶:低延迟,保护单节点;
- Redis/Lua:实现用户级或房间级共享限流;
- 流式计算:用于秒级统计和动态调整;
- 风控系统:用于长周期行为判断。
不要让每条心跳、每条普通弹幕都同步访问远程 Redis 多次,否则限流系统本身会变成瓶颈。
推荐:
- 网关先做本地粗限流;
- 接入服务做用户级共享限流;
- 房间处理器做房间级总量控制;
- 风控系统进行分钟级、小时级策略。
二十、缓存与热点 Key
20.1 热门房间 Redis Key
单房间 ZSet、Seq、在线人数都可能成为热点 Key。
常见应对方式:
- 热门房间使用独立 Redis 分片;
- 近期消息按时间桶切分;
- 在线人数改为分片计数后聚合;
roomSeq在房间处理器内生成;- 使用本地缓存和批量写;
- 热门房间提前预热;
- 不在主路径同步执行大范围删除;
- 避免把超大消息正文重复存入多个结构。
20.2 在线人数不要求强一致
直播间显示“30.2 万人观看”通常不要求精确到个位。
可以使用:
1 | |
而不是每个用户进出都对同一个 Redis Key 执行 INCR/DECR。
二十一、客户端渲染同样是系统的一部分
服务器成功推送 5000 条/秒,不代表产品成功。客户端可能直接掉帧、发热或崩溃。
客户端需要:
- 消息去重;
- 小窗口排序;
- 最大待渲染队列;
- 每帧创建对象上限;
- 文本测量缓存;
- 弹幕轨道复用;
- 对重复文案聚合;
- 根据帧率降低显示密度;
- 后台状态暂停普通弹幕;
- 弱网时请求低码率消息档位;
- 重要消息使用独立图层;
- 避免在主线程反复解析超大 JSON;
- 使用对象池降低 GC 压力。
flowchart LR
A[网络批量消息] --> B[协议解码]
B --> C[messageId 去重]
C --> D[小窗口排序]
D --> E[优先级与采样]
E --> F[渲染队列]
F --> G[轨道调度]
G --> H[屏幕展示]
可以让客户端上报:
- FPS;
- 待渲染队列长度;
- 解码耗时;
- 丢弃数量;
- 网络 RTT;
- 设备等级。
服务端据此选择推送档位。
二十二、可观测性设计
22.1 核心指标
连接层
- 当前连接数;
- 每秒建连数和断连数;
- 鉴权失败率;
- 心跳超时率;
- 重连次数;
- 每网关连接数;
- Event Loop 延迟;
- 文件描述符使用率;
- Direct Memory;
- TLS 握手耗时。
消息层
- 发送 QPS;
- 审核通过率、拒绝率;
- MQ 写入延迟;
- 消费 Lag;
- Room Processor 处理延迟;
- 每房间消息速率;
- 批量大小;
- 重复消息率;
- 缺口恢复率。
推送层
- Fanout 放大系数;
- 每秒发送字节;
- 每个房间目标网关数;
- 网关发送队列;
- 慢消费者比例;
- 普通消息丢弃率;
- 高优先级消息失败率;
- 端到端 P50/P95/P99 延迟。
存储层
- Redis 命中率;
- 热 Key;
- ZSet/Stream 长度;
- 历史写入 Lag;
- 数据库写入吞吐;
- 分区大小;
- 回放查询延迟。
22.2 端到端延迟
每条抽样消息记录时间点:
1 | |
最终可以拆出:
1 | |
只看服务端接口耗时,会遗漏真正影响用户体验的网络和客户端阶段。
22.3 日志与追踪
建议统一字段:
1 | |
普通弹幕量巨大,不能全量打印 INFO 日志。应使用:
- 指标聚合;
- 按比例采样;
- 错误消息全量;
- 重要消息全链路;
- 热门房间动态采样;
- 离线日志分析。
二十三、高可用与容灾
23.1 单节点故障
Connection Gateway 无需保存不可恢复业务状态。节点故障后:
- TCP 连接断开;
- 客户端指数退避重连;
- 新网关重新鉴权;
- 恢复房间订阅;
- 使用
lastSeq补最近消息。
23.2 Redis 故障
- 最近弹幕补拉暂时不可用;
- 主推送链路继续依赖 MQ;
- 路由表可使用本地缓存和 TTL;
- 在线人数展示降级为估算值;
- 不允许 Redis 故障阻断重要消息写入。
23.3 MQ 故障
- Ingress 快速失败或进入有限本地缓冲;
- 不应无限堆积在应用内存;
- 普通弹幕可提示发送失败;
- 付费消息走业务事务和补偿通道;
- Broker 恢复后按业务规则重放;
- 监控 Producer 错误率和 ISR/副本状态。
23.4 数据库故障
历史写入消费者停止提交 Offset,消息继续保留在 MQ;数据库恢复后继续消费。
如果 MQ 保留时间不足,应该先写对象存储或增加临时落盘,避免历史数据永久丢失。
23.5 多机房
flowchart TB
DNS[Global DNS / Anycast]
DNS --> R1[Region A]
DNS --> R2[Region B]
subgraph A[Region A]
G1[Gateway Cluster]
F1[Regional Fanout]
M1[(Regional MQ)]
end
subgraph B[Region B]
G2[Gateway Cluster]
F2[Regional Fanout]
M2[(Regional MQ)]
end
C[Global Room Event Backbone]
M1 --> C
M2 --> C
C --> F1
C --> F2
多机房需要回答:
- 用户优先接入最近区域;
- 房间主处理器位于哪个区域;
- 跨区域消息如何复制;
- 区域间网络中断时是否允许局部弹幕;
- 恢复后是否合并;
- 付费消息如何保证全局幂等;
- 房间顺序是全局顺序还是区域内顺序。
普通弹幕可以接受区域级最终一致;支付和礼物流水则需要独立的强一致业务系统。
二十四、部署与发布
24.1 网关滚动发布
长连接服务不能像普通 HTTP 服务一样直接杀 Pod。
正确流程:
sequenceDiagram
participant O as Orchestrator
participant G as Old Gateway
participant L as Load Balancer
participant C as Clients
O->>G: 标记 Draining
G->>L: Readiness = false
L-->>G: 不再分配新连接
G->>C: 可选:下发迁移/重连提示
G->>G: 等待连接自然退出
G->>C: 超时后优雅关闭
O->>G: 终止实例
需要配置:
- 足够长的
terminationGracePeriodSeconds; preStop钩子;- 停止接收新连接;
- 路由表及时注销;
- 连接排空超时;
- 客户端重连抖动;
- 发布期间限制同时下线节点比例。
24.2 HPA 指标
长连接网关不能只根据 CPU 扩容,还应考虑:
- 活跃连接数;
- 每秒发送字节;
- Event Loop 延迟;
- 出站队列长度;
- Direct Memory;
- 新建连接速率;
- 网卡利用率;
- 慢消费者比例。
CPU 很低不代表网关有余量,网卡和内存可能早已接近极限。
24.3 反向代理参数
以 Nginx 类反向代理为例,需要正确处理:
- HTTP Upgrade;
Connection头;- 空闲读取超时;
- 上游连接超时;
- 最大连接数;
- TLS;
- 日志采样。
WebSocket 是长连接,代理默认 60 秒无数据关闭连接会导致周期性断线,因此心跳和代理超时必须配套设计。
二十五、容量规划
25.1 单网关连接数
不要直接相信“单机可以支持百万连接”之类的宣传数字。真实容量取决于:
- 操作系统;
- 文件描述符;
- Event Loop 实现;
- TLS;
- 每连接内存;
- 心跳频率;
- 消息速率;
- 出站带宽;
- GC;
- 客户端网络质量;
- 消息编码;
- 是否启用压缩。
假设每连接综合占用 20 KB:
1 | |
这只是连接对象,不包含发送队列、Direct Buffer、线程栈、JVM、日志和系统页缓存。
生产容量应通过压测确定,并保留至少 30%~50% 安全余量。
25.2 网关数量估算
假设:
- 峰值连接 200 万;
- 单网关安全承载 5 万连接;
- 预留 40% 余量。
1 | |
还要按区域、可用区和故障域分散,不能全部集中在一个节点池。
25.3 带宽比连接数更容易成为瓶颈
假设单台网关 1 万名热门房间用户,每人平均收到 20 条/秒,每条线上的平均大小 220 字节:
1 | |
加上协议、TLS、重传和其他消息,可能很快接近 1 Gbps 网卡的安全上限。
因此:
- 批量;
- 二进制协议;
- 采样;
- 聚合;
- 区域就近推送;
- 合理压缩;
- 控制单机热门房间用户数;
往往比单纯增加 CPU 更重要。
二十六、压测与故障演练
26.1 压测场景
不能只测试平均流量,应至少覆盖:
- 200 万长连接稳定在线;
- 10 分钟内快速建立 100 万连接;
- 单热门房间 30 万在线;
- 热门房间 5000 条/秒;
- 100 个热点房间同时爆发;
- 50% 客户端为慢消费者;
- 网关节点随机下线;
- Redis 主从切换;
- MQ 消费延迟;
- 历史数据库写入暂停;
- 跨区域网络延迟和丢包;
- 大规模客户端同时重连;
- 内容审核服务超时;
- 滚动发布期间流量峰值。
26.2 验证指标
- 是否出现 OOM;
- Event Loop 是否阻塞;
- P99 延迟是否失控;
- 普通消息丢弃是否符合策略;
- 高优先级消息是否仍可到达;
- 重连是否出现雪崩;
- MQ Lag 是否可恢复;
- Redis 热 Key 是否打满;
- 单网关是否流量不均;
- 路由记录是否泄漏;
- 发布后连接是否正确迁移。
二十七、架构演进路线
阶段一:MVP
适用规模:
- 几千到几万在线;
- 弹幕量较低;
- 团队较小;
- 快速验证产品。
组件:
1 | |
阶段二:WebSocket 化
适用规模:
- 实时性要求提高;
- 轮询成本明显;
- 需要主动控制消息。
组件:
1 | |
阶段三:可靠消息与网关路由
适用规模:
- 百万级连接;
- 多个热门房间;
- 需要稳定扩容和回放。
组件:
1 | |
阶段四:超级热点与多地域
适用规模:
- 大型赛事;
- 大促;
- 跨国直播;
- 单房间数十万以上在线。
组件:
1 | |
flowchart LR
A[MVP<br/>轮询 + Redis] --> B[WebSocket<br/>长连接推送]
B --> C[MQ + 房间处理器<br/>可靠解耦]
C --> D[网关级广播<br/>路由表]
D --> E[热门房间专属集群<br/>动态采样]
E --> F[多地域 Relay<br/>分层扇出]
演进原则是:
先解决当前规模下最贵的问题,不要为了想象中的千万并发,把第一版做成一座没人敢改的分布式博物馆。
二十八、常见设计误区
误区一:每条消息直接循环所有用户
会造成用户级写扩散。应改为“中心到网关、网关到本机连接”的两级广播。
误区二:所有消息都必须严格不丢
普通弹幕、系统通知、付费消息的可靠性等级不同。统一最高可靠性会显著增加成本和延迟。
误区三:使用客户端时间排序
客户端时间可能错误、漂移或被篡改。应使用服务端时间和房间序号。
误区四:用时间戳当唯一 ZSet member
ZSet member 必须唯一,相同正文会覆盖。应加入消息 ID。
误区五:Kafka 一个 Consumer Group 可以广播给所有网关
Consumer Group 是负载均衡消费,不是广播。需要 Fanout Router 或独立发布订阅层。
误区六:长连接服务完全无状态
连接本身天然有状态。正确目标是业务状态可恢复,而不是假装连接不存在。
误区七:只按 CPU 做扩容
长连接服务常常先碰到带宽、文件描述符、内存、发送队列和 Event Loop 延迟。
误区八:所有用户必须看到完全相同的普通弹幕
热门房间中,客户端根本无法展示全部消息。合理采样比追求无意义的一致更重要。
误区九:Redis 是唯一持久化存储
Redis 可以承担最近消息和快速恢复,但重要消息、回放和审计需要可靠日志或数据库。
误区十:服务端推送成功就代表用户看到了
客户端可能排队、丢弃或渲染失败。真正的用户体验指标应覆盖客户端接收和展示阶段。
二十九、一套推荐的最终方案
综合前面的讨论,一套可落地的中大型直播弹幕架构如下:
flowchart TB
subgraph Clients[客户端层]
APP[App / Web / TV]
end
subgraph Access[接入层]
GLB[Global LB]
WS[WebSocket Gateway Cluster]
end
subgraph WritePath[写入链路]
ING[Danmaku Ingress]
RL[Rate Limit & Risk]
MOD[Content Moderation]
K1[(Kafka Ingress Topic)]
RP[Room Processor]
K2[(Kafka Fanout Topic)]
end
subgraph PushPath[推送链路]
FR[Global Fanout Router]
RR[Regional Relay]
GW[Target Gateways]
end
subgraph State[状态与缓存]
REG[(Room-Gateway Registry)]
REC[(Redis Recent Messages)]
ROOM[(Room State)]
end
subgraph Storage[存储与分析]
HW[History Writer]
DB[(History DB)]
OS[(Object Storage)]
RT[Realtime Analytics]
MON[Metrics / Logs / Traces]
end
APP --> GLB --> WS
WS --> ING
ING --> RL --> MOD --> K1
K1 --> RP
RP --> K2
RP --> REC
RP --> ROOM
K2 --> FR
FR --> REG
FR --> RR
RR --> GW
GW --> APP
K1 --> HW --> DB --> OS
K1 --> RT
WS --> MON
ING --> MON
RP --> MON
FR --> MON
关键决策:
- WebSocket 负责双向实时通信;
- 网关只维护连接和本机订阅;
- 写入先经过限流和快速审核;
- MQ 按
roomId分区; - Room Processor 维护房间内顺序;
- 全局路由只记录房间到网关;
- 每条消息对每个目标网关发送一份;
- 网关在本机批量广播;
- Redis 只保存最近窗口和可恢复状态;
- 历史存储异步写入;
- 热门房间使用独立分区、区域 Relay 和动态采样;
- 普通消息允许降级,重要消息使用独立可靠通道;
- 慢消费者必须有明确的丢弃和断开策略;
- HPA 同时关注连接、带宽、内存和发送队列;
- 客户端渲染能力纳入整体系统设计。
三十、面试中的精简回答
如果在系统设计面试中只有几分钟,可以这样概括:
直播弹幕本质是一个以直播间为会话空间的高扇出消息系统。早期可以使用 HTTP 轮询配合 Redis ZSet 保存最近消息,通过时间和序号游标增量拉取;规模增长后切换到 WebSocket 长连接。
写入链路先做鉴权、限流和内容审核,再按 roomId 写入 Kafka,同一房间进入同一分区以保证房间内顺序。Room Processor 分配 roomSeq,Fanout Router 查询 roomId 对应的网关集合,每条消息只向有该房间用户的网关发送一份,网关再向本机连接广播,从而避免中心服务对用户级写扩散。
Redis 保存最近几十秒消息,用于断线补拉;数据库或分布式存储异步保存回放和审计数据。热门房间要使用独立分区、批量推送、动态采样、重复消息聚合和区域级 Relay。网关对慢客户端设置有界发送队列,优先保留系统、礼物和付费消息,普通弹幕在过载时允许丢弃。
最终需要重点监控连接数、建连速率、端到端延迟、MQ Lag、路由放大系数、网关发送队列、丢弃率和出口带宽,并通过指数退避避免故障后的重连风暴。
参考资料
王帅真:《直播弹幕系统设计》
https://blog.qizong007.top/article/live-streaming-bullet-systemIETF RFC 6455:The WebSocket Protocol
https://www.rfc-editor.org/rfc/rfc6455IETF RFC 8441:Bootstrapping WebSockets with HTTP/2
https://datatracker.ietf.org/doc/html/rfc8441Redis Documentation:ZADD
https://redis.io/docs/latest/commands/zadd/Redis Documentation:Redis Streams
https://redis.io/docs/latest/develop/data-types/streams/Redis Documentation:Pub/Sub
https://redis.io/docs/latest/develop/interact/pubsub/Apache Kafka Documentation
https://kafka.apache.org/documentation/WHATWG HTML Standard:Server-sent events
https://html.spec.whatwg.org/multipage/server-sent-events.htmlNginx Documentation:WebSocket proxying
https://nginx.org/en/docs/http/websocket.html
总结
直播弹幕系统不是简单的“WebSocket + Redis”。
真正决定系统能否扛住热点流量的,是以下几件事:
- 是否完成用户级广播到网关级广播的转换;
- 是否把消息写入、顺序处理、推送和历史存储解耦;
- 是否承认客户端不可能展示全部消息;
- 是否针对热门房间实施独立容量和降级策略;
- 是否给慢消费者设置明确边界;
- 是否区分普通弹幕和重要业务消息的可靠性;
- 是否能在节点故障和大规模重连时保持稳定;
- 是否把出口带宽和客户端渲染纳入容量设计。
一个成熟的弹幕系统,并不是保证每一条“666”都以最高等级的可靠性抵达世界尽头,而是在流量洪峰中仍能保证直播可用、重要消息可靠、普通弹幕足够实时,并让系统有能力持续演进。