Celery任务如何仅对列表中执行失败的元素触发重试操作?
解决方案
你需求的实现有两种常用方案,可根据你的重试次数管控需求选择:
方案1:使用Celery自带的retry机制(推荐,支持自带的重试次数、间隔管控)
你提到的self.retry写法只需要调整参数传递格式即可,retry方法需要通过args/kwargs参数传递入参:
@app.task(bind=True, max_retries=3) # 可通过max_retries指定最大重试次数,避免无限重试 def my_test(self, my_list:list): new_list = [] for ele in my_list: try: do_something_may_fail(ele) except: new_list.append(ele) # 存在失败元素才触发重试 if new_list: # 重试时将new_list作为参数传入 raise self.retry(args=(new_list,), countdown=5)
注意:如果不设置
max_retries参数,任务会无限重试直到所有元素处理成功,可根据业务需求调整最大重试次数,超出次数后会抛出MaxRetriesExceededError异常,你也可以捕获该异常做失败兜底逻辑。
方案2:主动触发新的异步任务(适合不需要关联原任务重试上下文的场景)
如果你不需要用到Celery原生的重试次数统计、重试上下文关联,也可以直接调用apply_async提交新任务:
@app.task(bind=True) def my_test(self, my_list:list): new_list = [] for ele in my_list: try: do_something_may_fail(ele) except: new_list.append(ele) # 存在失败元素时提交新任务 if new_list: my_test.apply_async(args=(new_list,), countdown=5)
两种方案的差异:
- 使用
retry的方案会保留原任务的上下文,重试记录会关联到原任务ID,Celery的任务状态会显示为RETRY,符合任务重试的语义 - 使用
apply_async的方案会生成完全独立的新任务,和原任务没有关联,原任务执行完会直接标记为SUCCESS,适合不需要回溯原任务重试记录的场景
内容的提问来源于stack exchange,提问作者Mas Zero
相关产品推荐
相关产品推荐

