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

如何在Celery Chunks/Starmap中重试失败任务?配置无效排查

Celery Chunks中失败任务的重试配置方案

问题背景

使用Celery Chunks(推测由celery.starmap执行)时,单个任务失败会导致整个Chunk组直接失败,且已为任务配置自动重试规则,但失败任务未触发重试(重试次数为0)。

解决方案

1. 让Chunk组忽略单个任务失败,保留重试逻辑

Celery的Group默认会因组内单个任务失败而标记整个组为失败,直接终止后续流程,同时抑制了单个任务的重试逻辑。只需在创建Chunk对应的Group时添加errors='ignore'参数,即可让组内其他任务继续执行,失败任务则按自身配置触发重试。

2. 确保单个任务的重试配置精准有效

保留任务的autoretry_for、retry_kwargs核心配置,可额外添加retry_jitter=True,避免多个任务同时重试引发的资源竞争问题。

修改后的代码示例

@celery_app.task
def get_item_asteroids_callback(parent_item_id):
    pass


@celery_app.task(
    autoretry_for=(requests.exceptions.ReadTimeout, TimeoutError, urllib3.exceptions.ReadTimeoutError, InterfaceError),
    retry_kwargs={'max_retries': 3},
    retry_backoff=True,
    retry_jitter=True  # 可选,优化重试间隔的随机性
)
def get_item_asteroid(work_id):
    """
    执行可能抛出异常的业务操作
    """
    raise requests.exceptions.ReadTimeout('Test Exception')


@celery_app.task(ignore_result=False)
def get_item_asteroids(parent_item_id):
    item_ids = [1,2,3.......]
    chunk_size = min(len(item_ids), 128)
    # 创建Chunk组时指定errors='ignore',允许单个任务失败不影响全局流程
    chunk_group = get_item_asteroid.chunks(item_ids, chunk_size).group(errors='ignore')
    # 绑定Chord回调任务
    chord(chunk_group)(get_item_asteroids_callback.si(parent_item_id))

关键说明

  • errors='ignore'参数需Celery 4.x及以上版本支持,作用是让Group不因为单个任务失败而终止,保障其他任务正常执行的同时,触发失败任务的重试逻辑。
  • 任务的autoretry_for需覆盖所有需要触发重试的异常类型,确保异常被正确捕获并触发重试流程。

内容的提问来源于stack exchange,提问作者Eric O.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:16:09