弹幕系统
弹幕(Danmaku)系统针对「高并发写 + 时间轴读取」场景设计:写入实时可见(Redis ZSet)、持久化异步(Kafka 落库)、读取三级缓存(本地 LRU → Redis → DB)。
整体架构
text
┌──────────── 发送弹幕 ────────────┐
│ 1. 敏感词过滤(AC 自动机) │
│ 2. 写 Redis ZSet(实时可见) │
│ 3. 投 Kafka[danmaku](异步落库) │
└─────────────────┬────────────────┘
▼
┌──────────── 拉取弹幕(时间轴)───────────┐
│ 1. 本地 LRU(进程内,命中即返回) │
│ 2. Redis ZRangeByScore(时间范围) │
│ 3. DB 回源(兜底) │
└──────────────────────────────────────────┘发送链路
go
// internal/danmaku/danmaku.go
func (s *Service) Send(ctx context.Context, videoID, userID int64, content string, timeOffset float64, color string, mode int) (*entity.Danmaku, error) {
// 1. 敏感词过滤(命中返回 ErrSensitive,接口 4xx)
if s.filter.Contains(content) {
return nil, ErrSensitive
}
// 2. 实时写 Redis ZSet:member = 弹幕 JSON,score = 时间轴位置
s.rdb.ZAdd(ctx, danmakuKey(videoID), redis.Z{Score: timeOffset, Member: string(raw)})
// 3. 异步落库 Kafka(失败不阻断,弹幕已实时可见)
core.SendKafkaMessage(ctx, string(consts.KafkaTopicDanmaku), videoID, raw)
return d, nil
}写入即所见:弹幕先落 Redis ZSet(同视频的弹幕按时间轴排序),同一时刻的观众立即可见;落库完全异步(Kafka → worker → DB),即使落库失败也不影响实时体验。
读取链路(三级缓存)
go
func (s *Service) Fetch(ctx context.Context, videoID int64, start, end float64) ([]entity.Danmaku, error) {
// 1. 本地 LRU(key = videoID:start:end,TTL 2s)
if items, ok := s.local.Get(lkey); ok { return items, nil }
// 2. Redis ZSet 按时间范围
raws, err := s.rdb.ZRangeByScore(ctx, danmakuKey(videoID),
&redis.ZRangeBy{Min: start, Max: end}).Result()
if err == nil && len(raws) > 0 { s.local.Set(lkey, items); return items, nil }
// 3. DB 回源(兜底)
s.db.Where("video_id = ? AND time_offset BETWEEN ? AND ?", videoID, start, end).
Order("time_offset asc").Find(&list)
s.local.Set(lkey, list)
return list, nil
}| 层级 | 缓存 | TTL | 命中场景 |
|---|---|---|---|
| 本地 LRU | 进程内 local_cache_size(默认 1024)条 | local_cache_ttl(默认 2s) | 同一实例短时间重复拉取(播放器轮询) |
| Redis | ZSet 按时间范围 | 无(数据持续累积) | 跨实例共享、新观众首拉 |
| DB | 兜底回源 | - | Redis 数据丢失 / 冷门视频 |
为什么三级?
播放器通常按秒轮询弹幕,本地 LRU 吸收重复请求;Redis ZSet 天然按时间轴排序,ZRangeByScore 一次取回区间弹幕;DB 只做冷数据兜底。三级合起来把 QPS 峰值压到最低。
敏感词过滤
internal/danmaku/sensitive.go 使用 AC 自动机(Aho–Corasick) 多模式匹配:
- 启动时从 DB 加载全部敏感词(
LoadSensitiveWords)构建自动机; - 管理端可增删敏感词(
AddSensitiveWord/DeleteSensitiveWord),操作后重建自动机; - 发送弹幕 / 评论时
filter.Contains(content)快速判定,命中返回ErrSensitive。
异步落库(worker)
Kafka[danmaku] 消费者(internal/core/message_queue/danmaku/worker.go)批量写库,与主链路完全解耦。
接口与配置
接口
| 方法 | 路径 | 说明 |
|---|---|---|
| POST | /api/v1/danmaku | 发送弹幕 |
| GET | /api/v1/videos/{id}/danmaku?start=&end= | 按时间轴拉取 |
| GET/POST/DELETE | /api/v1/admin/sensitive-words | 敏感词管理 |
配置
toml
[danmaku]
enabled = true
local_cache_size = 1024 # 本地 LRU 容量
local_cache_ttl = 2 # 本地缓存 TTL(秒)
cache_control_max_age = 5 # HTTP 响应 Cache-Control演进方向
- 弹幕 WebSocket/长连接推送(当前为轮询拉取);
- 弹幕风控增强(频率限制、图灵测试);
- 弹幕 DB 冷热分离(历史弹幕归档)。
