【问题标题】:When and how to use asyncio queues?何时以及如何使用异步队列?
【发布时间】:2020-05-05 19:29:14
【问题描述】:

我有多个 api 路由,它们通过单独查询数据库来返回数据。

现在我正在尝试构建在 api 之上查询的仪表板。我应该如何将 api 调用放入队列中,以便它们异步执行?

我试过了

await queue.put({'response_1': await api_1(**kwargs), 'response_2': await api_2(**kwargs)})

似乎在将任务放入队列时返回了数据。

现在我正在使用

await queue.put(('response_1', api_1(**args_dict)))

在生产者和消费者中,我正在解析元组并进行 api 调用,我认为我做错了。

问题1 有没有更好的方法?

这是我用来创建任务的代码

producers = [create_task(producer(**args_dict, queue)) for row in stats]
consumers = [create_task(consumer(queue)) for row in stats]
await gather(*producers)
await queue.join()
for con in consumers:
    con.cancel()

Question2 我应该使用 create_task 还是 ensure_future?对不起,如果它是重复的,但我无法理解其中的区别,在网上搜索后我变得更加困惑。

我正在使用 FastAPI、数据库(异步)包。

我使用元组而不是字典,例如 await queue.put('response_1', api_1(**kwargs))

./app/dashboard.py:90: RuntimeWarning: coroutine 'api_1' was never awaited
item: Tuple = await queue.get_nowait()

我的消费者代码是

async def consumer(return_obj: dict, que: Queue):
    item: Tuple = await queue.get_nowait()
    print(f'consumer took {item[0]} from queue')
    return_obj.update({f'{item[0]}': await item[1]})
    await queue.task_done()

如果我不使用 get_nowait 消费者会因为队列可能为空而卡住, 但如果我使用 get_nowait 会显示上述错误。 我没有定义最大队列长度

------------编辑-----------

制片人

async def producer(queue: Queue, **kwargs):
    await queue.put('response_1', api_1(**kwargs))

【问题讨论】:

  • 您编辑的代码看起来无法运行,因为put 只接受一个参数,而您似乎给了它两个参数。请提供实际代码,如果可能,请提供一个最小但可运行的示例来说明问题。
  • 我添加了使用元组对象的生产者,其中第一个元素是 key,第二个元素是 coroutine @ user4815162342
  • 这不是问题中的代码所显示的内容 - 我猜它缺少一对括号。细节在这些事情中很重要,最好提供一个可运行的示例以获得帮助。

标签: python-asyncio fastapi


【解决方案1】:

您可以从您的第一个 sn-p 中删除 await 并在队列中发送 协程对象。协程对象是被调用但尚未等待的协程。

# producer:
await queue.put({'response_1': api_1(**kwargs),
                 'response_2': api_2(**kwargs)})
...

# consumer:
while True:
    dct = await queue.get()
    for name, api_coro in dct:
        result = await api_coro
        print('result of', name, ':', result)

我应该使用create_task 还是ensure_future

如果参数是调用协程函数的结果,则应使用create_task(请参阅 Guido 的this comment 以获得解释)。顾名思义,它将返回一个驱动该协程的Task 实例。该任务也可以等待,但它会继续在后台运行。

ensure_future 是一个更专业的函数,可以将各种等待对象转换为其对应的未来。当实现函数(如asyncio.gather())为了方便而接受不同类型的等待对象并且需要在使用它们之前将它们转换为期货时,它很有用。

【讨论】:

    猜你喜欢
    • 2011-11-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多