如何实现监控Celery任务执行成功后触发后续任务?非自动调用场景
解决方案:监控指定任务完成后触发回调任务
好问题!默认的Celery Chord机制会自动提交header中的任务,这显然不符合你「不自动调用任何任务,仅监控指定任务执行成功后再触发tasks.task_3」的需求。下面给你两种实用的实现方案:
方案一:使用Celery任务成功信号(推荐分布式场景)
Celery的task_success信号会在任意任务成功执行后触发,我们可以利用这个信号来跟踪目标任务的完成状态,当所有指定任务都成功后再调用回调任务。
步骤说明:
- 先启动(或记录)你要监控的任务实例,保存它们的任务ID
- 用共享存储(比如Redis,分布式环境必备)跟踪已完成的任务
- 绑定
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
相关产品推荐
相关产品推荐

