客户给了个死线:9000 多条历史录像,五天内全部转成 H.265 归档。
我先在一台 8 核机器上测了吞吐:
libx265 -crf 27 -preset slow,1080p 素材大概是 每分钟机器时间处理 8 分钟素材(也就是 8 倍速)。这批素材总时长约 42000 分钟,算下来单机要 5250 分钟 ≈ 3.6 天不停机。看着还行?但实测里还有读取、校验、重试、失败重来,再加上机器不能 100% 跑满,实际单机预估 20 天。只能上多台。我从云上临时开了 8 台同规格机器,搭了个分布式转码集群。第一次跑的结果很糟:有的机器跑满、有的空转、还有一批文件被转了两遍。修了几轮之后,最终 8 台机器用了 62 小时跑完,加速比 6.8 倍。
这篇文章写整个过程:算力怎么估算、任务怎么分发才不会重复、文件怎么共享、失败怎么处理、以及为什么 8 台机器只快了 6.8 倍而不是 8 倍。
TL;DR:分布式转码的核心是一个所有人都能抢的队列 + 原子领取 + 失败回收。用 Redis 的
BRPOPLPUSH(可靠队列模式)或者直接用 Celery(acks_late=True+visibility_timeout)。三个关键点:任务必须幂等(重复执行不能产生坏结果);转码一定要在本地磁盘上跑(在 NFS 上直接读写会慢好几倍);要有超时回收机制(worker 挂了之后它的任务要能被别人捡起来)。加速比达不到线性是正常的,瓶颈通常在共享存储和网络。
目录
- 一、先算清楚:到底需要多少算力
- 二、三种架构,我为什么选了队列
- 三、任务队列:原子领取与超时回收
- 四、文件怎么共享:NFS 的坑与对象存储
- 五、worker 设计:并发、能力声明、优雅退出
- 六、失败处理:重试、死信与坏文件隔离
- 七、进度与可观测
- 八、实测数据:为什么只有 6.8 倍
- 九、成本:这些机器花了多少钱
- 十、坑清单
一、先算清楚:到底需要多少算力
在开机器之前,先做一次实测——凭感觉估算力一定会翻车。
我的测试方法:挑 10 个有代表性的文件(覆盖不同分辨率、时长、编码),跑一遍真实命令,记录耗时和素材时长:
#!/bin/bash
# bench_single.sh
total_src=0
total_time=0
for f in samples/*.mp4; do
dur=$(ffprobe -v error -show_entries format=duration -of csv=p=0 "$f")
start=$(date +%s.%N)
ffmpeg -y -v error -i "$f" -c:v libx265 -crf 27 -preset slow -c:a aac -b:a 128k /tmp/out.mp4
end=$(date +%s.%N)
t=$(awk "BEGIN{print $end-$start}")
total_src=$(awk "BEGIN{print $total_src+$dur}")
total_time=$(awk "BEGIN{print $total_time+$t}")
echo "$f: ${dur}s 素材 / ${t}s 耗时"
done
echo "速度倍率: $(awk "BEGIN{printf \"%.2f\", $total_src/$total_time}")x"
实测结果(8 核 i7,libx265 CRF 27 preset slow):
| 素材类型 | 速度倍率 | 说明 |
|---|---|---|
| 480p 老片 | 26x | 分辨率低,很快 |
| 720p H.264 | 11x | |
| 1080p H.264 | 7.8x | 主力素材 |
| 1080p 高运动 | 5.2x | 慢 |
| 4K(少量) | 1.4x | 极慢,个别文件几乎 1:1 |
加权平均下来约 8x。所以:
总素材时长 42000 分钟 ÷ 8 = 5250 分钟 = 87.5 小时 = 3.6 天
但这是理想值。实际要打折:
| 折扣项 | 系数 |
|---|---|
| 文件读取/校验/写回的 IO 时间 | ×1.15 |
| 失败重试(约 3% 任务会重试一次) | ×1.03 |
| 机器不能满载(留 10% 给系统) | ×1.1 |
| 调度开销、任务启动 | ×1.05 |
| 合计 | ×1.37 |
所以单机实际 ≈ 3.6 × 1.37 ≈ 5 天。要在 5 天内完成,理论上一台就够——但没有任何余量,而且我算的时候还没算上一台机器跑 5 天不出问题的概率。
最终决策:开 8 台,目标 1 天内完成,留足余量应对意外。永远给批量任务留 2 倍以上的余量,这是我吃过亏之后的习惯。
二、三种架构,我为什么选了队列
可选的路子大概三种:
| 方案 | 怎么做 | 优点 | 缺点 |
|---|---|---|---|
| 手动分片 | 把文件列表按行数切成 8 份,每台机器跑一份 | 零基础设施 | 负载严重不均(慢文件集中在某一台)、失败没人管、加机器要重切 |
| 队列抢任务 | 一个中心队列,worker 抢着干 | 自动负载均衡、加机器随时加、失败可重试 | 需要一个可靠的队列 |
| K8s / Batch 服务 | 用现成调度平台 | 功能最全(自动扩缩容、重试、监控) | 运维成本高,为了一次性任务上 K8s 不划算 |
我第一版用的是手动分片,因为它看起来最简单。结果:
- 8 台里有 2 台提前 6 小时干完(分到的都是短文件),另外 2 台超时 3 小时还没完(分到了几个 4K 文件);
- 中途一台机器被云厂商回收(我用了抢占式实例),那 1/8 的任务全部丢失,第二天才发现;
- 想加机器帮忙,得重新切分,而且不知道哪些已经做完了。
结论:永远用队列。哪怕是最简单的 Redis List,也比手动分片强得多。队列带来的三个能力是刚需:
- 动态负载均衡:谁空谁干活,慢文件不会堵死一台机器;
- 失败可重试:任务没确认就从队列里再发放一次;
- 弹性:随时加机器、随时减机器。
至于 K8s——如果这是长期服务,值得上;一次性任务,Redis 队列 + 几台机器就够了。别为了"架构好看"引入自己hold不住的复杂度。
三、任务队列:原子领取与超时回收
用 Redis 做可靠队列
最朴素的做法是 LPUSH + RPOP(或 BRPOP 阻塞版),但有个致命问题:worker 拿到任务后挂了,这个任务就永久丢失了——它已经从队列里出来了,没人知道它去哪了。
解决办法是 可靠队列模式(reliable queue):用一个"处理中"的列表做中转。
import redis
import json
import time
r = redis.Redis(host='queue-host', port=6379, db=0)
QUEUE = 'transcode:queue'
PROCESSING = 'transcode:processing'
def claim(timeout=5):
"""原子地领取一个任务:从 QUEUE 移到 PROCESSING。"""
item = r.brpoplpush(QUEUE, PROCESSING, timeout=timeout)
if item:
return json.loads(item)
return None
def ack(item):
"""任务完成,从 PROCESSING 里删掉。"""
r.lrem(PROCESSING, 1, json.dumps(item, sort_keys=True))
def requeue_stale(max_age=3600):
"""回收超时未完成任务(一般是 worker 挂了)。"""
now = time.time()
for raw in r.lrange(PROCESSING, 0, -1):
item = json.loads(raw)
if now - item['claimed_at'] > max_age:
# 移回队列(用 LREM 保证原子性,避免重复入队)
if r.lrem(PROCESSING, 1, raw):
item['claimed_at'] = 0
item['retries'] = item.get('retries', 0) + 1
r.lpush(QUEUE, json.dumps(item))
print('requeued stale task:', item['id'])
要点:
BRPOPLPUSH是原子的,多个 worker 同时抢不会拿到同一个任务;- 完成后用
LREM从 PROCESSING 删除(不是LPOP,因为完成顺序和领取顺序不一定一致); requeue_stale要定期跑(我用一个单独的进程每 60 秒跑一次),回收那些 worker 挂掉的任务;json.dumps(..., sort_keys=True):保证序列化结果一致,LREM才能匹配上(Redis 的 LREM 是按值精确匹配的)。
直接用 Celery 也行
如果项目已经在用 Celery(我们就是),其实不用自己写队列——Celery 天然支持多机 worker,只要它们连同一个 broker:
# 机器 A
celery -A proj worker -Q transcode --concurrency=8 --hostname=workerA@%h
# 机器 B、C、...
celery -A proj worker -Q transcode --concurrency=8 --hostname=workerB@%h
配合这两个配置(前面那篇 Celery 的文章讲过,这里再强调一次,因为分布式场景下尤其重要):
CELERY_TASK_ACKS_LATE = True # 执行完才确认,挂了会重新投递
CELERY_BROKER_TRANSPORT_OPTIONS = {'visibility_timeout': 7200}
CELERY_WORKER_PREFETCH_MULTIPLIER = 1 # 长任务必须关掉预取
visibility_timeout 在分布式场景下更关键:一台机器挂了,消息要等超时后才会被别的 worker 看到。设太短(默认 1 小时)会导致"任务还在跑就被别人领走了",设太长会导致"机器挂了任务要等很久才被捞回来"。我按"最长任务 × 2"设的,这次是 2 小时。
我最后还是自己写了一层轻量队列,原因是:Celery 的 worker 部署要同步代码、配置、依赖,而这次是临时机器,我只想丢一个单文件脚本过去跑。如果是长期集群,用 Celery 更省事。
四、文件怎么共享:NFS 的坑与对象存储
任务分发解决了,下一个问题是:8 台机器怎么读同一批文件、写到同一个地方。
三个方案:
| 方案 | 优点 | 缺点 |
|---|---|---|
| NFS 共享目录 | 简单,挂载就能用 | 小文件性能差、容易 stale handle、带宽瓶颈 |
| 对象存储(S3/MinIO) | 扩展性好、不需要挂载 | 要先下载到本地再处理 |
| 每台机器本地存一份 | 最快 | 要先把 9000 个文件分发到 8 台,本身就是个大工程 |
我第一版用的 NFS,因为最省事。踩的坑:
坑 1:在 NFS 上直接跑 ffmpeg 极慢。
本地 SSD: 处理 1080p 素材 7.8x 速度
NFS 挂载: 处理同样素材 2.1x 速度 ← 慢了近 4 倍
原因是 ffmpeg 的读写模式(大量随机 seek + 持续写入)在 NFS 上开销巨大。正确做法:先把文件从共享存储拷到本地临时盘,处理完再拷回去。
import shutil
import tempfile
def process(item):
with tempfile.TemporaryDirectory(dir='/local_ssd') as tmp:
src = os.path.join(tmp, 'in.mp4')
dst = os.path.join(tmp, 'out.mp4')
download(item['src_path'], src) # 从 NFS/对象存储拉到本地
run_ffmpeg(src, dst)
upload(dst, item['dst_path']) # 传回去
加了这一步之后,吞吐立刻回到 7.5x 左右(比本地慢一点点,是拷贝的开销)。
坑 2:Stale file handle。
OSError: [Errno 116] Stale file handle: '/mnt/nfs/xxx.mp4'
这是因为 NFS 客户端持有文件句柄,但服务端上文件已经被别的进程删除/替换了。在分布式转码里特别常见——比如清理任务删了个旧文件,恰好有 worker 引用着它。
解决办法:
- 挂载时加 hard(默认)并在代码里重试;
- 更根本的是:不要在多个机器间共享"会被删除"的目录。输出目录和临时目录分开,清理任务只碰明确的过期目录。
坑 3:NFS 成为带宽瓶颈。 8 台机器同时读写同一个 NFS,服务端出口带宽打满,所有机器一起变慢。这个在第八节的实测数据里会看到。
我最后的方案:对象存储(MinIO)+ 本地盘。
MinIO (源桶) → 下载到本地 SSD → ffmpeg 转码 → 上传 MinIO (目标桶) → 删本地
对象存储的好处:没有 stale handle 问题、带宽可以水平扩展、天然支持重试(HTTP 请求)。代价是多了下载/上传两个步骤(约占总时间的 8~12%)。
如果你们的量不大、机器在同一个机房,用 NFS + 本地盘中转也完全可以。我只是提醒:别在 NFS 上直接跑 ffmpeg。
五、worker 设计:并发、能力声明、优雅退出
worker 脚本的核心循环:
#!/usr/bin/env python3
import os
import signal
import socket
import subprocess
import time
import json
import redis
HOST = socket.gethostname()
CONCURRENCY = int(os.getenv('WORKER_CONCURRENCY', '8'))
CAPABILITY = os.getenv('WORKER_CAP', 'cpu') # cpu / gpu
_running = True
def shutdown(signum, frame):
"""收到 SIGTERM:不再领新任务,把手上的干完再退出。"""
global _running
print(f'[{HOST}] received SIGTERM, finishing current task...')
_running = False
signal.signal(signal.SIGTERM, shutdown)
signal.signal(signal.SIGINT, shutdown)
def heartbeat(r, item):
"""上报心跳,让别人知道这个任务还活着。"""
r.setex(f'task:hb:{item["id"]}', 300, json.dumps({
'host': HOST, 'ts': time.time(), 'cap': CAPABILITY,
}))
def main():
r = redis.Redis(...)
while _running:
item = claim(timeout=5)
if not item:
continue
try:
hb_stop = start_heartbeat_thread(r, item) # 后台线程每 30 秒上报
process(item)
ack(item)
except Exception as e:
fail(item, str(e))
finally:
hb_stop.set()
print(f'[{HOST}] exited gracefully')
if __name__ == '__main__':
main()
几个设计点:
1. 优雅退出(SIGTERM 处理)
云厂商回收抢占式实例时会先发 SIGTERM(通常 30 秒到 2 分钟),所以必须处理 SIGTERM:停止领新任务,把当前任务做完(或者快速失败让它重新入队)。
不做的话:机器被回收 → 任务在 PROCESSING 里躺到超时才被回收 → 白白浪费一小时。
2. 能力声明(CPU / GPU)
不同机器的能力不同(有的带 GPU 能做 NVENC,有的只有 CPU)。用不同队列区分最简单:
QUEUE = f'transcode:{CAPABILITY}'
提交任务时按需求投到对应队列。不要试图在 worker 里做复杂的调度决策,队列划分是最简单可靠的。
3. 并发数
CPU 转码:concurrency = CPU 核数。我实测过 8 核机器上设 8 和 16:
| 并发数 | 总吞吐(文件/小时) |
|---|---|
| 4 | 21 |
| 8 | 34 |
| 12 | 35 |
| 16 | 32 |
超过核数之后吞吐反而下降(上下文切换开销)。所以别贪心,concurrency = nproc 就对了。
4. 每个 worker 一个独立临时目录
避免多个进程写同一个临时文件名。用 hostname + pid 做目录名:
tmpdir = tempfile.mkdtemp(prefix=f'{HOST}_{os.getpid()}_', dir='/local_ssd')
六、失败处理:重试、死信与坏文件隔离
跑 9000 个文件,一定会有失败的。我的失败分类和处理:
| 失败类型 | 占比(实测) | 处理 |
|---|---|---|
| 源文件损坏(解不开) | 1.8% | 不重试,直接标记 P0,移入隔离目录 |
| 超时(4K 素材估时不足) | 0.9% | 重试,降低 preset(slow → medium) |
| 网络错误(对象存储超时) | 0.7% | 重试,指数退避 |
| 磁盘满 | 0.3% | 重试(清理后可能恢复) |
| 未知错误 | 0.2% | 重试 2 次,之后进死信队列 |
MAX_RETRY = 2
def fail(item, reason):
retries = item.get('retries', 0)
kind = classify(reason)
if kind == 'corrupted' or retries >= MAX_RETRY:
# 进死信队列,人工处理
r.lpush('transcode:dead', json.dumps({**item, 'reason': reason}))
r.lrem(PROCESSING, 1, json.dumps(item, sort_keys=True))
return
item['retries'] = retries + 1
# 降低 preset 再试(超时场景下有效)
if kind == 'timeout' and item.get('preset') == 'slow':
item['preset'] = 'medium'
r.lpush(QUEUE, json.dumps(item))
r.lrem(PROCESSING, 1, json.dumps(item, sort_keys=True))
关键设计:区分"可重试"和"不可重试"。
源文件损坏这种,重试一万次也是失败,还会浪费算力。必须在失败分类上花点时间——我第一版没分类,所有失败都重试,结果 160 多个坏文件各重试了 3 次,浪费了几小时。
分类靠报错信息:
CORRUPT_PATTERNS = [
'Invalid data found when processing input',
'moov atom not found',
'error while decoding MB',
'Could not find codec parameters',
]
def classify(reason: str) -> str:
for p in CORRUPT_PATTERNS:
if p in reason:
return 'corrupted'
if 'timeout' in reason.lower() or 'Timeout' in reason:
return 'timeout'
if any(k in reason for k in ('ConnectionError', 'TimeoutError', 'S3')):
return 'network'
if 'No space left' in reason:
return 'disk'
return 'unknown'
死信队列要有人看。我每天看两次,把坏文件清单整理出来给客户——这批文件本来就要单独处理(那篇 ffprobe 质检的文章里讲过怎么提前发现它们)。
七、进度与可观测
跑 60 小时的任务,没有进度显示会让人焦虑。我做了个简单的进度脚本:
import redis
import time
r = redis.Redis(...)
TOTAL = 9234
while True:
pending = r.llen('transcode:queue')
processing = r.llen('transcode:processing')
dead = r.llen('transcode:dead')
done = TOTAL - pending - processing - dead
# 从心跳里统计活跃 worker
hosts = set()
for k in r.scan_iter('task:hb:*'):
hb = json.loads(r.get(k) or '{}')
if hb.get('host'):
hosts.add(hb['host'])
elapsed = time.time() - START
rate = done / elapsed * 3600 # 每小时处理数
eta = pending / rate if rate else 0
print(f'\r进度 {done}/{TOTAL} ({done/TOTAL*100:.1f}%) | '
f'待处理 {pending} | 进行中 {processing} | 失败 {dead} | '
f'worker {len(hosts)} | 速率 {rate:.0f}/h | '
f'预计剩余 {eta/3600:.1f}h', end='')
time.sleep(30)
输出:
进度 7421/9234 (80.4%) | 待处理 1683 | 进行中 24 | 失败 106 | worker 8 | 速率 121/h | 预计剩余 13.9h
"预计剩余时间"这个数字最有价值。它让你能在早期就发现问题——我有一次看到 ETA 突然从 8 小时涨到 20 小时,查下来是两台 worker 因为 NFS 报错退出了,及时重启避免了更大损失。
worker 数量的监控很重要:8 台应该一直有 8 个心跳。掉到 6 个就要立刻查。我把这条接到了告警(低于 7 个就发消息)。
八、实测数据:为什么只有 6.8 倍
最终结果:
| 配置 | 总耗时 | 加速比 | 效率 |
|---|---|---|---|
| 1 台(本地 SSD) | 423 小时(预估) | 1.0x | 100% |
| 8 台(NFS 直接读写) | 118 小时 | 3.6x | 45% |
| 8 台(本地盘中转 + NFS) | 71 小时 | 6.0x | 75% |
| 8 台(本地盘中转 + MinIO) | 62 小时 | 6.8x | 85% |
第一版只有 3.6 倍,原因就是前面说的在 NFS 上直接跑 ffmpeg。改成"下载到本地 SSD → 处理 → 传回"之后跳到 6.0 倍。再换成 MinIO 摆脱 NFS 的带宽瓶颈,到 6.8 倍。
剩下的 1.2 倍损失在哪:
| 损耗来源 | 占比 |
|---|---|
| 任务启动开销(每个任务约 2~4 秒) | 3% |
| 下载/上传文件的时间 | 9% |
| 队尾效应(最后几个长任务,其他机器空转) | 3% |
| 对象存储带宽争抢 | 2% |
"队尾效应"是分布式的固有损耗:跑到最后,队列里只剩几个 4K 大文件,7 台机器空着等 1 台干完。缓解办法是把大任务排在前面(长任务优先调度,LPT 算法):
# 提交时按预估时长降序入队(长任务先跑)
tasks.sort(key=lambda t: -t['estimated_duration'])
我第二版加了这个优化,队尾效应从 8 小时降到 2 小时。这是个很划算的优化,一行排序的事。
九、成本:这些机器花了多少钱
这次用的是云上的抢占式实例(spot),8 核 16G,跑了 62 小时。
| 方案 | 单价 | 总成本 |
|---|---|---|
| 按量付费 | 1.2 元/小时 | 8 × 62 × 1.2 = 595 元 |
| 抢占式实例 | 0.36 元/小时 | 8 × 62 × 0.36 = 179 元 |
| 包月(1 个月) | 380 元/台 | 3040 元(用不满,浪费) |
抢占式便宜 70%,代价是可能被回收(我这次 62 小时里被回收了 2 次,都是优雅处理了,只损失了两个任务的重做时间)。
用抢占式的前提是你的任务能容忍中断——也就是前面说的:SIGTERM 处理 + 任务可重试 + 幂等。如果这三件事没做好,抢占式会给你带来灾难;做好了,就是白捡 70% 的折扣。
另外一个省钱点:跑完立刻释放。我写了个检查脚本,队列空了就自动关机:
#!/bin/bash
# auto_shutdown.sh —— 队列空了就关掉自己(省最后一小时的钱)
while true; do
pending=$(redis-cli -h queue-host llen transcode:queue)
processing=$(redis-cli -h queue-host llen transcode:processing)
if [ "$pending" -eq 0 ] && [ "$processing" -eq 0 ]; then
echo "queue empty, shutting down"
sudo shutdown -h now
fi
sleep 60
done
别小看这个——8 台机器晚关 2 小时就是 6 块钱,而且人的注意力不可能刚好卡在结束那一刻。
十、坑清单
- 手动分片 → 负载严重不均、失败任务丢失。用队列。
- 在 NFS 上直接跑 ffmpeg → 慢 4 倍。拷到本地盘。
concurrency超过 CPU 核数 → 吞吐反而下降。- 没处理 SIGTERM → 抢占式实例被回收时任务全丢。
visibility_timeout/ 超时回收设太短 → 任务还在跑就被别人领走,重复执行。- 设太长 → 机器挂了之后任务要等几小时才被捞回。
- 没有幂等 → 重复执行的两个 worker 同时写同一个输出文件,产出损坏文件。用
.part+ 原子改名。 - 不做失败分类 → 坏文件重试 3 次,浪费算力。
- 任务没有排序 → 队尾效应,最后几小时只有 1 台在干活。长任务优先。
- 输出文件直接写共享目录 → 并发写同名文件、stale handle。用唯一临时名 + 原子改名(
os.replace)。 - 心跳没做 → 不知道 worker 死没死,只能干等超时。
- 磁盘没监控 → 某台机器本地盘写满,任务全失败。每台都要有磁盘告警。
- 跑完忘了关机 → 白烧钱。
- 代码依赖不一致 → 8 台机器上 Python/ffmpeg 版本不同,同样的命令有的成功有的失败。用容器或者统一镜像,这次我偷懒用脚本分发,代价是排查了两次"为什么只有 3 号机器会失败"(它的 ffmpeg 版本旧)。
- 没有基线测试 → 出问题不知道是变慢了还是本来就慢。开跑前先测单机吞吐,跑的过程中对比。
最后说说我对这件事的整体感受。
分布式转码听起来很唬人,但拆开看其实就是三件事:一个能抢的队列、一个能共享的文件存储、一批会干活的 worker。真正决定成败的不是架构有多 fancy,而是那些细节:
- worker 挂了任务能不能被捞回来(超时回收);
- 文件在哪读写(本地盘 vs 网络盘,差 4 倍);
- 失败了要不要重试(分类);
- 最后那几小时为什么只有一台在干活(任务排序)。
这四件事做对了,8 台机器能给你 6.8 倍;做错了,可能只有 3 倍,甚至跑出一堆坏文件。
还有一个我反复体会到的道理:批量任务一定要先小规模验证再全量跑。这次我先拿 50 个文件跑了一遍完整流程(包括失败、重试、回收),确认没问题才放全量。那 50 个文件花了 20 分钟,但帮我提前发现了 NFS 慢、并发设错、输出文件名冲突三个问题——如果直接跑全量,这三个问题会在第 30 小时集中爆发,那时候再改就要从头再来。