Skip to content

消息队列与任务编排

Vistack 使用 Kafka(segmentio/kafka-go) 作为主干消息队列,配合 Redis ZSet 延迟队列实现重试与兜底。

Topic 一览

Topic生产者消费者说明
transcodeapiworker转码任务(编排远程转码)
delete_fileapiworker文件/视频删除(异步清理对象与元数据)
danmakuapiworker弹幕异步落库
commentapiworker评论异步落库 / 审核

消费模型

go
// internal/core/kafka.go(核心逻辑)
core.StartKafkaConsumer(ctx, topic, handler)
  • 单 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 重新消费

重试派发器(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)   // 随机抖动,防止重试风暴
}
  • 每 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 消息可能丢失)

幂等三件套

机制作用
DB 状态机status=completed 的任务直接跳过
Redis 租约 lease:transcode:{id}SetNX 30min,同一任务同一时刻只有一个 worker 处理
尝试计数 attempts:transcode:{id}控制重试上限(7 次)

删除任务(delete_file)

Kafka[delete_file] 消费者:

  1. 幂等校验(DB + 租约);
  2. 删除 MinIO 对象(原片 + dash/{videoID} 前缀分片 + 封面);
  3. 事务清理 DB 元数据与关联互动数据。

配置

配置说明默认
[kafka].brokersbroker 地址(容器内 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 事务与消息投递原子化,严格保证不丢消息。

基于 MIT License 发布