Skip to content

弹幕系统

弹幕(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)同一实例短时间重复拉取(播放器轮询)
RedisZSet 按时间范围无(数据持续累积)跨实例共享、新观众首拉
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 冷热分离(历史弹幕归档)。

基于 MIT License 发布