做过在线教育的大概都遇到过这个场面:学期初老师集中上传课程录像,晚上八九点更是上传高峰,一节几十分钟的视频要转出 720p、480p 好几路码率,单节转码动辄几十秒。我们最早图省事,在上传接口里同步触发转码,结果上传请求大量超时、worker 被瞬时任务压满,还时不时有视频永远停在"转码中",或者同一节被重复转了两遍、白白烧 CPU。这篇讲我们后来怎么用消息队列把这条链路重做一遍:上传和转码解耦削峰、生产消费两端保证消息不丢、任务幂等保证不重、状态机加死信队列兜住失败,附上关键代码和五个踩坑。
同步转码的死结在于,转码是 CPU 密集的慢活,却被塞进了一个要求秒级返回的上传请求里。上传高峰一来,请求线程全卡在等转码,连接池和 worker 一起被拖死。正确方向是异步:上传接口只负责把视频落存储、写一条任务记录、发一条消息,真正的转码交给独立 worker 按自己的节奏消费。
但异步不是免费的,它把一个问题换成了四个:消息可能丢(发出去 broker 没收到、worker 处理到一半挂了),消息可能重(网络抖动导致重复投递),任务可能失败(坏片、参数异常、机器宕机),大量任务同时到会把 worker 压垮。后面所有设计,都是在逐个回答这四件事。
上传侧保持极轻,不做任何重活:
def on_upload(course_id, video_key):
# 1. 先把任务落库,状态=待转码(这一步也是后面不丢消息的关键)
task_id = insert_transcode_task(course_id, video_key, status="pending")
# 2. 再发消息,带上任务主键
mq.send("transcode.todo", {"task_id": task_id})
return {"task_id": task_id, "status": "pending"}消费侧不要有多少拉多少,要按 worker 真实处理能力控制预取数量,否则任务全堆在单个 worker 内存里,等于没削峰:
# 每台 worker 同时只预取 2 个任务,处理完再拉下一批
channel.basic_qos(prefetch_count=2)
channel.basic_consume(queue="transcode.todo", on_message_callback=handle)prefetch_count 是最容易被忽略的旋钮。我们一开始用默认值,结果任务被不均匀地堆给先连上的几台 worker,那几台 CPU 跑满、其余机器却闲着。按单机能并行转码的路数设成 2 之后,负载才均匀下来。
消息丢失通常发生在两个位置,要分别堵。
生产端,发消息要开持久化加确认机制,并且更稳的做法是用上一步那张任务表做"本地消息表"兜底:消息先和业务数据落在同一个数据库事务里,再由一个后台进程扫描"待发送"的任务补发,这样即使发消息那一下 broker 抖动,任务也不会凭空消失,最多是重复发一次(重复由下一节的幂等兜住)。
消费端,必须手动 ACK,并且只在转码真正成功、状态落库之后才确认:
def handle(ch, method, properties, body):
task = json.loads(body)
try:
run_transcode(task["task_id"]) # 转码 + 状态机推进
ch.basic_ack(delivery_tag=method.delivery_tag) # 成功才确认
except Exception:
# 不 ACK,让消息重新入队或进死信,绝不能吞掉异常假装成功
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)如果用自动 ACK,消息刚被取走 broker 就认为投递完成,此时 worker 崩溃,这条任务就永远丢了——这正是早期那些"永远转码中"视频的来源。
MQ 的"至少一次"语义决定了重复投递是常态而不是异常,所以转码逻辑必须幂等。我们给每个任务建了明确状态:待转码、转码中、成功、失败,并用条件更新做状态流转的闸门:
-- 只有处于"待转码"的任务能被推进到"转码中"
UPDATE transcode_task
SET status = 'running', worker = ?
WHERE id = ? AND status = 'pending';
-- 影响行数 = 0,说明已经被别的 worker 领走或处理过,本次直接跳过重复消息进来时,状态已经不是 pending,条件更新影响 0 行,worker 就直接 ACK 跳过,不会再烧一遍 CPU。输出路径也用任务 ID 固定,重复执行写同一份目标文件,不会产生一堆重复产物。靠"状态机条件流转 + 固定输出",重复投递从故障变成了无害操作。
转码失败的原因分两类,要区别对待。一类是可恢复的,比如 worker 临时内存不足、对象存储偶发抖动,这类值得重试;另一类是不可恢复的,比如源文件本身损坏、编码参数非法,重试一百次还是失败,还会反复占着队列。
可恢复失败用指数退避重试,并设最大次数,避免无效空转:
RETRY_BACKOFF = [10, 30, 60, 180] # 第n次重试前等待秒数,逐步拉长
def handle(task):
if task["retry"] > len(RETRY_BACKOFF):
send_to_dead_letter(task) # 超过上限进死信队列,人工/告警介入
return
try:
run_transcode(task["id"])
except RecoverableError:
delay = RETRY_BACKOFF[min(task["retry"], len(RETRY_BACKOFF)-1)]
republish_with_delay(task, delay) # 延迟后重新投递,retry+1超过最大重试次数的进死信队列,单独可见、触发告警,而不是在主队列里无限循环。另外还要有一个超时回收:任务被推进"转码中"后若 worker 直接宕机,状态会卡住,需要一个定时扫描,把停留在 running 超过阈值(比如单节转码 P99 的三倍)的任务重置回 pending 重新调度。
异步任务的可靠性,说到底是在跟"丢、重、失败、压垮"这四件事反复较劲:用 MQ 加可控预取削峰,用本地消息表和手动 ACK 保证不丢,用状态机和幂等保证重复也无害,用退避重试、死信队列和超时回收兜住各种失败。这套模式并不只属于视频转码,订单履约、消息推送、报表生成这些"慢活异步化"的场景,踩的几乎是同一批坑。下次要把一个重操作从主请求里拆出去时,不妨先对着这四点问一遍自己,能少走不少弯路。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。