24. Design Bullet Screen System 弹幕系统
本节将介绍弹幕系统的设计,包括其架构、关键组件以及实现细节。
aka Danmaku System
00 Resources
01 Background 背景
用户通过发送弹幕、送礼等,可以实时在直播画面上展现自己的想法、评论和互动内容,从而丰富了用户观看体验。可以类比的是:
- 一种 "群" 消息,将某个用户 (也有系统) 产生的消息实时广播给全体 "群成员"
- 这个群动辄数百万成员
- 这个群是临时的 (进了某个直播间,一会儿又退出了)
- 消息基于直播间,进直播间才能收到消息
- 这种群消息可分布在多端。APP、Web、H5
- 这种消息阅后即焚,因为和视频画面呼应才有意义 (当然服务端会存档)
简单点讲,弹幕消息就是一种群消息,所以技术界是将弹幕归类为 IM(即时通讯) 范畴,它背后是实时消息。弹幕是实时消息中能被用户看到的内容,还有一部分是系统消息,用来主动触发业务逻辑、行为等。
在直播界,弹幕是一个举足轻重的角色,作为视觉和信息的直接载体,它身上承接了极重的业务形态:
- 广播出和画面内容呼应的内容
- 弹幕内容根据用户身份做特殊呈现 (给了钱的更显眼)
- 弹幕触发其他业务,例如抽奖
- 内容、频率要做管控
type Bullet struct {
ID string // 弹幕唯一ID,用于去重/追踪
UserId int // 用户ID
Nickname string // 用户昵称(展示用,可缓存)
Content string // 文本内容
Timestamp int64 // 绝对时间戳(发出时的毫秒)
Offset int // 相对直播/视频的时间(秒或毫秒,方便回放)
Extra *Extra // 样式 & 效果(颜色、位置、动画等)
// 扩展
RoomId int // 所属直播间/视频ID
Type string // 弹幕类型: text/gift/like/emoji/system
FontSize int // 字体大小
Color string // 颜色值,如 "#FF0000"
Position string // 位置: top/bottom/scroll
Device string // 发送设备信息 (web/ios/android)
IsVip bool // 是否会员/大航海
IsAdmin bool // 是否管理员/房管
IsPinned bool // 是否置顶/特殊弹幕
}
1.1 挑战和需求
直播弹幕是一个读写 QPS 要求都很高,假设一个直播间有 100w 用户同时在线观看,假设弹幕的提交频率为有 10000条/秒,那么需要每秒同时推送给在线用户的次数为 100w * 10000。由此可见,读请求的吞吐量需要远大于写请求,这点类似于 IM 实时聊天。
架构设计考虑以下几个场景:
- 支持直播弹幕回放 -> 意味着需要持久化 DB
- 用户进入直播间可以推送最新几秒的弹幕数据 -> 意味着需要多级缓存
- 长连模式和短连模式可以做降级切换 -> 保证高可用
- 尽可能地减少消息体的大小,节省带宽 -> 批请求 + 消息体压缩
其实这也是 IM 的主要挑战。
02 系统架构
2.1 MVP 版本
为了不影响读写的性能,采用读写分离架构。
-
写服务
- 若不考虑历史弹幕可回放,可以直接使用 Redis 作为唯一存储。
- 若考虑支持弹幕的回放,数据还是需要持久化,可以考虑使用 MySQL 或者 TiDB,暂且认为写入不是较大的瓶颈。
- 如果有更高性能的写需求,HBase、OpenTSDB 等都可以解决问题。
-
读服务
- Redis 主要用于读缓存,缓存直播间最新的弹幕数据,采用直播间 ID 作为 Key。
- 系统读服务最大 QPS = Redis 集群QPS。
Redis 存储结构选择 -> SortedSet
- 提交弹幕:
ZADD. score 设置为时间戳。进一步优化可以只存储时间的 delta 值,减少数据存储量。 - 弹幕查询:
ZRANGEBYSCORE定时轮询弹幕数据。
2.2 缓存优化
如果能让最新的实时弹幕数据都能命中本地缓存,那性能是最高的,同时大幅度降低了 Redis 的读取压力。所以弹幕读服务可以每秒轮询 Redis 数据,构建本地缓存。
热点问题
- 假设同时在线的直播间有 10000 个,读服务机器有 50 台,那么每秒轮询 Redis 的 QPS = 10000 * 50 = 50w,读取请求线性膨胀。
- 本地内存的使用量也随直播间的数量增长而膨胀,每个直播间的缓存的数据量降低,导致本地缓存的命中率降低,容易导致 GC 频繁。
2.3 热点优化
如何降低本地缓存的使用量?
- 因为火爆的直播间会占据整个平台大部分的流量,可以只针对火爆的直播间开启本地缓存。
- 通过 路由控制 同一个直播间的请求分发到固定的几台机器,例如一致性 Hash 算法。通过减少读服务机器上的直播间数量,达到降低本地缓存使用量的目的。
2.4 客户端长连接推送
为了保障客户端消息的推送性能和实时性,长连接基本是必备的,最新的消息可以直接采用长连接实时推送。
- Push Server 从 Redis 中获取用户和直播间的订阅关系以及长连接信息。
- 连接代理只负责与客户端保持长连接。
- 海量的消息推送需要批量压缩。
04 Deep Dive - 本地缓存技术方案
4.1 方案目标
通过读服务节点本地内存缓存热点直播间最新弹幕,减少分布式缓存(Redis)轮询压力,降低网络开销,提升弹幕查询响应速度,支撑高并发读场景。核心平衡 性能、内存成本 、数据实时性 三大维度。
4.2 方案设计 - 缓存数据模型
本地缓存无需存储直播间全量弹幕(全量依赖 Redis/DB),仅需缓存最近 N 条实时弹幕(如最近 10 秒、50 条,可按业务配置),且需剔除冗余字段以节省内存。
// 本地缓存用的轻量弹幕结构体(剔除存储/追溯字段,保留展示核心信息)
type LocalCacheBullet struct {
Content string // 核心文本内容
Nickname string // 展示昵称(已脱敏/格式化)
Timestamp int64 // 发送时间戳(用于排序/过期)
// 样式字段按需保留,避免大字段(如Color可存16进制短码,而非完整字符串)
Color uint32 // 颜色用uint32存储(如0xFF0000替代"#FF0000")
FontSize int8 // 字体大小用int8(仅几种可选值)
Position int8 // 位置用枚举值(0:scroll,1:top,2:bottom)
}
4.3 方案设计 - 本地缓存结构
本地缓存需支持快速插入、范围查询、过期清理,且需保证多线程安全 -> 读服务可能多协程处理用户请求。
工程首选: 环形缓冲区 -> 平衡性能与内存
- 用固定长度的切片(如cap=50)存储弹幕,尾部插入新数据,满了则覆盖头部旧数据
- 用
RWMutex做读写锁(读多写少场景,读锁不互斥,性能更优) - 额外记录count(实际存储条数)和latestTS(最新弹幕时间戳),方便快速判断是否需要更新
- 插入 / 查询
O(1),内存连续,GC 压力小
4.4 数据同步 - Pull Model
本地缓存的数据来源于分布式缓存 Redis,同步机制是核心,需解 决 "轮询压力、数据一致性、更新效率" 问题。
4.4.1 同步模式: "主动拉取 + 批量合并" 替代 "单直播间轮询"
针对 “10000 个直播间 + 50 台机器 = 50w Redis QPS” 的轮询膨胀问题,采用 直播间分片 + 批量拉取 策略
- 直播间分片绑定
- 给读服务集群的每个节点分配固定的 “直播间分片”(如按
RoomId % 机器数分片),每个节点仅负责自己分片内的直播间缓存更新. - 例:50 台机器,RoomId=1001 → 1001%50=1 → 仅由节点 1 负责同步该直播间的缓存.
- 优势:轮询 QPS 从
直播间数×机器数降至直播间数,彻底解决轮询膨胀.
- 给读服务集群的每个节点分配固定的 “直播间分片”(如按
- 批量拉取 Redis 数据
- 每个节点按分片聚合直播间,批量执行 Redis 的
ZRANGEBYSCORE(SortedSet 按时间戳查最新数据),而非单直播间单独查询; - 例:节点 1 负责 100 个直播间,每 1 秒执行 100 个
ZRANGEBYSCORE key latestTS +inf(仅查上次同步后的新数据),而非 100 次单独请求。 - 优化:使用 Redis Pipeline 批量发送命令,进一步减少网络往返开销。
- 每个节点按分片聚合直播间,批量执行 Redis 的
4.4.2 同步频率 - 动态调整 避免无效轮询
同步频率并非固定 1 秒,需结合直播间热度动态调整。读服务节点维护 直播间热度表(记录每个直播间的在线人数 / 弹幕频率),每 5 秒更新一次热度等级
- 热点直播间 如
10w +在线:- 同步频率 100ms~500ms(保证实时性)
- 中冷直播间
1000~10w在线:- 同步频率 1s~2s
- 冷直播间
<1000在线- 同步频率
5s~10s,甚至 “有新弹幕才触发更新” - 通过 Redis 的KEYSPACE_NOTIFY事件优化
- 同步频率
4.4.3 数据合并 - 去重 + 增量更新
增量更新的核心逻辑是:仅从 Redis 拉取 “上次同步后新增的弹幕”,并过滤掉重复数据,最终只将真正的新弹幕写入本地缓存。整个流程可拆解为 增量拉取 -> 精准去重 -> 安全写入 三个核心步骤,以下是工程级具体实现细节。
为实现增量更新,每个直播间的本地缓存需维护 3 个核心状态
// 单个直播间的本地缓存实体
type RoomLocalCache struct {
Buffer *RingBuffer // 环形缓冲区(存弹幕数据)
DuplicateMap map[string]bool // 去重索引(key:弹幕唯一ID)
LastSyncTS int64 // 上次从Redis拉取数据的最大时间戳
RWMutex sync.RWMutex // 读写锁(保护上述字段)
}
同步协程针对分片内的每个直播间,按以下规则构造 Redis 查询参数
- Key: 直播间弹幕的 Redis 键(如bullet:room:1001)
- MinScore: LastSyncTS + 1(仅拉取比上次更新更新的弹幕,避免重复拉取)
- MaxScore: 当前时间戳(拉取截止到此刻的所有新弹幕)
- Limit: 0, 100(单次拉取上限 100 条,避免单批数据过大阻塞同步)
// 伪代码:批量拉取某分片内的直播间新弹幕
pipeline := redisClient.Pipeline()
for _, roomCache := range shardRoomCaches {
roomCache.RLock() // 读锁,不阻塞用户查询
lastTS := roomCache.LastSyncTS
roomCache.RUnlock()
// 构造命令:ZRANGEBYSCORE key min max WITHSCORES LIMIT 0 100
cmd := pipeline.ZRangeByScoreWithScores(
fmt.Sprintf("bullet:room:%d", roomId),
&redis.ZRangeBy{
Min: strconv.FormatInt(lastTS+1, 10),
Max: strconv.FormatInt(time.Now().UnixMilli(), 10),
Count: 100,
},
)
// 绑定房间ID与命令结果,后续处理
pendingCmds[roomId] = cmd
}
// 执行批量查询
_, err := pipeline.Exec()
解析 Redis 返回结果,转换为 “弹幕 ID + 时间戳 + 原始数据” 的结构化列表。
针对拉取到的弹幕列表,通过 “两层校验” 过滤重复数据,确保写入缓存的弹幕唯一。
- 第一层 -> 基于 弹幕ID 去重
- 这是最核心的去重逻辑,利用
DuplicateMap快速判断
- 这是最核心的去重逻辑,利用
- 第二层 -> 极端场景兜底去重
- 若弹幕 ID 生成逻辑存在漏洞(如 ID 重复),可增加
UserId + Timestamp兜底 - 组合key 精确到毫秒,避免同一用户同一时刻发多条的重复
- 若弹幕 ID 生成逻辑存在漏洞(如 ID 重复),可增加
// 伪代码:去重逻辑
var newBullets []LocalCacheBullet
var maxTS int64 = roomCache.LastSyncTS
for _, bullet := range redisBullets {
// 1. 校验弹幕ID是否已在本地缓存中
if roomCache.DuplicateMap[bullet.ID] {
continue // 已存在,跳过
}
// 2. 记录当前批次的最大时间戳(用于更新LastSyncTS)
if bullet.Timestamp > maxTS {
maxTS = bullet.Timestamp
}
// 3. 转换为轻量本地缓存模型,加入新弹幕列表
newBullets = append(newBullets, convertToLocalModel(bullet))
}
// 生成组合唯一键(UserId+Timestamp精确到毫秒,避免同一用户同一时刻发多条的重复)
comboKey := fmt.Sprintf("%d_%d", bullet.UserId, bullet.Timestamp)
if roomCache.DuplicateMap[bullet.ID] || roomCache.DuplicateMap[comboKey] {
continue
}
// 同时将ID和组合键存入去重索引(双重保障)
roomCache.DuplicateMap[bullet.ID] = true
roomCache.DuplicateMap[comboKey] = true
仅将去重后的新弹幕写入环形缓冲区,并同步更新 LastSyncTS 和 DuplicateMap,保证状态一致性。
- 1.加写锁保护写入
- 由于同步协程和用户查询协程可能同时操作缓存,需加写锁确保原子性
- 2.写入环形缓冲区
- 将
newBullets列表按时间戳顺序插入环形缓冲区(若 Redis 返回无序,需先排序)
- 将
- 3.更新同步状态
- 将
LastSyncTS更新为当前批次的最大时间戳,确保下次拉取从正确的起点开始
- 将
- 4.清理过期去重索引 -> 防内存泄漏
DuplicateMap会随弹幕增多而膨胀,需同步清理 “已从环形缓冲区淘汰的旧弹幕 ID”
// 方法1:基于时间窗口清理(推荐)
expireTS := time.Now().UnixMilli() - 10*1000 // 清理10秒前的旧ID
for id := range roomCache.DuplicateMap {
// 从缓冲区中查询该ID对应的弹幕时间戳(或在DuplicateMap存ID→TS的映射)
if bulletTS, exists := roomCache.IdToTSMap[id]; exists && bulletTS < expireTS {
delete(roomCache.DuplicateMap, id)
delete(roomCache.IdToTSMap, id)
}
}
// 方法2:基于缓冲区大小清理(简单)
if len(roomCache.DuplicateMap) > roomCache.Buffer.Capacity()*2 {
// 当索引大小超过缓冲区容量2倍时,清空并重建(依赖缓冲区中的弹幕ID)
newDuplicateMap := make(map[string]bool, roomCache.Buffer.Capacity())
for _, bullet := range roomCache.Buffer.All() {
newDuplicateMap[bullet.ID] = true
}
roomCache.DuplicateMap = newDuplicateMap
}
4.4.4 异常处理与边界场景
-
网络重试导致的重复拉取
- 若 Redis 查询超时后重试,可能导致同一批弹幕被多次拉取。
- 此时
DuplicateMap会直接过滤重复 ID,且LastSyncTS未更新(因拉取的弹幕时间戳≤当前LastSyncTS),不会重复写入 。
-
直播间热度突变(冷→热)
- 当冷流直播间突然变为热点(如开播),
LastSyncTS可能滞后较久,首次拉取会返回大量历史弹幕。此时: - 单次拉取按
Limit 0 100分批处理,避免阻塞; - 环形缓冲区会覆盖旧数据,仅保留最新 N 条,符合 “本地缓存聚焦实时” 的目标。
- 当冷流直播间突然变为热点(如开播),
-
节点重启后的初始化
- 节点重启后,
LastSyncTS重置为 0,首次拉取会获取该直播间最近的弹幕 (如 Redis 中保留的最近 100 条) - 批量写入缓冲区并重建
DuplicateMap - 后续同步恢复正常增量逻辑
- 节点重启后,
4.5 数据同步 - Push Model
在本地缓存(Pull 模型)基础上,对高热度直播间(如在线超 10 万)新增 Push 主动推送能力,减少客户端轮询开销,将弹幕实时性从 “秒级” 压缩到 “百毫秒级”,同时保留 Pull 模型对冷流直播间的成本优势。

4.5.1 拓扑与路由
- 入口聚合
- 写服务将每条弹幕写入 MQ
- 按 roomId 做一致性分区,保证同房间同分区。
- 分片消费
- 消息分发层维护
roomId -> edge shard映射 - 每个读节点仅订阅自己负责的房间分片分区。
- 消息分发层维护
- 更新服务器本地缓存
- 读节点消费到“房间帧(frame)/增量消息”后直接写入 RoomLocalCache 即 Server的本地缓存
- 随后广播到所有连接到这台 Server的客户端的 send buffer。
send buffer 定义
- 就是一个待发送消息队列,绑定在每个客户端长连接上。
- Push 线程把要发给该用户的弹幕帧放入队列,而不是立刻写 socket。
- 专门的 写线程/协程 会从队列里取数据,批量 flush(writev/io_uring) 到网络。
- 假设房间里有 10 万个用户,读节点拉到一帧新弹幕后,会遍历房间内的所有连接,将这帧数据拷贝/引用到每个连接的 send buffer 中。
- 这样每个连接都有一个 send buffer as 中间缓存队列,保证消息不会因为 socket 写阻塞而影响全局。
- flush 时把队列里多条消息合并成一个大包,减少系统调用。
4.5.2 核心链路和服务流程
1.服务端 - 热点直播间识别与标记// 读服务节点定期更新直播间热度
function 刷新直播间热度() {
遍历所有分片内的直播间:
从统计服务获取 在线人数、弹幕频率
if 在线人数 > 10万 或 弹幕频率 > 100条/秒:
标记为"Push优先"
else:
维持"Pull优先"
}
// 本地缓存新增弹幕时的处理逻辑
function 本地缓存_新增弹幕(roomID, 新弹幕):
// 1. 写入环形缓冲区(原Pull逻辑)
环形缓冲区.添加(新弹幕)
去重索引.记录(新弹幕.ID)
更新最后同步时间戳
// 2. 若为热点直播 间,触发Push
if 直播间[roomID].类型 == "Push优先":
生成推送包 = {
msgID: 全局唯一ID,
roomID: roomID,
弹幕内容: 新弹幕,
时间戳: 当前毫秒数
}
发送到内部消息队列(主题: "push_room_" + roomID)
// 推送层:处理待推送弹幕
function 推送层_消息处理():
订阅"push_room_*"主题:
收到推送包时:
查房间对应的在线用户连接列表
对每个连接:
Async 异步发推送包
记到未确认表(存msgID、推送包、重试次数0)
// 推送层:重试没收到确认的消息
function 推送层_重试未确认消息():
遍历未确认表:
if 消息超时>5秒 且 重试<3次:
重发消息,重试次数+1
else if 重试≥3次:
删掉这条记录(放弃)
// 推送层:收到客户端确认
function 推送层_接收Ack(msgID):
从未确认表删掉该msgID
// 客户端接收推送
function 客户端_处理推送(推送包):
// 1. 去重
if 本地去重集合.包含(推送包.msgID):
return
本地去重集合.添加(推送包.msgID)
// 2. 加入渲染队列(批量渲染优化)
渲染队列.添加(推送包.弹幕内容)
if 渲染队列.长度 >= 20 或 距上次渲染已过50ms:
批量渲染弹幕(渲染队列)
清空渲染队列
// 3. 回复Ack
发送Ack到服务端(msgID)
4.6 内存管控 - 避免 OOM 与 GC 风暴
本地缓存本质是 “用内存换性能”,若不加管控会导致内存泄漏或 GC 频繁,需从 “容量、过期、淘汰” 三方面限制。
- 1.单节点缓存总量上限
- 按机器内存配置全局上限(如单节点内存 8G,分配 1G 给本地缓存);
- 按直播间分配内存配额:热点直播间配额 50 条,冷直播间配额 20 条,避免单个直播间占用过多内存;
- 维护 “缓存内存计数器”,每插入一条弹幕累加内存占用(按结构体大小估算),达到上限时触发全局淘汰。
- 2.数据过期策略: 时间窗口 + 惰性删除
- 主动过期 -> 同步线程每轮询时,清理环形缓冲区中 “超过缓存窗口” 的旧数据(如缓存窗口 10 秒,删除
Timestamp < 当前时间-10s的弹幕); - 惰性删除 -> 用户请求查询弹幕时,先过滤掉过期数据再返回,避免返回无效数据。
- 主动过期 -> 同步线程每轮询时,清理环形缓冲区中 “超过缓存窗口” 的旧数据(如缓存窗口 10 秒,删除
- 3.淘汰策略:针对冷数据优先释放
- 当内存达到上限,采用 LRU (最近最少使用) + 热度加权 淘汰
- 为每个直播间的缓存记录
lastAccessTime(最后一次被用户查询的时间) - 淘汰时优先选择 “lastAccessTime最早、热度最低” 的直播间缓存,清空其环形缓冲区和去重索引
- 优势:保留热点数据,释放冷数据内存,最大化缓存命中率