Celery任务流中如何将前置myname任务结果作为参数传给reverse任务
Celery Group任务结果传递给下游Signature的实现
首先修正3个会直接导致传参失败的基础问题:
- 所有需要传递返回值的任务不能配置
ignore_result=True,该配置开启后Celery不会存储任务返回值,下游完全拿不到上游输出,需要将add、myname、reverse三个任务的该配置移除或设为False - Celery的signature单元素参数必须加尾逗号,否则会被识别为普通变量而非参数元组,现有代码中
args=(n)要改成args=(n,) - group执行完成后,会按照内部任务的书写顺序,把所有任务的返回值打包成列表传给紧邻的下游任务。你的代码里group先写的add、后写的myname,所以结果列表里索引0是add的返回值,索引1是myname的返回值
核心注意点
你预留的reverse的args位置不需要填写任何静态取值表达式:非immutable模式(即不用si()定义)的signature会自动接收上游返回值作为运行时参数,如果硬填静态值,反而会和上游传入的参数叠加,触发参数数量不匹配报错。
可直接运行的修正代码
方式1:不新增任务,直接调整reverse的入参逻辑(改动最小)
from app import app from time import sleep from celery.utils.log import get_task_logger import os from celery import signature, chain, group, chord from celery.result import allow_join_result MyQUEUE = os.getenv("SCANS_QUEUE") logger = get_task_logger(__name__) @app.task(queue=MyQUEUE) def reverse(group_result): # 从group结果中取索引1的myname返回值,再取name字段作为待反转文本 text = group_result[1]["name"] logger.info('reverse order {}'.format(text)) return {"reversename": str(text[::-1])} @app.task(queue=MyQUEUE, ignore_result=True) # add结果不需要传给下游,可以开ignore_result def add(a,b): logger.info('Addition --> a : {0} & b : {1} '.format(a,b)) return {"addition": str(a+b)} @app.task(queue=MyQUEUE) def myname(a): logger.info('Name --> a : {0}'.format(a)) return {"name": str(a)} @app.task(queue=MyQUEUE) def run_pipeline(a,b,n): resultchain = chain([ group([ signature( add, args=(a,b), queue=MyQUEUE ), signature( myname, args=(n,), # 单元素元组补尾逗号 queue=MyQUEUE ) ]), signature( reverse, # args位置留空,自动接收group返回的结果列表 queue=MyQUEUE ) ]).apply_async() with allow_join_result(): results = resultchain.join() return results
方式2:保留reverse原有逻辑,新增中间提取任务(适合reverse被多处复用、不能改入参的场景)
如果reverse已经在其他业务逻辑中使用,入参固定为待反转文本,可以新增一个极轻量的中间任务做结果提取,不需要修改reverse代码:
# 保留原reverse逻辑不变 @app.task(queue=MyQUEUE) def reverse(text): logger.info('reverse order {}'.format(text)) return {"reversename": str(text[::-1])} # 新增中间提取任务 @app.task(queue=MyQUEUE) def extract_myname_text(group_result): # 提取myname返回的name字段 return group_result[1]["name"] # run_pipeline中的chain调整为: resultchain = chain([ group([ signature(add, args=(a,b), queue=MyQUEUE), signature(myname, args=(n,), queue=MyQUEUE) ]), signature(extract_myname_text, queue=MyQUEUE), signature(reverse, queue=MyQUEUE) # reverse自动接收extract任务返回的文本参数 ]).apply_async()
可选优化
如果要让group+回调的语义更清晰,可以用Celery专门为该场景设计的chord写法替换chain接group的写法,执行逻辑完全一致:
resultchain = chord( header=group([ signature(add, args=(a,b), queue=MyQUEUE), signature(myname, args=(n,), queue=MyQUEUE) ]), # body支持直接写任务/chain,逻辑和chain下游一致 body=signature(reverse, queue=MyQUEUE) ).apply_async()
内容的提问来源于stack exchange,提问作者Faisal Shahbaz
相关产品推荐
相关产品推荐

