VOD 与转码流水线
上传(分片 + 秒传)
前端将原片分片上传到 MinIO(Multipart 直传,预签名 URL),支持断点续传;上传前计算文件哈希用于秒传去重(已存在则直接引用,不重复存储)。
上传完成后,api 的 CompleteVideoUpload 处理器:
- 写入
video_sources(原始视频上传记录)与files(通用文件表); - 投递
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 的删除消费者:
- 删除 MinIO 对象(原片 + DASH 分片 + 封面,递归清理
dash/{videoID}前缀); - 事务清理 DB 元数据(video / sources / transcodes / manifests / 互动数据)。
删除任务同样带幂等租约,多 worker 不会重复删除。
状态机
video_transcodes 状态流转:
pending ──> processing ──> completed
│ │
└── failed ←─┘(重试后回到 pending/processing)videos 状态流转:
uploading ──> processing ──> published
│
└── failed相关配置
| 配置 | 说明 | 默认 |
|---|---|---|
[kafka].concurrency | worker 并发消费者数 | 4 |
[transcoder].listen_addr | transcoder gRPC 监听 | :50051 |
[transcoder].use_etcd | 是否走 etcd 服务发现 | true |
[minio].bucket | 对象存储桶 | vistack |
