客户是做在线培训的,每周几场直播,要求"直播结束后半小时内能看到回放"。
原来的流程:直播结束 → 运营去后台下载录制文件 → 用转码工具转 → 上传到 CDN → 手动发通知。平均两小时,而且经常出问题(录制的 playlists 不全、转码参数选错、忘了通知)。
我们做了条自动化流水线:录制结束事件触发 → 修复录制文件 → 转码多档 → 生成 HLS → 抽封面 → 发布 → 通知。现在平均 18 分钟,失败的环节会自动重试并告警。
这篇写完整设计。
TL;DR:核心是事件驱动 + 状态机。流水线:
录制结束 → 修复(remux) → 转码多档 → HLS 打包 → 封面 → 发布 → 通知。三个最容易出问题的地方:直播录制的 MP4 常常没有完整索引(因为录制进程被中断),必须先 remux 修复;HLS 切片要关键帧对齐(-g+-sc_threshold 0+-hls_time是 GOP 整数倍);每一步都要有状态和重试(不能一个脚本跑到底,失败了不知道在哪一步)。另外:录制时优先录 TS/fMP4 分片而不是单个 MP4,能从根上避免索引问题。
目录
- 一、为什么必须自动化
- 二、流水线整体设计
- 三、第一步:拿到能用的录制文件
- 四、修复:remux 与时间戳重生
- 五、转码多档(ABR)
- 六、HLS 打包
- 七、封面与章节
- 八、发布与通知
- 九、状态机与可观测
- 十、坑清单
一、为什么必须自动化
时效性就是点播的价值——直播结束后半小时内的观看量占了回放总观看量的相当比例(观众的记忆和兴趣还在)。两小时后才上线,价值大打折扣。
人工流程的问题:
| 问题 | 后果 |
|---|---|
| 慢 | 两小时 |
| 依赖人 | 晚上/周末的直播没人处理 |
| 易错 | 参数选错、忘了上传、忘了通知 |
| 不可观测 | 失败了不知道,等用户问 |
自动化的收益不只是"快",更是"可靠"和"可预期"。
二、流水线整体设计
直播结束(录制进程退出 / 收到结束回调)
↓
[事件] recording.finished
↓
① REMUX 修复录制文件(重建索引、重整时间戳)
↓
② PROBE 提取元信息(时长、分辨率、编码、是否有音轨)
↓
③ TRANSCODE 转码多档(1080p / 720p / 480p)
↓
④ PACKAGE 生成 HLS(fMP4 切片 + master playlist)
↓
⑤ THUMBNAIL 抽封面(自动选帧)
↓
⑥ UPLOAD 上传到对象存储 / CDN
↓
⑦ PUBLISH 更新数据库状态、发布
↓
⑧ NOTIFY 回调 / 通知
↓
DONE
用 Celery 的链式任务或者一个状态机任务(前面 Celery 那篇讲过任务编排)。我用状态机而不是 chain,因为每一步失败后可以单独重试。
三、第一步:拿到能用的录制文件
录制方式的选择很重要:
| 方式 | 产出 | 问题 |
|---|---|---|
| 录成单个 MP4 | recording.mp4 |
中断时没有 moov,文件损坏 |
| 录成 TS 分片 | seg_001.ts, ... |
分片多,但每个分片独立完整 |
| 录成 fMP4 分片 | init.mp4 + seg_*.m4s |
同上,且更适合 HLS |
| 直接用 HLS 录制 | .m3u8 + .ts |
本身就是分发格式 |
我的建议:录制阶段就用分片格式(-f segment),因为它对"录制中断"免疫:
ffmpeg -i rtmp://live.example.com/stream \
-c copy \
-f segment -segment_time 600 -segment_format mpegts \
-strftime 1 "records/%Y%m%d_%H%M%S.ts"
如果已经录成了 MP4(很多平台是这样),下一步的修复就必不可少。
四、修复:remux 与时间戳重生
症状:
[mov,mp4,m4a,3gp,3g2,mj2 @ 0x...] moov atom not found
这是因为录制进程被 kill 时,ffmpeg 没能写 moov(索引)。前面那篇 moov 修复的文章讲过原理。
修复命令:
# 1. 尝试忽略错误重新封装
ffmpeg -err_detect ignore_err -i broken.mp4 -c copy -movflags +faststart fixed.mp4
# 2. 如果上面失败,强制重新生成时间戳
ffmpeg -fflags +genpts+igndts -i broken.mp4 -c copy -avoid_negative_ts make_zero fixed.mp4
# 3. 极端情况:只救数据(可能丢一点结尾)
ffmpeg -i broken.mp4 -c copy -fflags +igndts -movflags +frag_keyframe+empty_moov fixed.mp4
修复后必须验证:
def validate_media(path, min_duration=60):
"""验证文件可用:能 probe、时长合理、有视频流、能解码。"""
info = probe(path)
if not info.get('streams'):
return False, '没有流'
v = next((s for s in info['streams'] if s['codec_type'] == 'video'), None)
if v is None:
return False, '没有视频流'
dur = float(info['format'].get('duration', 0))
if dur < min_duration:
return False, f'时长过短: {dur}'
# 解码验证(抽样,不整片)
err = decode_check(path, duration=30)
if err:
return False, f'解码错误: {err[:200]}'
return True, 'ok'
这一层验证是整个流水线的"安全门"——坏文件进到后面只会浪费转码算力,而且会产出废片。
五、转码多档(ABR)
ffmpeg -i fixed.mp4 \
-map 0:v -map 0:a -map 0:v -map 0:a -map 0:v -map 0:a \
-c:v libx264 -profile:v high -preset slow \
-c:a aac -ar 48000 -b:a 128k \
-b:v:0 5000k -maxrate:0 6000k -bufsize:0 10000k -s:v:0 1920x1080 \
-b:v:1 2500k -maxrate:1 3000k -bufsize:1 5000k -s:v:1 1280x720 \
-b:v:2 1000k -maxrate:2 1200k -bufsize:2 2000k -s:v:2 854x480 \
-g 120 -keyint_min 120 -sc_threshold 0 \ # 关键:固定 GOP,保证对齐
-f hls -var_stream_map "v:0,a:0 v:1,a:1 v:2,a:2" \
-hls_time 4 -hls_segment_type fmp4 \
-hls_playlist_type vod \
-master_pl_name master.m3u8 \
-hls_segment_filename 'v%v/seg_%06d.m4s' 'v%v/index.m3u8'
要点(前面 ABR 那篇详细讲过,这里只列):
-g+-keyint_min+-sc_threshold 0:固定 GOP,保证所有档位关键帧对齐;-hls_time是 GOP 时长的整数倍(30fps、-g 120 = 4 秒,hls_time=4 刚好);-var_stream_map:一次输出多档 + master playlist;-hls_segment_type fmp4:现代格式(前面 fMP4 那篇讲过)。
耗时优化:用切片并行(前面那篇讲过)——一场 2 小时的直播,单进程转码要 40 分钟,切成 8 段并行只要 6 分钟。这是把"两小时"压到"20 分钟"的关键一步。
六、HLS 打包
上面那条命令已经完成了打包。补充几个实用点:
1. 分片命名与目录:按档位分目录(v0/ v1/ v2/),方便管理和 CDN 缓存。
2. master playlist 的 BANDWIDTH 要写实际值:
def fix_bandwidth(master_path, measured):
"""把 ffmpeg 生成的 BANDWIDTH 替换成实测的峰值码率。"""
text = Path(master_path).read_text()
for idx, kbps in measured.items():
text = re.sub(rf'(name:v{idx}.*?BANDWIDTH=)(\d+)',
rf'\g<1>{int(kbps * 1000)}', text, flags=re.S)
Path(master_path).write_text(text)
(ffmpeg 写的是你给的 -b:v,实际峰值可能更高。写低了客户端会低估导致缓冲。)
3. 校验:解析一遍 master.m3u8 和每个子 playlist,确认分片都存在、时长一致。
def validate_hls(master_path):
"""校验 HLS:分片存在、数量与 EXTINF 一致。"""
master = Path(master_path).read_text()
subs = re.findall(r'^(?!#)(.+m3u8)$', master, re.M)
ok = True
for s in subs:
p = Path(master_path).parent / s
if not p.exists():
ok = False
continue
content = p.read_text()
segs = re.findall(r'^(?!#)(.+)$', content, re.M)
for seg in segs:
if not (p.parent / seg).exists():
ok = False
return ok
七、封面与章节
封面:复用自动选帧(第 30 篇)——跳过片头片尾,抽候选帧打分,选最高的。
from thumbnail import pick_best_frame
frame, stamp = pick_best_frame(fixed_path)
save_as_jpeg(frame, 'cover.jpg')
章节(对培训/会议类直播特别有用):
- 用场景检测找画面切换点(第 25 篇);
- 或者用音频静音段找"段落间隔"(演讲中的停顿往往对应换话题);
- 或者用ASR 文本关键词("下面我们讲第二章")。
直播回放最实用的是"音频停顿 + 画面变化"的组合:
def find_chapter_points(video_path, min_gap=120, silence_thresh=-40, min_silence=1.5):
"""用静音检测找章节点(演讲中的长停顿通常对应换主题)。"""
cmd = ['ffmpeg', '-v', 'error', '-i', video_path, '-af',
f'silencedetect=noise={silence_thresh}dB:d={min_silence}', '-f', 'null', '-']
p = subprocess.run(cmd, capture_output=True, text=True)
starts = [float(m) for m in re.findall(r'silence_start: ([\d.]+)', p.stderr)]
# 过滤:间隔太近的合并
points = []
for s in starts:
if not points or s - points[-1] > min_gap:
points.append(s)
return points
章节的命名:自动只能叫"第 N 节",如果能拿到 ASR 文本就能好很多(项目里那篇字幕文章讲过)。
八、发布与通知
上传:
import boto3
def upload_directory(local_dir, bucket, prefix):
s3 = boto3.client('s3', endpoint_url=S3_ENDPOINT, ...)
for p in Path(local_dir).rglob('*'):
if p.is_file():
key = f'{prefix}/{p.relative_to(local_dir).as_posix()}'
ct = guess_content_type(p)
# HLS 的 m3u8 要设正确的 Content-Type 和缓存策略
s3.upload_file(str(p), bucket, key,
ExtraArgs={'ContentType': ct,
'CacheControl': cache_policy_for(p)})
Content-Type 很关键:
| 文件 | Content-Type |
|---|---|
.m3u8 |
application/vnd.apple.mpegurl 或 audio/mpegurl |
.m4s / .mp4 |
video/mp4 |
.ts |
video/mp2t |
设错了播放器可能拒绝播放。
通知:回调客户的 Webhook + 内部群消息。
def notify(vod_id, status, extra=None):
payload = {'vod_id': vod_id, 'status': status, 'ts': int(time.time())}
if extra:
payload.update(extra)
# 回调(要有重试)
for attempt in range(3):
try:
requests.post(WEBHOOK_URL, json=payload, timeout=5)
break
except Exception:
time.sleep(2 ** attempt)
# 内部通知
send_group_message(f'点播 {vod_id} 已{status}')
九、状态机与可观测
状态定义:
STATES = ['PENDING', 'REMUX', 'PROBE', 'TRANSCODE', 'PACKAGE',
'THUMBNAIL', 'UPLOAD', 'PUBLISH', 'NOTIFY', 'DONE', 'FAILED']
def advance(vod_id, next_state):
with transaction.atomic():
v = VodJob.objects.select_for_update().get(id=vod_id)
if v.state == 'FAILED':
return
v.state = next_state
v.state_started_at = timezone.now()
v.save(update_fields=['state', 'state_started_at'])
每一步的重试:用 Celery 的 autoretry_for + retry_backoff(前面 Celery 那篇)。转码这种耗时的步骤不要整体重试(太贵),要在步骤内部做检查点——比如"已经转完的档位不重转"。
监控指标:
| 指标 | 说明 |
|---|---|
| 端到端耗时(直播结束 → 可播放) | 核心 SLA |
| 每一步的耗时分布 | 找瓶颈 |
| 失败率(按步骤) | 哪一步最容易出问题 |
| 卡住的任务数(某状态超过 N 分钟) | 必须有(否则任务卡死没人知道) |
"卡住的任务"告警特别重要——流水线里最糟糕的情况不是失败,而是静静地卡在某一步(比如上传卡住),既不成功也不失败。所以我有个定时任务扫"状态时间超过阈值"的任务并告警。
十、坑清单
- 录制直接录成单个 MP4 → 中断就没索引。用分片格式录。
- 不修复就直接转码 → 转码失败或者产出坏文件。先验证再转码。
- 忘了
-sc_threshold 0→ 场景切换插关键帧,档位不对齐,切换花屏。 hls_time不是 GOP 整数倍 → 切片时长不均。- master playlist 的 BANDWIDTH 用
-b:v的值 → 低估峰值,客户端缓冲。写实测值。 - Content-Type 设错 → 播放器不认。
- 一个脚本跑到底 → 失败了不知道在哪一步。用状态机。
- 不做卡死检测 → 任务静静卡住。扫"状态超时"。
- 转码失败就整体重跑 → 浪费。做检查点,已完成的步骤跳过。
- 不统计端到端耗时 → 不知道有没有达成 SLA。
- 上传没做并发/分片 → 大文件慢。用 multipart 上传。
- 忘记清理本地临时文件 → 磁盘被录制文件塞满。
- 通知没有重试 → 回调失败,客户不知道已经好了。
- 封面用了片头黑帧 → 难看。跳过片头。
- 章节点太碎 → 没有导航价值。最小间隔约束。
- 不考虑"直播中途断流再续" → 录制会有多段,要合并处理(时间戳要重整)。
最后说说这个项目的核心价值。
技术上没有难点——每一步都是成熟的命令和工具。价值在于把它们串成一条"可靠、可观测、可重试"的流水线。
我见过很多"半自动"的方案:写了个脚本,手动跑,出问题手动改。这类方案的问题不是"慢",而是它不可预期——你不知道这次会不会成功、要多久、失败了要怎么救。
而流水线方案的价值是:同样的输入,总能在预期时间内得到预期结果,出问题时有明确的告警和状态。这对业务来说是质变——运营可以承诺"半小时上线",而不是"看情况"。
还有一个我觉得值得强调的设计:状态机 + 卡死检测。
大部分人做流水线时会关注"失败处理",但"卡住"比"失败"更常见也更危险——失败会抛异常、会告警,而卡住是"什么都没发生"。我们第一版就遇到过一次上传卡住 6 小时(网络问题),因为没有卡死检测,直到客户问"回放呢"才发现。
所以我现在做流水线,第一个加的不是功能,是"每个状态的超时检测"。 这是性价比最高的一条保险。