【问题标题】:How can Celery workflows include dynamically generated groups?Celery 工作流程如何包含动态生成的组?
【发布时间】: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/…

标签: python celery


【解决方案1】:

这是您的示例,工作方式有所不同,但在我看来,实现了预期的结果。

@app.task
def collect_items(n):
    logger.info("collect %r items", n)
    items = list(range(n))
    return items


@app.task
def schedule_task_group(items):
    logger.info("group get_details tasks & pass results to publish_items")
    return (
        group(get_details.s(i) for i in items) | publish_items.s()
    ).delay()


@app.task
def get_details(item):
    logger.info("get item detail for = %r", item)
    return (item, item * item)


@app.task
def publish_items(items):
    logger.info("publish items = %r", items)
    return items


print('schedule collect_items & pass result to schedule_task_group with n = 5')
(collect_items.s(5) | schedule_task_group.s()).delay()

与您的代码的主要区别在于,我将get_details 组与publish_items 任务链接在一起,有效地使其成为一个和弦。这在文档中提到并且是必需的,因为您希望在将其传递给 publish_items 之前安排整个任务组并完成运行。

请查看@quantoid 并告诉我您的想法。 请注意,使用 -l INFO 标志运行 celery 将更容易可视化工作人员中实际发生的情况。

参考: - http://docs.celeryproject.org/en/latest/userguide/canvas.html - https://stackoverflow.com/a/15147171/484127

【讨论】:

    猜你喜欢
    • 2020-07-06
    • 2010-12-15
    • 2021-10-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多