【发布时间】: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