Celery嵌套Chord使用报错:无法将Chord输出传入另一Chord
Celery嵌套Chord报错:AttributeError: 'AsyncResult' object has no attribute 'AsyncResult'
问题原因
你犯了Celery Canvas用法的典型错误:把已执行任务的AsyncResult对象当作了Chord的header参数。
当你执行chord(i.s(1))(group(i.s(), i.s()))时,返回的是一个AsyncResult(代表这个Chord任务的执行结果),而外层chord的header需要的是未执行的任务签名(Signature)或任务组(Group),而非已提交的任务结果对象,这直接触发了属性错误。
这不是Celery的Bug,而是写法不符合其设计规则——Celery的Canvas结构是声明式的任务编排,需要传递的是任务的定义,而非运行中的结果。
正确实现「单任务→任务组→单任务」流程
你的业务流不需要嵌套Chord,用chain结合chord就能完美实现:
完整代码示例
from celery import Celery, chain, chord, group app = Celery('canvastest', backend='redis://', broker='redis://') @app.task def i(x): return x # 编排任务流 if __name__ == "__main__": workflow = chain( i.s(1), # 第一步:执行单个任务 chord( group(i.s(), i.s()), # 第二步:任务组自动接收上游结果作为参数 i.s() # 第三步:最终任务接收任务组的汇总结果 ) ) # 获取执行结果 print(workflow.get(timeout=5)) # 输出:[[1], [1]]
逻辑拆解
chain:按顺序执行任务链,前一个任务的输出自动作为后一个任务的输入。chord:等待header中的任务组全部完成后,将任务组的结果集合传给body中的最终任务。- 任务组
group(i.s(), i.s())会自动接收上游i.s(1)的输出,作为每个子任务的参数。
验证方式
按你原来的运行步骤启动Redis、Celery Worker,然后在Celery Shell中执行上述代码即可得到预期结果。
内容的提问来源于stack exchange,提问作者ajn
相关产品推荐
相关产品推荐

