You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确实现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()

关键逻辑说明

  1. 单一音频任务调度:通过chain先执行t3,确保音频下载仅触发一次,执行完成后自动将结果传递给后续任务。
  2. 并行合并执行:group(merge_task1, merge_task2)会并行处理两个视频合并流程,每个流程先完成对应视频的下载,再接收t3传递的音频结果完成合并。
  3. 自动结果传递: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 04:57:37