【发布时间】:2014-10-28 15:09:58
【问题描述】:
我要做的是执行一组任务,并将组中每个任务的结果发送到另一个任务,然后将所有这些任务的结果发送到一个最终任务。
例如
jobgroup = (
group(tasks.task1.s(p) for p in params) |
tasks.dmap.s(tasks.task2.s())
)
这部分对我来说没问题,但我遇到的问题是,如果我试图将所有结果放入一个最终任务中,例如
mychain = chain(jobgroup | tasks.task3.s())
我看到的是 task3() 正在被调用,而来自 task2 的任务仍处于挂起状态,方法是在它们进入时打印出状态(以下 sn-p 是我在 task3 函数中打印的内容:
@task()
def task3(input):
for item in input.results:
logger.info(item.status)
记录结果
[2014-09-04 10:48:41,905: INFO/Worker-7] tasks.task3[27053688-3c5c-4ca5-975f-356a66d55364]: PENDING
[2014-09-04 10:48:41,905: INFO/Worker-7] tasks.task3[27053688-3c5c-4ca5-975f-356a66d55364]: PENDING
那么我该如何设置,以便在作业组中的所有任务完成之前不调用 task3?
【问题讨论】:
-
我在这里也注意到的其他一点是父任务完成并按预期返回所有结果但是我的任务3仍然有待处理的项目