【问题标题】:Python Asyncio producer-consumer workflow congestion / growing queuePython Asyncio生产者-消费者工作流拥塞/增长队列
【发布时间】:2021-03-10 10:38:35
【问题描述】:

我一直在编写一个 Python 应用程序,其中:

  • 有一个async 函数producer 通过websocket 监听传入的items,并将这些items 放入queue = asyncio.Queue()。
  • 有一个async 函数consumer 执行queue.get(),并通过不同的websocket 连接查询item_details。

问题: 将传入的items 放入queue 的平均速度远高于consumer 到get 项从queue 的速度,并且因此,queue 会在一段时间后堆积起来。

问题:扩大consumer 不进行多处理和不限制传入连接的正确方法是什么?我对asyncio 和threading 还不是很精通。我想过在单独的工作人员中运行consumer,但据我了解,asyncio 的run_in_executor 不能用于async 函数,还有asyncio.Queue() 不是线程安全的。

【问题讨论】:

  • 如果您使用异步,您必须使用信号量或异步限制器,它们实际上可以根据您的要求限制传入连接的数量。
  • @mac_online 抱歉,我在这一点上不清楚。对于要处理的传入项目,速度对于这个应用程序来说真的很重要,所以对我来说,不可能限制传入的连接。我会将这个添加到问题中。
  • 你可以在同一个线程中为同一个队列运行多个消费者
  • 是的,类似的。如果您的任务是 IO 绑定的,您可以创建额外的消费协程,它应该会有所帮助。但是你最好把这些任务放在某种集合中——集合或列表,这样你就可以实现优雅的关闭或类似的东西。不要无处运行任务,它是不可管理的。
  • 谢谢@NobbyNobbs!我很快用create_task 对其进行了测试,这实际上解决了n 的数量。我还将检查正常关机,感谢您的建议。如果您写了答案,我会将其标记为已接受。

标签: python multithreading python-asyncio python-multithreading producer-consumer


【解决方案1】:

如果消费者执行 IO 密集型工作,您可以扩展其计数。而且你并不关心多线程,因为asyncio 基于非阻塞IO 的思想,并且设计为在单线程中工作。如果没有本机异步替代方案,例如,您可以甚至必须使用线程来处理阻塞 IO。对于文件 IO,但这是一个单独的故事。

这里有一个简单的例子来说明生产者创建任务的速度快于单个消费者处理它们的速度。我用 asyncio.sleep 模拟 IO 工作负载。

import asyncio
import itertools

async def producer(queue: asyncio.Queue):
    """producer emulator, creates ~ 10 tasks per second"""
    sleep_seconds=0.1
    counter = itertools.count(1)
    while True:
        await queue.put(next(counter))
        await asyncio.sleep(sleep_seconds)


async def consumer(queue: asyncio.Queue, index):
    """slow io-bound consumer emulator, process ~ 5 tasks per second"""
    sleep_seconds=0.2
    while True:
        task = await queue.get()
        print(f"consumer={index}, task={task}, queue_size={queue.qsize()}")
        await asyncio.sleep(sleep_seconds)


async def main():
    q = asyncio.Queue()
    concurrency = 2  # consumers count
    tasks = [asyncio.create_task(consumer(q, i)) for i in range(concurrency)]
    tasks += [asyncio.create_task(producer(q))]
    await asyncio.wait(tasks)


if __name__ == "__main__":
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        pass

单个消费者的输出,队列大小不断增长

consumer=0, task=1, queue_size=0
consumer=0, task=2, queue_size=0
consumer=0, task=3, queue_size=1
consumer=0, task=4, queue_size=2
consumer=0, task=5, queue_size=3
consumer=0, task=6, queue_size=4
consumer=0, task=7, queue_size=5
consumer=0, task=8, queue_size=6
consumer=0, task=9, queue_size=7
consumer=0, task=10, queue_size=8

两个消费者的输出,队列为空

consumer=0, task=1, queue_size=0
consumer=1, task=2, queue_size=0
consumer=0, task=3, queue_size=0
consumer=1, task=4, queue_size=0
consumer=0, task=5, queue_size=0
consumer=1, task=6, queue_size=0
consumer=0, task=7, queue_size=0
consumer=1, task=8, queue_size=0
consumer=0, task=9, queue_size=0
consumer=1, task=10, queue_size=0

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-01-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多