从 MVP 到百万并发:直播弹幕系统完整设计

本文参考并延伸了王帅真《直播弹幕系统设计》中的演进思路,在其“轮询 MVP → 长连接推送 → MQ 削峰 → 热门直播间优化”的基础上,补充容量估算、协议设计、房间路由、顺序一致性、可靠性、慢消费者治理、内容安全、容灾和运维方案。

本文不是对某一家直播平台内部实现的复刻,而是一套可以用于系统设计面试、技术方案评审和真实项目落地的通用架构。

一、问题定义:直播弹幕系统到底在解决什么

从表面看,直播弹幕只是用户发送一段文本,然后让直播间里的其他用户看到。

但当一个直播间从几十人增长到几十万人时,问题就不再是“保存一条消息”,而是:

  1. 如何维持海量客户端长连接;
  2. 如何让一条弹幕尽快到达同一房间内的用户;
  3. 如何避免一条消息被直接复制几十万次,形成巨大的写扩散;
  4. 如何处理热门直播间的瞬时流量;
  5. 如何保证同一直播间内消息大体有序;
  6. 如何在断线重连后补齐少量遗漏消息;
  7. 如何过滤违规内容、广告、刷屏和恶意请求;
  8. 如何在系统过载时优先保住直播,而不是让弹幕拖垮整个平台;
  9. 是否需要保存历史弹幕,以及如何支持直播回放;
  10. 如何控制带宽、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
100000 × 220 B ≈ 22 MB/s

22 MB/s 的入口流量对一个经过扩容的消息集群并不夸张。

真正危险的是热门房间的广播。假设一个房间每秒产生 5000 条弹幕,在线用户 30 万,每条消息 220 字节,如果对每个用户逐条复制:

1
2
5000 × 300000 × 220 B
≈ 330 GB/s

这还没有计算 TCP、TLS、WebSocket 帧和公网传输开销。

因此,大型弹幕系统绝不能简单实现成:

1
2
3
收到 1 条弹幕
→ 查询房间 30 万用户
→ 循环发送 30 万次

这段代码在功能测试里可能很优雅,在生产环境里则像一个精心编写的自毁按钮。

2.3 系统目标分级

不同类型的消息,可靠性和优先级不应相同。

消息类型 示例 可靠性 优先级
普通弹幕 “主播好”“666” 允许极少量丢失
系统通知 开播、封禁、房间状态变化 应可靠到达
礼物消息 用户赠送礼物 必须可追踪、可补偿
付费弹幕 醒目留言、超级留言 不应丢失 最高
互动事件 点赞聚合、投票变化 可合并、可采样
风控指令 禁言、踢出、关闭房间 必须及时执行 最高

这意味着系统不应只有一个统一的“消息通道”。至少要在协议、队列或调度层面区分优先级。


三、核心业务流程

一条弹幕从用户 A 到达同房间用户 B,至少经过以下阶段:

  1. 客户端建立连接并完成鉴权;
  2. 客户端加入直播间;
  3. 用户发送弹幕;
  4. 接入层校验参数、身份、频率和房间状态;
  5. 内容安全系统进行快速审核;
  6. 消息写入可靠消息通道;
  7. 房间处理器分配顺序标识;
  8. 路由系统找到当前订阅该房间的网关节点;
  9. 网关在本机对连接进行广播;
  10. 客户端去重、排序、限速并渲染;
  11. 历史写入服务异步落库;
  12. 实时分析系统消费消息,更新热度和运营指标。
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 万份。


四、消息数据模型

直播弹幕看似只有 roomIduserIdcontent,但真正上线后会迅速出现去重、顺序、重试、审核、样式和回放问题。

推荐的数据结构如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
{
"messageId": "01JZ8X9YQ4J9K1P7X2R5F8A6BC",
"clientMsgId": "8e88b9d0-97a4-4a87-b5d5-4c6b991f0c67",
"roomId": "live_9527",
"userId": "user_10086",
"messageType": "DANMAKU",
"content": "这波操作太强了",
"serverTime": 1785297600123,
"roomSeq": 928381,
"priority": 10,
"style": {
"color": "#FFFFFF",
"fontSize": 14,
"position": "SCROLL"
},
"moderation": {
"status": "PASS",
"riskLevel": "LOW"
},
"extension": {
"userLevel": 12,
"badge": "VIP"
},
"traceId": "6bb32f92f35c48cf"
}

4.1 字段说明

字段 作用
messageId 服务端生成的全局唯一消息 ID
clientMsgId 客户端生成,用于请求重试时幂等去重
roomId 直播间标识,也是消息分区和路由的重要依据
messageType 普通弹幕、系统消息、礼物、禁言指令等
serverTime 服务端接收或确认时间,避免依赖不可信的客户端时间
roomSeq 房间内顺序号,用于排序、断线补偿和缺口检测
priority 推送和丢弃策略的依据
moderation.status 内容审核状态
extension 可扩展业务属性,避免频繁修改主协议
traceId 链路追踪标识

4.2 直播是否需要视频时间偏移

纯实时直播中,弹幕通常按服务器时间到达即可,并不一定需要点播视频那样的固定偏移量。

但以下场景仍然需要视频时间轴字段:

  • 直播回放;
  • DVR 时移播放;
  • 用户主动拖动直播进度;
  • 主播端和观众端存在不同播放缓冲;
  • 多 CDN 节点造成播放延迟差异;
  • 需要把弹幕重新挂载到录播文件。

因此可以增加可选字段:

1
2
3
4
{
"videoPts": 186320,
"streamEpoch": "2026-07-29T11:30:00Z"
}

其中 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
2
3
Key: danmaku:room:{roomId}:recent
Score: serverTimeMillis
Member: roomSeq|messageId|payload

写入:

1
ZADD danmaku:room:{roomId}:recent <serverTimeMillis> <member>

查询:

1
ZRANGEBYSCORE danmaku:room:{roomId}:recent <lastTimestamp> +inf LIMIT 0 200

Redis Sorted Set 允许不同成员使用相同分值,但成员本身必须唯一。同分值成员会进一步按字典序排列。因此,不能只把弹幕正文作为 member,否则相同内容会互相覆盖;应把 messageIdroomSeq 编入 member。

5.2 时间戳游标的缺陷

如果客户端只记录 lastTimestamp,同一毫秒内存在多条消息时容易出现:

  • 重复拉取;
  • 边界消息遗漏;
  • 同分值顺序不稳定;
  • 客户端重试时难以精确恢复。

更稳妥的方式是返回复合游标:

1
2
3
4
5
6
{
"cursor": {
"serverTime": 1785297600123,
"roomSeq": 928381
}
}

服务端按照 (serverTime, roomSeq) 做稳定排序。

如果系统已经具备连续房间序号,也可以直接使用 roomSeq 作为增量游标。

5.3 清理旧消息

Redis 只保留最近 30 秒或最近 5000 条消息:

1
ZREMRANGEBYSCORE key -inf <expireTimestamp>

或者:

1
ZREMRANGEBYRANK key 0 -(maxCount + 1)

清理动作可以由写入 Lua 脚本、定时任务或后台消费者完成,避免每次读请求都执行昂贵清理。

5.4 本地缓存降低重复读取

同一直播间的客户端会在相近时间请求几乎相同的数据。如果每次都回源 Redis,会造成大量重复读取。

查询服务可以缓存最近 1~5 秒的房间消息:

1
roomId -> recentMessageBatch

优化原则:

  1. 只缓存活跃房间;
  2. 设置最大房间数量和最大内存;
  3. 使用 Caffeine 等带权重淘汰的本地缓存;
  4. 不要因为房间多就无限创建缓存项;
  5. 对热门房间采用请求合并,避免同一时刻大量回源;
  6. 使用一致性哈希或房间路由,使同一房间的读请求尽量命中相同节点。

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 请求头。可以采用以下方式:

  1. 客户端先通过 HTTPS 使用登录凭证换取一个短期、一次性的 WebSocket Ticket;
  2. 客户端使用 Ticket 建立连接;
  3. 网关校验 Ticket,绑定用户、设备和会话;
  4. Ticket 使用后立即失效。
1
2
3
4
5
6
7
8
POST /api/ws/ticket
Authorization: Bearer <access-token>

Response:
{
"ticket": "short-lived-one-time-ticket",
"expiresIn": 30
}
1
wss://push.example.com/ws?ticket=<ticket>

不要把长期 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
2
3
4
{
"type": "PING",
"clientTime": 1785297600000
}

服务端响应:

1
2
3
4
{
"type": "PONG",
"serverTime": 1785297600026
}

心跳间隔不能过短,否则百万连接会产生大量无意义请求。例如 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
2
3
4
5
6
7
{
"type": "SEND_DANMAKU",
"requestId": "req-101",
"clientMsgId": "local-9527",
"roomId": "live-9527",
"content": "主播这波操作太秀了"
}

服务端确认:

1
2
3
4
5
6
7
8
9
{
"type": "SEND_ACK",
"requestId": "req-101",
"clientMsgId": "local-9527",
"messageId": "01JZ8X9YQ4J9K1P7X2R5F8A6BC",
"roomSeq": 928381,
"status": "ACCEPTED",
"serverTime": 1785297600123
}

客户端可以先本地乐观展示,再根据 ACK 将临时消息替换为正式消息。如果服务端拒绝,则标记发送失败,而不是悄悄消失。

9.3 批量推送

逐条发送 WebSocket 帧会增加系统调用、对象创建和协议开销。网关可以按 10~50 毫秒窗口聚合:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
{
"type": "BATCH_MESSAGES",
"roomId": "live-9527",
"fromSeq": 928381,
"toSeq": 928390,
"messages": [
{
"messageId": "m1",
"roomSeq": 928381,
"content": "第一条"
},
{
"messageId": "m2",
"roomSeq": 928382,
"content": "第二条"
}
]
}

批量窗口过大,会增加延迟;过小,则节省不了多少开销。需要结合消息速率动态调整。

9.4 JSON 还是 Protobuf

方案 优点 缺点
JSON 易调试、前端接入简单 体积较大、解析成本较高
Protobuf 体积小、解析快、Schema 明确 调试和版本治理更复杂
MessagePack 比 JSON 紧凑 生态和 Schema 约束弱于 Protobuf

早期可以使用 JSON,规模上升后切换 Protobuf,并保留协议版本:

1
magic | version | command | requestId | payloadLength | payload

十、消息写入链路

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
(userId, clientMsgId) -> messageId

重复请求直接返回原 ACK。

幂等记录可以放在 Redis,并设置几分钟过期。对于付费消息,还应使用数据库唯一约束或业务流水号,不能只依赖短期缓存。


十一、消息队列与房间顺序

11.1 按 roomId 分区

消息队列的分区键应使用 roomId

1
partition = hash(roomId) % partitionCount

这样同一房间的消息进入同一个分区,并由一个房间处理器顺序消费。

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
INCR danmaku:room:{roomId}:seq

优点是简单;缺点是热门房间会形成单 Key 热点,Redis 故障和迁移也需要处理序号连续性。

方案二:房间处理器内存计数

同一房间始终由同一个顺序消费者处理,在内存中递增,并周期性持久化检查点。

优点:

  • 无需每条消息访问 Redis;
  • 延迟低;
  • 与分区消费模型自然结合。

缺点:

  • 消费者切换时需要恢复;
  • 检查点落后可能产生重复或跳号;
  • 需要使用 Epoch 区分不同处理实例。

可使用:

1
roomSeq = epoch + localCounter

更准确地说,可以把 epoch 放在高位,把本地计数器放在低位,避免旧消费者恢复后与新消费者产生序号冲突。

方案三:直接使用消息队列 Offset

队列 Offset 在分区内单调递增,但一个分区通常承载多个房间,因此对某个房间来说会出现大量跳号。

它可以用于排序,但不适合直接作为“连续房间序号”进行缺口判断。

11.4 是否需要 Exactly Once

普通弹幕没有必要追求端到端 Exactly Once。

更现实的目标是:

  • MQ 至少一次投递;
  • 消息处理器允许重复消费;
  • 使用 messageId 幂等;
  • 网关和客户端去重;
  • 顺序以 roomSeq 为准;
  • 重要消息使用业务流水和持久化补偿。

Exactly Once 往往只是把复杂性转移到了事务、状态和边界定义中。对于普通弹幕,少量重复比大幅降低吞吐更容易处理。


十二、房间与网关路由

12.1 不要在中心存储每个用户的连接

一种直觉设计是:

1
2
3
roomId -> userId 列表
userId -> connectionId
connectionId -> gatewayId

当房间有几十万用户时,每条消息都遍历用户列表,中心路由系统会成为瓶颈。

更好的模型是:

1
2
roomId -> gatewayId 集合
gatewayId + roomId -> 本机 connection 集合

全局只记录“哪些网关有这个房间的用户”,具体用户连接由网关本地维护。

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
2
3
4
5
6
7
8
9
10
11
12
13
{
"roomId": "live-9527",
"gateways": [
{
"gatewayId": "gw-sg-az1-001",
"region": "sg",
"az": "az1",
"localSubscribers": 12000,
"expireAt": 1785297660000
}
],
"version": 381
}

12.3 路由存储选型

可选方案包括:

  • Redis Set / Hash;
  • etcd;
  • 服务注册中心;
  • 专用路由服务;
  • 网关主动订阅 Fanout Router;
  • 基于 Gossip 的节点状态同步。

需要注意:etcd 更适合低频配置和服务发现,不适合每个用户进出房间都直接写。通过“首个用户注册、最后一个用户注销”可以把写入频率从用户级降低到网关级。


十三、避免广播风暴:两级与三级扇出

13.1 两级扇出

1
2
3
4
Room Processor
→ Fanout Router
→ 目标 Gateway
→ Gateway 本机连接

中心层按网关发送,网关按本地连接发送。

13.2 三级扇出

跨地域或超热门直播间可以增加区域中继:

1
2
3
4
Global Fanout
→ Region Relay
→ Gateway
→ Client
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
subPartition = hash(messageId) % N

客户端按到达顺序展示,或在 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
2
3
4
5
6
7
8
High Priority Queue:
系统指令、付费消息、禁言、房间关闭

Normal Priority Queue:
普通弹幕、点赞聚合

Threshold:
最大 256 条或 512 KB

15.2 丢弃策略

当连接拥塞时:

  1. 优先丢弃过期普通弹幕;
  2. 丢弃重复、低权重弹幕;
  3. 合并点赞和同文案消息;
  4. 降低普通弹幕推送频率;
  5. 保留系统、礼物和付费消息;
  6. 队列持续增长时断开连接,让客户端指数退避后重连。

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
2
3
4
roomId
lastReceivedSeq
lastAckedClientMsgId
sessionId

重连后发送:

1
2
3
4
5
{
"type": "RESUME",
"roomId": "live-9527",
"lastSeq": 928381
}

服务端处理:

  1. 检查最近消息缓存是否仍包含 lastSeq 之后的数据;
  2. 如果缺口较小,返回增量消息;
  3. 如果缺口太大,只返回最新窗口;
  4. 重要状态通过房间快照重新同步;
  5. 客户端按 messageId 去重;
  6. 如果 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
2
3
4
5
1s + random(0~500ms)
2s + random(0~1s)
4s + random(0~2s)
...
最大 30s

十七、存储设计

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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
CREATE TABLE live_danmaku (
room_id VARCHAR(64) NOT NULL,
bucket_date DATE NOT NULL,
room_seq BIGINT NOT NULL,
message_id VARCHAR(32) NOT NULL,
user_id VARCHAR(64) NOT NULL,
message_type SMALLINT NOT NULL,
content VARCHAR(500) NOT NULL,
server_time TIMESTAMP(3) NOT NULL,
video_pts BIGINT NULL,
audit_status SMALLINT NOT NULL,
extension JSON NULL,
PRIMARY KEY (room_id, bucket_date, room_seq),
UNIQUE KEY uk_message_id (message_id)
);

单库 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 多次,否则限流系统本身会变成瓶颈。

推荐:

  1. 网关先做本地粗限流;
  2. 接入服务做用户级共享限流;
  3. 房间处理器做房间级总量控制;
  4. 风控系统进行分钟级、小时级策略。

二十、缓存与热点 Key

20.1 热门房间 Redis Key

单房间 ZSet、Seq、在线人数都可能成为热点 Key。

常见应对方式:

  • 热门房间使用独立 Redis 分片;
  • 近期消息按时间桶切分;
  • 在线人数改为分片计数后聚合;
  • roomSeq 在房间处理器内生成;
  • 使用本地缓存和批量写;
  • 热门房间提前预热;
  • 不在主路径同步执行大范围删除;
  • 避免把超大消息正文重复存入多个结构。

20.2 在线人数不要求强一致

直播间显示“30.2 万人观看”通常不要求精确到个位。

可以使用:

1
2
3
4
gateway local count
→ 周期上报
→ 区域聚合
→ 全局估算

而不是每个用户进出都对同一个 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
2
3
4
5
6
7
8
9
T1 客户端发送
T2 网关接收
T3 Ingress 接收
T4 MQ 写入成功
T5 Room Processor 处理
T6 Fanout Router 发出
T7 目标网关写 Socket
T8 客户端接收
T9 客户端渲染

最终可以拆出:

1
2
3
4
5
6
7
8
上行网络
+ 接入处理
+ 审核
+ MQ
+ 房间处理
+ 路由
+ 下行网络
+ 客户端渲染

只看服务端接口耗时,会遗漏真正影响用户体验的网络和客户端阶段。

22.3 日志与追踪

建议统一字段:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
traceId
messageId
clientMsgId
roomId
userId
gatewayId
region
partition
offset
roomSeq
messageType
priority
moderationStatus
costMs
dropReason

普通弹幕量巨大,不能全量打印 INFO 日志。应使用:

  • 指标聚合;
  • 按比例采样;
  • 错误消息全量;
  • 重要消息全链路;
  • 热门房间动态采样;
  • 离线日志分析。

二十三、高可用与容灾

23.1 单节点故障

Connection Gateway 无需保存不可恢复业务状态。节点故障后:

  1. TCP 连接断开;
  2. 客户端指数退避重连;
  3. 新网关重新鉴权;
  4. 恢复房间订阅;
  5. 使用 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
100000 连接 × 20 KB ≈ 2 GB

这只是连接对象,不包含发送队列、Direct Buffer、线程栈、JVM、日志和系统页缓存。

生产容量应通过压测确定,并保留至少 30%~50% 安全余量。

25.2 网关数量估算

假设:

  • 峰值连接 200 万;
  • 单网关安全承载 5 万连接;
  • 预留 40% 余量。
1
2
基础节点数 = 2000000 / 50000 = 40
考虑余量后 ≈ 56~60 台

还要按区域、可用区和故障域分散,不能全部集中在一个节点池。

25.3 带宽比连接数更容易成为瓶颈

假设单台网关 1 万名热门房间用户,每人平均收到 20 条/秒,每条线上的平均大小 220 字节:

1
2
3
10000 × 20 × 220 B
≈ 44 MB/s
≈ 352 Mbps

加上协议、TLS、重传和其他消息,可能很快接近 1 Gbps 网卡的安全上限。

因此:

  • 批量;
  • 二进制协议;
  • 采样;
  • 聚合;
  • 区域就近推送;
  • 合理压缩;
  • 控制单机热门房间用户数;

往往比单纯增加 CPU 更重要。


二十六、压测与故障演练

26.1 压测场景

不能只测试平均流量,应至少覆盖:

  1. 200 万长连接稳定在线;
  2. 10 分钟内快速建立 100 万连接;
  3. 单热门房间 30 万在线;
  4. 热门房间 5000 条/秒;
  5. 100 个热点房间同时爆发;
  6. 50% 客户端为慢消费者;
  7. 网关节点随机下线;
  8. Redis 主从切换;
  9. MQ 消费延迟;
  10. 历史数据库写入暂停;
  11. 跨区域网络延迟和丢包;
  12. 大规模客户端同时重连;
  13. 内容审核服务超时;
  14. 滚动发布期间流量峰值。

26.2 验证指标

  • 是否出现 OOM;
  • Event Loop 是否阻塞;
  • P99 延迟是否失控;
  • 普通消息丢弃是否符合策略;
  • 高优先级消息是否仍可到达;
  • 重连是否出现雪崩;
  • MQ Lag 是否可恢复;
  • Redis 热 Key 是否打满;
  • 单网关是否流量不均;
  • 路由记录是否泄漏;
  • 发布后连接是否正确迁移。

二十七、架构演进路线

阶段一:MVP

适用规模:

  • 几千到几万在线;
  • 弹幕量较低;
  • 团队较小;
  • 快速验证产品。

组件:

1
HTTP API + Redis ZSet + 定时轮询 + MySQL 异步落库

阶段二:WebSocket 化

适用规模:

  • 实时性要求提高;
  • 轮询成本明显;
  • 需要主动控制消息。

组件:

1
WebSocket Gateway + Redis 最近消息 + 简单 Pub/Sub + 数据库

阶段三:可靠消息与网关路由

适用规模:

  • 百万级连接;
  • 多个热门房间;
  • 需要稳定扩容和回放。

组件:

1
2
3
4
5
6
WebSocket Gateway
+ Kafka/Pulsar
+ Room Processor
+ Fanout Router
+ Room-Gateway Registry
+ 分层存储

阶段四:超级热点与多地域

适用规模:

  • 大型赛事;
  • 大促;
  • 跨国直播;
  • 单房间数十万以上在线。

组件:

1
2
3
4
5
6
独立热点房间集群
+ 区域 Relay
+ 分层广播
+ 动态采样
+ 多 Lane 并行
+ 多地域容灾
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

关键决策:

  1. WebSocket 负责双向实时通信;
  2. 网关只维护连接和本机订阅;
  3. 写入先经过限流和快速审核;
  4. MQ 按 roomId 分区;
  5. Room Processor 维护房间内顺序;
  6. 全局路由只记录房间到网关;
  7. 每条消息对每个目标网关发送一份;
  8. 网关在本机批量广播;
  9. Redis 只保存最近窗口和可恢复状态;
  10. 历史存储异步写入;
  11. 热门房间使用独立分区、区域 Relay 和动态采样;
  12. 普通消息允许降级,重要消息使用独立可靠通道;
  13. 慢消费者必须有明确的丢弃和断开策略;
  14. HPA 同时关注连接、带宽、内存和发送队列;
  15. 客户端渲染能力纳入整体系统设计。

三十、面试中的精简回答

如果在系统设计面试中只有几分钟,可以这样概括:

直播弹幕本质是一个以直播间为会话空间的高扇出消息系统。早期可以使用 HTTP 轮询配合 Redis ZSet 保存最近消息,通过时间和序号游标增量拉取;规模增长后切换到 WebSocket 长连接。

写入链路先做鉴权、限流和内容审核,再按 roomId 写入 Kafka,同一房间进入同一分区以保证房间内顺序。Room Processor 分配 roomSeq,Fanout Router 查询 roomId 对应的网关集合,每条消息只向有该房间用户的网关发送一份,网关再向本机连接广播,从而避免中心服务对用户级写扩散。

Redis 保存最近几十秒消息,用于断线补拉;数据库或分布式存储异步保存回放和审计数据。热门房间要使用独立分区、批量推送、动态采样、重复消息聚合和区域级 Relay。网关对慢客户端设置有界发送队列,优先保留系统、礼物和付费消息,普通弹幕在过载时允许丢弃。

最终需要重点监控连接数、建连速率、端到端延迟、MQ Lag、路由放大系数、网关发送队列、丢弃率和出口带宽,并通过指数退避避免故障后的重连风暴。


参考资料

  1. 王帅真:《直播弹幕系统设计》
    https://blog.qizong007.top/article/live-streaming-bullet-system

  2. IETF RFC 6455:The WebSocket Protocol
    https://www.rfc-editor.org/rfc/rfc6455

  3. IETF RFC 8441:Bootstrapping WebSockets with HTTP/2
    https://datatracker.ietf.org/doc/html/rfc8441

  4. Redis Documentation:ZADD
    https://redis.io/docs/latest/commands/zadd/

  5. Redis Documentation:Redis Streams
    https://redis.io/docs/latest/develop/data-types/streams/

  6. Redis Documentation:Pub/Sub
    https://redis.io/docs/latest/develop/interact/pubsub/

  7. Apache Kafka Documentation
    https://kafka.apache.org/documentation/

  8. WHATWG HTML Standard:Server-sent events
    https://html.spec.whatwg.org/multipage/server-sent-events.html

  9. Nginx Documentation:WebSocket proxying
    https://nginx.org/en/docs/http/websocket.html


总结

直播弹幕系统不是简单的“WebSocket + Redis”。

真正决定系统能否扛住热点流量的,是以下几件事:

  • 是否完成用户级广播到网关级广播的转换;
  • 是否把消息写入、顺序处理、推送和历史存储解耦;
  • 是否承认客户端不可能展示全部消息;
  • 是否针对热门房间实施独立容量和降级策略;
  • 是否给慢消费者设置明确边界;
  • 是否区分普通弹幕和重要业务消息的可靠性;
  • 是否能在节点故障和大规模重连时保持稳定;
  • 是否把出口带宽和客户端渲染纳入容量设计。

一个成熟的弹幕系统,并不是保证每一条“666”都以最高等级的可靠性抵达世界尽头,而是在流量洪峰中仍能保证直播可用、重要消息可靠、普通弹幕足够实时,并让系统有能力持续演进。


从 MVP 到百万并发:直播弹幕系统完整设计
https://allendericdalexander.github.io/2026/07/29/archtect/design/live-streaming-danmaku-system-design/
作者
AtLuoFu
发布于
2026年7月29日
许可协议