【问题标题】:Get all items in an asyncio.Queue and return them获取 asyncio.Queue 中的所有项目并返回它们
【发布时间】:2020-09-18 23:05:06
【问题描述】:

如何从asyncio.Queue 实例中收集所有项目并将它们作为结果返回?

注意事项:

  • 无法提前知道放入队列的项目总数(但它是有限的,可以放入内存中)。
  • 多个生产者可以将项目添加到队列中;生产者自己不一定知道所有项目何时都已添加到队列中。

The examples I have seen for consumers of an asyncio.Queue 不必收集并返回结果;他们使用队列中的项目但没有返回值。他们依靠副作用来完成工作,而不关心返回结果。

具体来说,下面是一个简单的例子。我弄清楚如何完成这项工作的唯一方法是向从队列中收集结果的协程提供一个名为 itemsoutput parameter/output argument

import asyncio
import random


async def add_queue_item(item, queue):
    # simulate some work
    sleep_interval = random.randint(0, 3)
    await asyncio.sleep(sleep_interval)
    output_item = item + 1
    await queue.put(output_item)


async def get_all_queue_items(queue, items):
    while True:
        items.append(await queue.get())
        queue.task_done()


async def main():
    queue = asyncio.Queue()
    items = []
    producer_tasks = [asyncio.create_task(add_queue_item(item, queue)) for item in range(5)]
    collect_queue_items_task = asyncio.create_task(get_all_queue_items(queue, items))
    await queue.join()
    await asyncio.gather(*producer_tasks)
    collect_queue_items_task.cancel()
    print(items)
    assert sorted(items) == [1, 2, 3, 4, 5]


asyncio.run(main())

有没有办法实现上面的get_all_queue_items,这样我们就可以await <something> 来获取所有项目——明确其目的是什么?即,

    …
    await queue.join()
    await asyncio.gather(*producer_tasks)
    items = await <something>
    print(items)
    assert sorted(items) == [1, 2, 3, 4, 5]

【问题讨论】:

  • 我不确定你的问题是否完全有道理。如果生产者永远不确定他们是否完成了,那么消费者怎么可能知道是否有更多的物品来。消费者必须通过某种方式确定“生产者不再生产”。标准解决方案要求每个生产者通过信号量或标志以某种方式指示它已完成,并且当所有生产者承诺他们已完成并且没有更多输入时,消费者停止。

标签: python asynchronous queue python-asyncio


【解决方案1】:

我能够使用sentinel value 来提醒消费者get_all_queue_items 没有更多的值可以从队列中获得,从而将其从循环中断开。调度get_all_queue_items 的任务可以等待,收集的物品会在里面。

import asyncio
import random


SENTINEL = object()


async def add_queue_item(item, queue):
    # simulate some work
    sleep_interval = random.randint(1, 3)
    await asyncio.sleep(sleep_interval)
    output_item = item + 1
    await queue.put(output_item)


async def get_all_queue_items(queue):
    items = []
    item = await queue.get()
    while item is not SENTINEL:
        items.append(item)
        queue.task_done()
        item = await queue.get()
    queue.task_done()
    return items


async def main():
    queue = asyncio.Queue()
    producer_tasks = [asyncio.create_task(add_queue_item(item, queue)) for item in range(5)]
    collect_queue_items_task = asyncio.create_task(get_all_queue_items(queue))
    await asyncio.gather(*producer_tasks)
    await queue.put(SENTINEL)
    await queue.join()
    items = await collect_queue_items_task
    print(items)
    assert sorted(items) == [1, 2, 3, 4, 5]


asyncio.run(main())

【讨论】:

  • 在这个设计中你根本不需要await queue.join()(或queue.task_done())。 queue.join 用于当您停止生产并且没有有哨兵时,您需要一种方法来了解消费者何时处理(而不仅仅是出列)所有项目。
猜你喜欢
  • 1970-01-01
  • 2015-05-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-03-02
  • 2022-11-01
  • 1970-01-01
  • 2017-10-27
相关资源
最近更新 更多