Celery回调参数传递不一致,如何替换任务链路后续入参
问题根因
- Celery的
link参数传入任务序列时,所有任务为并行触发,均直接接收前序主任务的返回值,任务间不存在参数传递关系,因此alt_task生成的新id无法传给同序列的task_one - 你当前的写法中,link里的第二个task_one接收的始终是第一次调用task_one返回的初始id=1,自然输出重复的1
修复方案
推荐使用Celery内置的chain串行任务组实现参数透传,chain会自动将上一个任务的返回值作为下一个任务的第一个入参,刚好匹配你的业务需求,修改后代码如下:
from celery import chain @shared_task() def alt_task(): flow = Flow.objects.create() return flow.id @shared_task() def task_one(flow_id): flow = Flow.objects.get(id=flow_id) print(flow.id) return flow.id @shared_task() def main_task(flow_id): flow = Flow.objects.create() # 构造串行执行链路,参数逐环传递 flow_chain = chain( task_one.s(flow.id), alt_task.s(), task_one.s() ) flow_chain.apply_async(link_error=log_error.s())
运行后输出就会符合你的预期:
1 2
补充说明
如果你的alt_task不需要接收上游任务的传入参数,也可以用不可变签名alt_task.si()替代alt_task.s(),不会影响后续的参数传递逻辑。
内容的提问来源于stack exchange,提问作者af3ld
相关产品推荐
相关产品推荐

