You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 16:03:31