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

Celery Canvas多层任务未执行回调及后续任务求助

Celery多层Canvas任务不触发后续任务的问题修复

问题核心

你的代码中chain嵌套chord的写法存在逻辑错误,导致update_job_status和final_follow_up_job任务无法被触发,具体原因:

  • chord的回调任务用si()创建了不可变任务,无法接收group的执行结果;且on_error仅在回调任务自身执行失败时触发,并非group中任务失败时触发,完全不符合你预期的"group执行完后无论成败都更新状态"逻辑。
  • 错误的chord结构导致chain无法正确识别前序任务的完成状态,进而中断了后续final_follow_up_job的执行。

修复方案

方案1:在回调任务中处理成败逻辑

修改update_job_status任务,让它接收group的执行结果,自行判断任务组的成败状态后执行对应逻辑:

@task
def update_job_status(job_id_list, task_results):
    # 判断group内所有任务是否执行成功
    all_succeeded = all(not isinstance(res, Exception) for res in task_results)
    # 执行你的状态更新逻辑
    print(f"Job IDs {job_id_list} status updated to: {'success' if all_succeeded else 'failed'}")
    # 返回结果给chain的下一个任务(可选)
    return all_succeeded

# 重构Canvas结构
job = chain(
    chord(
        group(job_list),
        update_job_status.s(job_id_list)  # 使用s()而非si(),接收group的执行结果
    ),
    final_follow_up_job.si(job_id_list)
)
job.apply_async()

如果需要分开处理group的成功和失败场景,可直接给group绑定link(成功触发)和link_error(失败触发),再接入chain:

# 定义任务组并绑定成败回调
job_group = group(job_list)
job_group.link(update_job_status.si(job_id_list, success=True))
job_group.link_error(update_job_status.si(job_id_list, success=False))

# 构建chain,确保final_follow_up_job在group执行完成后触发
job = chain(
    job_group,
    final_follow_up_job.si(job_id_list)
)
job.apply_async()

关键注意点

  • 避免给chord的回调任务使用si():不可变任务无法接收前序任务的结果,会破坏Canvas的执行链。
  • 明确on_error的作用范围:它是给单个任务绑定错误处理,仅在该任务自身执行失败时触发,不要用来处理前置任务的错误。

内容的提问来源于stack exchange,提问作者ProgramSpree

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 20:05:28