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

如何实现监控Celery任务执行成功后触发后续任务?非自动调用场景

解决方案:监控指定任务完成后触发回调任务

好问题!默认的Celery Chord机制会自动提交header中的任务,这显然不符合你「不自动调用任何任务,仅监控指定任务执行成功后再触发tasks.task_3」的需求。下面给你两种实用的实现方案:

方案一:使用Celery任务成功信号(推荐分布式场景)

Celery的task_success信号会在任意任务成功执行后触发,我们可以利用这个信号来跟踪目标任务的完成状态,当所有指定任务都成功后再调用回调任务。

步骤说明:

  1. 先启动(或记录)你要监控的任务实例,保存它们的任务ID
  2. 用共享存储(比如Redis,分布式环境必备)跟踪已完成的任务
  3. 绑定task_success信号,每次有任务成功时检查是否是目标任务,当全部完成后触发tasks.task_3

代码示例:

import redis
from celery import current_app
from celery.signals import task_success

# 初始化Redis客户端(根据你的配置调整参数)
redis_client = redis.Redis(host="localhost", port=6379, db=0)

# 定义存储任务ID的Redis键名
TARGET_TASKS_KEY = "monitored_task_ids"
COMPLETED_TASKS_KEY = "completed_task_ids"

# ----------------------
# 第一步:启动你要监控的任务(或者记录已启动任务的ID)
# ----------------------
task1 = current_app.tasks["tasks.task_1"].delay()  # 手动启动任务,而非Chord自动触发
task2 = current_app.tasks["tasks.task_2"].delay()
# 将目标任务ID存入Redis集合
redis_client.sadd(TARGET_TASKS_KEY, task1.id, task2.id)

# ----------------------
# 第二步:定义信号处理函数
# ----------------------
@task_success.connect
def trigger_callback_on_completion(sender=None, result=None, **kwargs):
    current_task_id = sender.request.id
    # 检查当前成功的任务是否在我们的监控列表中
    if redis_client.sismember(TARGET_TASKS_KEY, current_task_id):
        # 标记该任务已完成
        redis_client.sadd(COMPLETED_TASKS_KEY, current_task_id)
        # 检查是否所有监控任务都已完成
        target_count = redis_client.scard(TARGET_TASKS_KEY)
        completed_count = redis_client.scard(COMPLETED_TASKS_KEY)
        if target_count == completed_count:
            # 触发回调任务
            current_app.tasks["tasks.task_3"].delay()
            # 清理Redis中的临时数据
            redis_client.delete(TARGET_TASKS_KEY, COMPLETED_TASKS_KEY)
            # 断开信号连接,避免重复触发
            trigger_callback_on_completion.disconnect()

方案二:使用AsyncResult轮询等待(适合简单单进程场景)

如果你的场景比较简单(比如单进程、不需要异步监控),可以直接用Celery的AsyncResult来等待所有目标任务完成,然后再触发回调。

代码示例:

from celery import current_app
from celery.result import AsyncResult
from celery import group

# 启动(或获取)要监控的任务实例
task1 = current_app.tasks["tasks.task_1"].delay()
task2 = current_app.tasks["tasks.task_2"].delay()

# 创建任务结果组,等待所有任务成功完成
task_group = group([AsyncResult(task1.id), AsyncResult(task2.id)])
# 这里会阻塞当前进程,直到所有任务成功(任务失败会抛出异常)
task_group.get()

# 所有任务完成后,触发回调任务
current_app.tasks["tasks.task_3"].delay()

注意事项:

  • 方案二的get()方法会阻塞当前进程,不适合在Web请求或需要异步处理的场景中使用
  • 如果任务可能失败,建议添加异常处理逻辑,比如捕获Exception或者提前检查任务状态

为什么不能用默认的Chord?

默认的chord(header)(callback)写法会自动提交header中的任务,这和你「不自动调用任务,仅监控已有任务」的需求完全冲突,所以必须用上面的替代方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:41:15