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

如何在API注册中实现Celery的条件任务执行与回调

Celery条件工作流与回调实现方案

核心实现代码

任务定义(my_tasks.py)

from celery import Celery, chain, maybe

app = Celery('my_tasks', broker='redis://localhost:6379/0')

@app.task
def task_a():
    # 模拟业务处理,返回True/False
    result = True  # 根据实际逻辑调整
    return result

@app.task
def task_b():
    return "Task B executed."

@app.task
def task_c():
    return "Task C executed."

Flask API实现

from flask import Flask
from my_tasks import task_a, task_b, task_c, chain, maybe

app = Flask(__name__)

@app.route('/register-task', methods=['POST'])
def register_task():
    # 构建顺序执行的工作流:Task A → 条件执行Task B → 总是执行Task C
    workflow = chain(
        task_a.s(),
        maybe(task_b.s()),
        task_c.s()
    )

    # 启动工作流
    result = workflow.delay()
    return {"task_id": result.id}, 202

问题解答

1. 基于Task A结果条件执行Task B

使用Celery原生的maybe原语即可实现简单布尔条件的分支逻辑:

  • maybe(task_b.s())会自动检查前序任务(task_a)的返回值,仅当返回值为True时,才触发task_b的执行。
  • 如果需要更复杂的条件判断(比如匹配特定字符串、数值范围),可以自定义路由任务:
    @app.task
    def route_task(result):
        if result == "特定值":  # 替换为你的业务条件
            task_b.delay()
    
    之后在工作流中替换为chain(task_a.s(), route_task.s())即可。

2. 确保Task C在Task A成功后无论结果都运行

有两种可靠方案,根据执行顺序需求选择:

  • 顺序执行:将task_c放在chain的最后位置(如上述示例),这样无论task_b是否执行,task_c都会在task_a完成后(或task_b完成后)启动。
  • 并行执行:使用task_a的link参数,让task_c在task_a成功后立即触发,与条件分支并行运行:
    workflow = task_a.s().link(task_c.s()).link(maybe(task_b.s()))
    
    这种方式下,task_c和task_b(如果满足条件)会同时启动,无需等待对方完成。

3. API中构建工作流的最佳方式

推荐使用Celery官方提供的工作流原语组合任务,优势在于:

  • 原子性:整个工作流通过一次delay()提交,避免分散调用导致的状态不一致。
  • 可追踪性:通过返回的task_id可以追踪整个工作流的执行状态,包括每个子任务的结果。
  • 可维护性:原语(chain/maybe/group)语义清晰,便于后续修改工作流逻辑。

避免在任务内部直接调用delay(),这种方式会割裂工作流的整体结构,增加调试和维护成本。


内容的提问来源于stack exchange,提问作者Mohammad_Moataz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 23:30:58