提示

返回博客列表

一台机器转不完:把 ffmpeg 任务分发到 8 台机器的分布式转码实践

客户给了个死线: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 挂了之后它的任务要能被别人捡起来)。加速比达不到线性是正常的,瓶颈通常在共享存储和网络。

目录

一、先算清楚:到底需要多少算力

在开机器之前,先做一次实测——凭感觉估算力一定会翻车。

我的测试方法:挑 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,也比手动分片强得多。队列带来的三个能力是刚需:

  1. 动态负载均衡:谁空谁干活,慢文件不会堵死一台机器;
  2. 失败可重试:任务没确认就从队列里再发放一次;
  3. 弹性:随时加机器、随时减机器。

至于 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'])

要点:

  1. BRPOPLPUSH 是原子的,多个 worker 同时抢不会拿到同一个任务;
  2. 完成后用 LREM 从 PROCESSING 删除(不是 LPOP,因为完成顺序和领取顺序不一定一致);
  3. requeue_stale 要定期跑(我用一个单独的进程每 60 秒跑一次),回收那些 worker 挂掉的任务;
  4. 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 块钱,而且人的注意力不可能刚好卡在结束那一刻。

十、坑清单

  1. 手动分片 → 负载严重不均、失败任务丢失。用队列。
  2. 在 NFS 上直接跑 ffmpeg → 慢 4 倍。拷到本地盘。
  3. concurrency 超过 CPU 核数 → 吞吐反而下降。
  4. 没处理 SIGTERM → 抢占式实例被回收时任务全丢。
  5. visibility_timeout / 超时回收设太短 → 任务还在跑就被别人领走,重复执行。
  6. 设太长 → 机器挂了之后任务要等几小时才被捞回。
  7. 没有幂等 → 重复执行的两个 worker 同时写同一个输出文件,产出损坏文件。用 .part + 原子改名。
  8. 不做失败分类 → 坏文件重试 3 次,浪费算力。
  9. 任务没有排序 → 队尾效应,最后几小时只有 1 台在干活。长任务优先。
  10. 输出文件直接写共享目录 → 并发写同名文件、stale handle。用唯一临时名 + 原子改名(os.replace)。
  11. 心跳没做 → 不知道 worker 死没死,只能干等超时。
  12. 磁盘没监控 → 某台机器本地盘写满,任务全失败。每台都要有磁盘告警。
  13. 跑完忘了关机 → 白烧钱。
  14. 代码依赖不一致 → 8 台机器上 Python/ffmpeg 版本不同,同样的命令有的成功有的失败。用容器或者统一镜像,这次我偷懒用脚本分发,代价是排查了两次"为什么只有 3 号机器会失败"(它的 ffmpeg 版本旧)。
  15. 没有基线测试 → 出问题不知道是变慢了还是本来就慢。开跑前先测单机吞吐,跑的过程中对比。

最后说说我对这件事的整体感受。

分布式转码听起来很唬人,但拆开看其实就是三件事:一个能抢的队列、一个能共享的文件存储、一批会干活的 worker。真正决定成败的不是架构有多 fancy,而是那些细节:

  • worker 挂了任务能不能被捞回来(超时回收);
  • 文件在哪读写(本地盘 vs 网络盘,差 4 倍);
  • 失败了要不要重试(分类);
  • 最后那几小时为什么只有一台在干活(任务排序)。

这四件事做对了,8 台机器能给你 6.8 倍;做错了,可能只有 3 倍,甚至跑出一堆坏文件。

还有一个我反复体会到的道理:批量任务一定要先小规模验证再全量跑。这次我先拿 50 个文件跑了一遍完整流程(包括失败、重试、回收),确认没问题才放全量。那 50 个文件花了 20 分钟,但帮我提前发现了 NFS 慢、并发设错、输出文件名冲突三个问题——如果直接跑全量,这三个问题会在第 30 小时集中爆发,那时候再改就要从头再来。

想亲手试试?用 VidDown 一键解析下载

粘贴视频链接即可解析,多平台支持、网页端即用;下载桌面客户端解锁海外平台本地解析,开通会员更享不限次下载。

顶部