Celery如何在指定taskA完成后调度taskB并传入其执行结果
Celery任务调度:指定taskA完成后自动触发taskB
问题场景
我有如下Celery任务定义(文件tasks.py):
# tasks.py @app.task def taskA(): time.sleep(10) return "result" @app.task def taskB(result): print("got result")
程序会定期启动新的taskA任务,并保存任务ID供后续使用:
res = taskA.apply_async() keep_for_later(res.id)
后续程序运行过程中,当满足特定条件时,需要实现:在指定的taskA任务完成后,立即让Celery异步运行taskB,并将taskA的返回结果作为taskB的输入。当前触发逻辑的代码片段如下:
taskA_id = get_task_id() if condition_is_met(): res = app.AsyncResult(taskA_id) # 此处需要实现需求逻辑
运行Celery Worker的命令:
celery -A tasks worker
要求两个任务均在独立于主程序的Worker上异步执行,环境为Ubuntu系统,Celery版本v5.3.6。
可行实现方案
修改任务定义
更新tasks.py中的taskB,让它接收taskA的任务ID,主动获取taskA的执行结果:
# tasks.py @app.task def taskA(): time.sleep(10) return "result" @app.task def taskB(task_id): res = celery_app.AsyncResult(task_id) with allow_join_result(): taskA_result = res.get() # 在这里使用taskA_result进行业务操作
触发taskB
当条件满足时,直接异步启动taskB并传入目标taskA的ID:
taskA_id = get_task_id() if condition_is_met(): taskB.apply_async((taskA_id,))
内容的提问来源于stack exchange,提问作者agrav
相关产品推荐
相关产品推荐

