如何让Celery中第二个Group在第一个Group完成后再执行?
嘿,我来帮你搞定这个Celery的问题!
你遇到的核心问题是:直接用chain(g1, g2).delay()的时候,Celery只会把第一个group提交到队列就认为它完成了,不会等group里所有任务都执行完毕,所以第二个group会提前启动,这完全不是你想要的效果。
要实现「第一组所有任务完成后再启动第二组」的需求,Celery专门提供了chord这个工具——它就是为“等待一组任务全部结束后再执行后续逻辑”的场景设计的,完美匹配你的需求。
修正后的代码应该是这样的(先帮你修正了代码里的笔误,比如tasks1应该是task1这类小问题):
from celery import group, chord # 假设你已经定义了task1和task2这两个Celery任务 g1 = group([task1.si(1), task1.si(2)]) g2 = group([task2.si(3), task2.si(4)]) # 使用chord:等待g1所有任务完成后,再启动g2 chord(g1)(g2) # 或者用delay()的写法也可以 chord(header=g1, body=g2).delay()
简单解释下:
g1作为chord的header,会先启动里面的所有异步任务,同一组内的任务会并行执行- 只有当
g1里的所有任务都执行完毕,chord才会启动作为body的g2,保证两组任务的严格顺序
如果你非要用chain实现,也可以,但需要加一个中间任务来等待第一组完成(不过这种方法不推荐,会阻塞worker),代码大概是这样:
from celery import chain, group @celery.task def wait_for_group(group_sig): # 等待传入的group所有任务完成 group_result = group_sig.apply_async() group_result.get() # 同步等待,会阻塞当前worker进程 # 用chain串联等待任务和第二组 chain(wait_for_group.s(g1), g2).delay()
显然第一种用chord的方法更简洁、更符合Celery的设计理念,推荐你用这个方案。
内容的提问来源于stack exchange,提问作者user1247196
相关产品推荐
相关产品推荐

