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

Celery含Group的Chain异常:Group任务失败后后续任务未执行

问题根源

你碰到的问题是Celery的默认行为导致的:只要Group里有一个任务失败,整个Group就会被标记为失败状态,而Chain会在遇到失败的节点时直接中断执行,所以后续的task_final和Chain级别的on_error都不会触发。

解决方案

以下两种方案可以满足你的需求:让单个失败任务触发自身task_error、完成所有Group任务后执行Chain后续逻辑。

方案一:让Group忽略失败,继续执行Chain后续步骤

给Group添加propagate=False配置,这样即使Group内有任务失败,Group本身不会向上抛出异常,Chain会认为Group执行完成,继续走后续的task_final。同时保留单个任务的on_error处理自身失败逻辑。

修改后的Canvas代码:

id = 1  # 数据库主键ID
# 注意:循环变量不能用id,会覆盖外层的id值
chunk = [
          task_to_group.si(i).on_error(task_error.s(i)) 
          for i in range(0,5)
        ]
c = chain(
      task_initial.si(id),
      # 关键:添加propagate=False,让Group不传播内部任务的失败
      group(chunk).set(propagate=False),
      task_final.si(id)
    ).on_error(task_error.s(id))
c.apply_async()

方案二:按需执行Chain的成功/失败逻辑

如果需要在Group有任何任务失败时,最终触发Chain级别的task_error,同时保证所有Group任务都执行完毕,可以用chord替代Chain里的Group,通过中间任务判断Group执行结果,再决定走task_final还是task_error。

示例代码:

id = 1
chunk = [
          task_to_group.si(i).on_error(task_error.s(i))
          for i in range(0,5)
        ]

# 新增中间任务:判断Group执行结果
@app.task
def check_group_result(group_result, main_id):
    # 检查Group内是否有失败任务
    failed_tasks = [res for res in group_result if not res.successful()]
    if failed_tasks:
        return task_error.si(main_id)
    else:
        return task_final.si(main_id)

# 用chord组织逻辑:先执行所有Group任务,再执行结果判断
c = chain(
    task_initial.si(id),
    chord(chunk)(check_group_result.s(main_id=id))
).on_error(task_error.s(id))
c.apply_async()
关键说明
  • propagate=False:让Group不把内部任务的失败向上传播,Chain会认为Group已完成,继续执行后续步骤。
  • 单个任务的on_error:依然会正常触发,处理单个任务的失败逻辑(比如更新数据库状态)。
  • 循环变量修正:原代码中for id in range(0,5)会覆盖外层的id变量,导致所有Group任务用的是循环的id值,必须修正循环变量名称。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 06:43:17