【问题标题】:Can Celery pass a Status Update to a non-Blocking Caller?Celery 可以将状态更新传递给非阻塞调用者吗?
【发布时间】:2020-01-03 21:38:54
【问题描述】:

我正在使用Celery 异步执行一组操作。有很多这样的操作,每一个都可能需要很长时间,所以我不想将结果发送回 Celery 工作函数的返回值,而是一次将它们发送回作为自定义状态更新。这样调用者就可以实现一个带有改变状态回调的进度条,并且工作函数的返回值可以是恒定大小而不是操作数的线性。

这是一个简单的示例,我使用 Celery 工作函数 add_pairs_of_numbers 添加数字对列表,为每个添加的数字对发送自定义状态更新。

#!/usr/bin/env python

"""
Run worker with:

    celery -A tasks worker --loglevel=info
"""
from celery import Celery

app = Celery("tasks", broker="pyamqp://guest@localhost//", backend="rpc://")

@app.task(bind=True)
def add_pairs_of_numbers(self, pairs):
    for x, y in pairs:
        self.update_state(state="SUM", meta={"x":x, "y":y, "x+y":x+y})
    return len(pairs)

def handle_message(message):
    if message["status"] == "SUM":
        x = message["result"]["x"]
        y = message["result"]["y"]
        print(f"Message: {x} + {y} = {x+y}")

def non_looping(*pairs):
    task = add_pairs_of_numbers.delay(pairs)
    result = task.get(on_message=handle_message)
    print(result)

def looping(*pairs):
    task = add_pairs_of_numbers.delay(pairs)
    print(task)
    while True:
        pass

if __name__ == "__main__":
    import sys

    if sys.argv[1:] and sys.argv[1] == "looping":
        looping((3,4), (2,7), (5,5))
    else:
        non_looping((3,4), (2,7), (5,5))

如果您只运行./tasks,它会执行non_looping 函数。这完成了标准的 Celery 事情:对工作函数进行延迟调用,然后使用 get 等待结果。 handle_message 回调函数打印每条消息,并返回添加的对数作为结果。这就是我想要的。

$ ./task.py
Message: 3 + 4 = 7
Message: 2 + 7 = 9
Message: 5 + 5 = 10
3

虽然对于这个简单的示例来说,非循环场景就足够了,但我要完成的实际任务是处理一批文件,而不是添加成对的数字。此外,客户端是一个Flask REST API,因此不能包含任何阻塞get 调用。在上面的脚本中,我使用looping 函数模拟了这个约束。此函数启动异步 Celery 任务,但不等待响应。 (随后的无限 while 循环模拟 Web 服务器继续运行并处理其他请求。)

如果您使用参数“循环”运行脚本,它将运行此代码路径。在这里它会立即打印 Celery 任务 ID,然后进入无限循环。

$ ./tasks.py looping
a39c54d3-2946-4f4e-a465-4cc3adc6cbe5

Celery worker 日志显示执行了 add 操作,但是调用者没有定义回调函数,所以它永远不会得到结果。

(我意识到这个特定的例子是令人尴尬的并行,所以我可以使用chunks 将它划分为多个任务。但是,在我的非简化现实世界案例中,我有无法并行化的任务。)

我想要的是能够在looping 场景中指定回调。像这样。

def looping(*pairs):
    task = add_pairs_of_numbers.delay(pairs, callback=handle_message) # There is no such callback.
    print(task)
    while True:
        pass

在 Celery 文档和我可以在网上找到的所有示例(例如 this)中,无法将回调函数定义为 delay 调用或其等效 apply_async 的一部分。您只能指定一个作为get 回调的一部分。这让我觉得这是一个有意的设计决定。

在我的 REST API 场景中,我可以通过让 Celery 工作进程以 HTTP 帖子的形式将“状态更新”发送回 Flask 服务器来解决此问题,但这似乎很奇怪,因为我开始复制消息传递HTTP 中的逻辑已经存在于 Celery 中。

有没有办法编写我的looping 场景,以便调用者在不进行阻塞调用的情况下接收回调,或者在 Celery 中明确禁止这样做?

【问题讨论】:

    标签: python asynchronous flask celery


    【解决方案1】:

    这是一种 celery 不支持的模式,尽管您可以(在某种程度上)通过将自定义状态更新发布到您的任务 as described here 来欺骗它。

    使用 update_state() 更新任务的状态:

    def upload_files(self, filenames):
        for i, file in enumerate(filenames):
            if not self.request.called_directly:
                self.update_state(state='PROGRESS',
                    meta={'current': i, 'total': len(filenames)})```
    

    celery 不支持这种模式的原因是任务生产者(调用者)与任务消费者(工作者)是强解耦的,两者之间的唯一通信是代理,以支持从生产者到消费者的通信以及结果后端支持从消费者到生产者的通信。您目前可以获得的最接近的方法是轮询任务状态或编写自定义结果后端,这将允许您通过 AMP RPC 或 redis 订阅发布事件。

    【讨论】:

      猜你喜欢
      • 2015-11-27
      • 1970-01-01
      • 2018-06-11
      • 1970-01-01
      • 1970-01-01
      • 2012-01-20
      • 2012-04-30
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多