【发布时间】:2018-08-07 20:37:16
【问题描述】:
考虑一下这个 Celery 工作流程:
wf = collect_items.s() | add_details.s() | publish_items.s()
它收集一些项目,为每个项目并行添加额外的细节,然后在某处发布修饰信息。
我想要add_details 充当一组任务,每个项目一个,并行获取每个项目的详细信息。显然这个组必须是由collect_items输出的数据生成的。
这是我尝试过的,使用默认的 rabbitmq 代理:
app = Celery(backend="rpc://")
@app.task
def collect_items(n):
return range(n)
@app.task
def add_details(items):
return group(get_details.s(i) for i in items).delay()
@app.task
def get_details(item):
return (item, item * item)
@app.task
def publish_items(items):
print("items = %r" % items)
我希望输出是数字 0-9,并用它们的正方形装饰,所有这些都是同时计算的:
>>> wf.delay(10).get()
items = [(0, 0), (1, 1), (2, 4), ... (8, 64), (9, 81)]
这确实调用了预期的任务,但不幸的是,将结果作为一组 GroupResults 传递给 publish_items,其中包含处于 PENDING 状态的 AsyncResults,即使任务似乎已经完成。
我等不及publish_items 中的这些结果,因为您不能在任务中使用get()(死锁风险等)。我认为 Celery 会识别像 add_details 这样的任务何时返回 GroupResult 并在其上执行 get ,然后返回该值以传递给链中的下一个任务。
这似乎是一种常见的模式,在 Celery 中是否可以这样做?
我在这里看到过类似的问题,但答案似乎假设对 Celery 如何在幕后工作有很多深入的了解,但无论如何它们对我不起作用。
【问题讨论】:
-
简要看一下您的示例,它应该可以正常工作。我唯一想知道的是
add_details中的那个组是否应该调用delay()...这有效地安排了任务,而您可能想要返回的是签名...查看之前的示例链:docs.celeryproject.org/en/latest/userguide/…