消息队列与任务编排
Vistack 使用 Kafka(segmentio/kafka-go) 作为主干消息队列,配合 Redis ZSet 延迟队列实现重试与兜底。
Topic 一览
| Topic | 生产者 | 消费者 | 说明 |
|---|---|---|---|
transcode | api | worker | 转码任务(编排远程转码) |
delete_file | api | worker | 文件/视频删除(异步清理对象与元数据) |
danmaku | api | worker | 弹幕异步落库 |
comment | api | worker | 评论异步落库 / 审核 |
消费模型
go
// internal/core/kafka.go(核心逻辑)
core.StartKafkaConsumer(ctx, topic, handler)1
2
2
- 单 goroutine 顺序消费:
ReadMessage → handler → Commit(at-least-once 语义); [kafka].concurrency控制每个 worker 实例启动的消费者数(默认 4);- 消费组
vistack-consumer-group共享,多 worker 实例自动分担分区。
转码任务的完整生命周期
┌─────────┐ Kafka[transcode] ┌────────┐ gRPC ProcessVideo ┌────────────┐
│ api │ ───────────────────▶ │ worker │ ────────────────────▶ │ transcoder │
└─────────┘ └────────┘ └────────────┘
│ │ │
│ ├─ 幂等校验(DB 状态 + Redis 租约) │
│ ├─ processing ──────────────────────▶ │
│ │ │ MinIO 下载→转码→上传
│ ◀────────── 返回(时长/清单/封面)───────┘
│ └─ 事务写库 → published
│
└── 失败 ──▶ Redis ZSet 延迟队列 ──▶ retry dispatcher 回投 Kafka ──▶ worker 重新消费1
2
3
4
5
6
7
8
9
10
11
2
3
4
5
6
7
8
9
10
11
重试派发器(retry dispatcher)
go
// internal/core/message_queue/transcode/retry.go
// 指数退避 + jitter:1min → 2min → 4min → … → 8h 封顶
func scheduleDelay(attempt int) time.Duration {
d := time.Duration(1<<uint(attempt-1)) * time.Minute
if d > 8*time.Hour { d = 8 * time.Hour }
return d + jitter(d/5) // 随机抖动,防止重试风暴
}1
2
3
4
5
6
7
2
3
4
5
6
7
- 每 5s 扫描 ZSet(
transcode:retry:zset),取出 score ≤ now 的任务; - 重新投递到
Kafka[transcode],投递成功才ZRem(失败保留,下个周期再试); - 尝试次数 ≤ 7(
attempts:transcode:{id}计数,TTL 24h),超过则放弃。
Watchdog(超时兜底)
go
// internal/core/message_queue/transcode/watchdog.go(leader 运行,每分钟)
1. processing 超 15 分钟且无 Redis 租约 → 重新入重试队列(消费者可能已挂)
2. pending 超 10 分钟 → 重新入重试队列(Kafka 消息可能丢失)1
2
3
2
3
幂等三件套
| 机制 | 作用 |
|---|---|
| DB 状态机 | status=completed 的任务直接跳过 |
Redis 租约 lease:transcode:{id} | SetNX 30min,同一任务同一时刻只有一个 worker 处理 |
尝试计数 attempts:transcode:{id} | 控制重试上限(7 次) |
删除任务(delete_file)
Kafka[delete_file] 消费者:
- 幂等校验(DB + 租约);
- 删除 MinIO 对象(原片 +
dash/{videoID}前缀分片 + 封面); - 事务清理 DB 元数据与关联互动数据。
配置
| 配置 | 说明 | 默认 |
|---|---|---|
[kafka].brokers | broker 地址(容器内 kafka:29092,宿主机 localhost:9092) | - |
[kafka].group_id | 消费组 | vistack-consumer-group |
[kafka].concurrency | 每个 worker 的并发消费者数 | 4 |
演进方向(roadmap)
- Kafka 并发消费:worker 内部支持有界 worker pool 并发处理,提升单实例吞吐;
- DLQ 死信队列:重试 7 次后进入死信 topic + 告警,而非静默丢弃;
- 消息契约版本化:Kafka JSON 消息增加
v版本字段,支持跨版本演进; - Outbox 模式:DB 事务与消息投递原子化,严格保证不丢消息。
