【发布时间】:2013-11-17 23:48:02
【问题描述】:
我有一些这样的任务:
@celery.task
def generate():
sleep(1.0)
print "Generate done!"
return 'result'
@celery.task
def lower(result):
sleep(1.0)
print "Lower done!"
return result.lower()
@celery.task
def upper(result):
sleep(1.0)
print "Upper done!"
return result.upper()
@celery.task
def upload(result):
sleep(1.0)
print "Upload done for: %s!" % (result,)
return 'upload'
@celery.task
def callback(results):
print "It's all done! %s" % (results,)
我正在创建一个如下所示的和弦:
chord(
header=chain(
generate.s(),
group(
chain(lower.s(), upload.s()),
chain(upper.s(), upload.s())
)
), body=callback.s()
).delay()
我遇到的问题是我的回调(应该在所有任务完成后触发)似乎在 generate 之后立即触发。
如果不清楚,工作流程是这样的:
- 生成一个结果,然后将其结果传递给一个组的成员,从而实现并行:
- 第一组将从
generate获取结果,使用lower将其转换为小写,然后使用upload上传结果。 - 第二组将从
generate获取结果,使用upper将其转换为大写,然后使用“上传”上传结果。
- 第一组将从
- 完成所有这些后,应该调用
callback任务回调。
预期
callback 任务将在启动后至少 3 秒被调用。
实际
callback 任务在启动后大约 1 秒被调用,并且不等待组成员完成执行。
以下是证明它不等待组的日志:
[2013-11-17 18:20:40,447: WARNING/PoolWorker-8] Generate done!
[2013-11-17 18:20:41,493: WARNING/PoolWorker-6] Upper done!
[2013-11-17 18:20:41,493: WARNING/PoolWorker-1] Lower done!
[2013-11-17 18:20:41,535: WARNING/PoolWorker-6] It's all done! [('e0016a35-d538-4e96-ad86-6ddf91ef4a09', [('b1af78a9-7935-4037-84e4-9fae6d7c027e', None), ('d69c4c99-af9c-476f-af7d-7f647c4d9c83', None)])]
[2013-11-17 18:20:42,522: WARNING/PoolWorker-7] Upload done for: result!
[2013-11-17 18:20:42,523: WARNING/PoolWorker-5] Upload done for: RESULT!
似乎 Celery 不等待组。有没有办法让 Celery 等到所有任务(包括组成员)执行完毕?
【问题讨论】:
标签: celery