如何正确实现Celery Chain?解决音频任务重复执行问题
解决Celery任务重复执行的问题
要让download_audio任务仅执行一次,核心是调整任务依赖结构,让音频下载任务只被调度一次,然后将其结果共享给两个视频合并任务,而非在每个合并流程中重复触发。
修改后的代码
from celery import group, chain # 定义基础下载任务 t1 = download_video.si(url=video1) t2 = download_video.si(url=video2) t3 = download_audio.si(url=audio1) # 调整合并任务参数顺序:音频路径在前,方便接收上游传递的音频结果 def merge_video_and_audio(audio_path, video_path): # 音视频合并逻辑 pass # 定义两个合并任务:先完成对应视频下载,再接收音频结果执行合并 merge_task1 = t1 | merge_video_and_audio.s() merge_task2 = t2 | merge_video_and_audio.s() # 构建工作流:先完成音频下载,再并行执行两个合并任务,最后发送报告 workflow = chain( t3, group(merge_task1, merge_task2), send_report.s() ) # 启动工作流 workflow.apply_async()
关键逻辑说明
- 单一音频任务调度:通过
chain先执行t3,确保音频下载仅触发一次,执行完成后自动将结果传递给后续任务。 - 并行合并执行:
group(merge_task1, merge_task2)会并行处理两个视频合并流程,每个流程先完成对应视频的下载,再接收t3传递的音频结果完成合并。 - 自动结果传递:Celery会自动将上游
t3的结果作为参数传递给group内的每个合并任务,无需手动重复调度音频下载。
另一种简洁写法(无需调整参数顺序)
如果不想修改合并任务的参数顺序,可以用partial预先绑定音频结果:
from functools import partial from celery import group, chain t1 = download_video.si(url=video1) t2 = download_video.si(url=video2) t3 = download_audio.si(url=audio1) def merge_video_and_audio(video_path, audio_path): # 音视频合并逻辑 pass # 构建工作流:先下载音频,再并行执行绑定了音频结果的合并任务 workflow = chain( t3, group( t1 | partial(merge_video_and_audio, audio_path=t3.result), t2 | partial(merge_video_and_audio, audio_path=t3.result) ), send_report.s() ) workflow.apply_async()
内容的提问来源于stack exchange,提问作者Martin256
相关产品推荐
相关产品推荐

