Skip to content

VOD 与转码流水线

上传(分片 + 秒传)

前端将原片分片上传到 MinIO(Multipart 直传,预签名 URL),支持断点续传;上传前计算文件哈希用于秒传去重(已存在则直接引用,不重复存储)。

上传完成后,apiCompleteVideoUpload 处理器:

  1. 写入 video_sources(原始视频上传记录)与 files(通用文件表);
  2. 投递 Kafka[transcode] 任务(消息体含 video_id / transcode_id / object_key)。

转码流水线

上传分片 ──> api(CompleteVideoUpload) ──> Kafka[transcode]
  ──> worker 消费
      ├─ 幂等校验:已完成则跳过
      ├─ Redis SetNX 租约 lease:transcode:{id}(30min,防重复执行)
      ├─ 置 video_transcodes.status = processing
      └─ etcd 发现 transcoder ──> gRPC ProcessVideo
          └─ transcoder:MinIO 下载原片 → ffprobe 探测 → 抽封面
             → DASH 多档位转码(240p–4K)→ MinIO 上传
  ──> 返回 {时长 / 清单 / 封面 / 档位}
  ──> worker 事务写库:
      ├─ files(manifest 文件)
      ├─ video_transcodes(status=completed + 分辨率 + codec)
      ├─ video_manifest(protocol=dash + profiles)
      ├─ files(封面,若有)
      └─ videos(status=published + duration + cover_file_id)

关键设计

1. 幂等(防重复处理)

go
// 已完成则跳过
if tc.Status == TranscodeStatusCompleted { return nil }

// Redis 租约:同一时刻只有一个 worker 处理该任务
leaseKey := fmt.Sprintf("lease:transcode:%d", msg.TranscodeID)
ok, _ := core.Redis.SetNX(ctx, leaseKey, "1", 30*time.Minute).Result()
if !ok { return nil }   // 已被其他实例持有
defer core.Redis.Del(ctx, leaseKey)

即使 Kafka 重复投递(at-least-once),DB 状态机 + Redis 租约保证任务只执行一次。

2. 失败重试(指数退避)

失败时 markFailed

  • video_transcodes.status = failed
  • Redis INCR attempts:transcode:{id} 记录尝试次数(TTL 24h);
  • 次数 ≤ 7 时写入 Redis ZSet 延迟队列,下次投递时间 = 指数退避 + 随机 jitter:
1min → 2min → 4min → 8min → … → 8h(封顶) + jitter(0~20%)

3. Watchdog 超时兜底

每分钟扫描(由 leader 运行,避免多实例重复):

  • processing 状态且 updated_at 超过 15 分钟、且无 Redis 租约(说明消费者挂了)→ 重新入重试队列;
  • pending 状态且超过 10 分钟(说明 Kafka 消息丢失)→ 重新入重试队列并触碰 updated_at 防重复投递。

4. 远程转码(gRPC)

  • transcoder 无状态:下载 → 转码 → 上传全部走 MinIO,天然可水平扩展;
  • worker 经 etcd 发现所有 transcoder 实例,gRPC round_robin 负载均衡;
  • 单次调用超时 25 分钟(transcodeCallTimeout)。

播放(DASH ABR)

转码产出 DASH 流:manifest.mpd + init-*.m4s / chunk-*.m4s 分片(按 240p–4K 多档位)。前端 dash.js 根据带宽自适应切换清晰度。

播放鉴权:MinIO 预签名 URL / STS 临时凭证,防止盗链。

删除(异步清理)

删除视频时投递 Kafka[delete_file],worker 的删除消费者:

  1. 删除 MinIO 对象(原片 + DASH 分片 + 封面,递归清理 dash/{videoID} 前缀);
  2. 事务清理 DB 元数据(video / sources / transcodes / manifests / 互动数据)。

删除任务同样带幂等租约,多 worker 不会重复删除。

状态机

video_transcodes 状态流转:

pending ──> processing ──> completed
    │            │
    └── failed ←─┘(重试后回到 pending/processing)

videos 状态流转:

uploading ──> processing ──> published

                    └── failed

相关配置

配置说明默认
[kafka].concurrencyworker 并发消费者数4
[transcoder].listen_addrtranscoder gRPC 监听:50051
[transcoder].use_etcd是否走 etcd 服务发现true
[minio].bucket对象存储桶vistack

基于 MIT License 发布