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

Celery v5.3.6异步执行group任务时部分任务未执行问题

问题原因

当你在process_task中直接调用task_group()时,这是同步执行子任务:

  • 同步调用process_task()时,子任务在本地执行,Celery的GroupResult会遍历所有任务并收集结果,即便有任务抛出异常,也不会终止后续任务的执行(因为未调用get()触发异常抛出)。
  • 异步调用process_task.apply_async()时,process_task在Worker进程中执行,子任务同步执行时抛出的异常会直接终止process_task本身,导致后续子任务根本没机会被执行。
解决方案

把task_group()改为异步提交group任务,让子任务独立进入Celery队列并行执行,这样即便某个子任务失败,其余任务仍会正常调度:

@shared_task
def success_task(num):
    return num

@shared_task
def failed_task(num):
    raise Exception

@shared_task
def process_task():
    task_group = group(success_task.s(1), failed_task.s(2), failed_task.s(3), success_task.s(4))
    # 改用apply_async异步提交group
    task_group.apply_async()
额外说明

如果需要等待group中所有任务完成后执行后续逻辑,可以使用chord(需确保你的Celery后端支持):

from celery import chord

@shared_task
def callback(results):
    # 处理所有任务的结果
    pass

@shared_task
def process_task():
    task_group = group(success_task.s(1), failed_task.s(2), failed_task.s(3), success_task.s(4))
    chord(task_group)(callback.s())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 17:25:18