如何在一组Celery任务全部执行完成后自动触发action()方法
Celery任务组全部完成后触发指定方法的实现
最简便的方式是直接用Celery原生提供的group任务原语,不需要自己维护任务ID、轮询任务状态,原生支持完成回调触发,稳定无额外依赖。
注意所有方案都需要提前在Celery配置中开启result_backend(可选用Redis、数据库、RabbitMQ等作为结果存储),否则Celery无法追踪单任务状态,无法判定任务组是否全部执行完成。
方案1:回调在Worker端执行(无阻塞,性能最优)
这个方案不需要主进程常驻等待,所有任务执行完成后由Worker自动调度回调方法,适合异步发完任务主进程就可以退出的场景。
- 先把需要触发的
action()也注册为Celery任务,注意回调函数会默认接收任务组所有子任务的返回值组成的列表作为入参:
#tasks.py @app.task def rank(item): # Update database pass @app.task def action(task_results): print('Tasks has been finished.')
- 修改主程序发任务的逻辑,把所有待执行的rank任务打包为group,发任务时通过
link参数绑定完成回调:
#main.py from celery import group from tasks import rank, action import tqdm # 构造所有待执行的任务签名 task_sigs = [ rank.s([{"_id": item["_id"], "max": item["max"]}]) for item in tqdm.tqdm(all_items) ] # 异步提交任务组,绑定完成回调,提交后主进程可以继续做其他事或者直接退出 group(task_sigs).apply_async(link=action.s())
如果需要任务组内有任务执行失败时也触发对应逻辑,可以额外给
apply_async传link_error参数绑定错误回调任务。
方案2:回调在主进程执行(适合需要主进程上下文的场景)
如果不想把action()注册为Celery任务,需要在运行main.py的主进程里执行回调,可以选择阻塞等待或者异步监听任务组完成状态:
阻塞等待版本(逻辑最简单)
#main.py from celery import group from tasks import rank import tqdm def action(): print('Tasks has been finished.') task_sigs = [ rank.s([{"_id": item["_id"], "max": item["max"]}]) for item in tqdm.tqdm(all_items) ] job = group(task_sigs).apply_async() # 阻塞等待所有任务执行完成,返回值是所有子任务的结果列表 job_results = job.join() # 所有任务跑完后触发action action()
非阻塞监听版本(主进程不用卡着等)
如果不想阻塞主进程执行其他逻辑,可以给任务组绑定就绪监听,任务完成后自动触发回调:
#main.py from celery import group from celery.result import GroupResult from tasks import rank import tqdm def action(): print('Tasks has been finished.') def on_tasks_done(sender, result, **kwargs): # 任务组全部完成后会自动调用这个方法 action() task_sigs = [ rank.s([{"_id": item["_id"], "max": item["max"]}]) for item in tqdm.tqdm(all_items) ] job = group(task_sigs).apply_async() # 绑定完成监听,主进程可以继续执行其他逻辑 GroupResult.then(job, on_tasks_done)
不要自己单独循环发apply_async再手动维护任务ID轮询状态,group已经把状态维护、异常判定、回调触发的逻辑做了完整封装,重复造轮子很容易出现状态判定不准、回调漏触发的问题。
内容的提问来源于stack exchange,提问作者2017561-1
相关产品推荐
相关产品推荐

