来自 github 问题 #1881 如果回调设置了 link_error 选项集,它采用任务名称列表,那么当和弦的任务失败时,link_error 任务将被执行。
@task(name='super_task.good')
def good():
return True
@task(name='super_task.raise_exception')
def raise_exception():
raise ValueError('error')
@task(name='super_task.callback')
def callback(*args, **kwargs):
logger.info('callback')
logger.info(args)
logger.info(kwargs)
return 'finished'
@task(name='super_task.error_callback')
def error_callback(*args, **kwargs):
logger.info('error_callback')
logger.info(args)
logger.info(kwargs)
return 'error'
>>> c = chord(
[raise_exception.s(), good.s(), raise_exception.s()],
callback.s().set(link_error=['super_task.error_callback'])
)
>>> result = c()
这将执行和弦,在你的 celery 日志中,你会看到 raise_exception 任务失败,并且 error_callback 的执行将在它的 args 中收到 callback 的 task_id。
此时result 的值将包含callback 的AsyncResult实例,并且因为在和弦中错误传播到执行result.get() 的回调将引发任务异常和result.traceback给你回溯。
如果你想要单个回调,只需将和弦回调的名称传递给link_error
callback.s().set(link_error='super_task.callback')
注意
另一个选项是设置CELERY_CHORD_PROPAGATES = False,它将恢复到前 celery 3.1 的行为并始终执行回调。
但这不是推荐的方法,因为您可以在 github 问题#1349 中找到
Celery 3.1 定义了如何处理和弦错误,以前的行为从未记录过
以及更多的意外,因为它从来没有打算那样工作。
我们无法更改错误修复版本中的行为,因此必须改用设置,
但绝不是有人故意禁用新行为。
新行为是为了防止发生此类问题,并且可能会删除向后兼容的设置。我建议你在这里找到一些其他方法来处理错误(如果你可以为它发明一个好的 api,我不介意提出建议)