asyncio.queue不能用list替代,因其是协程安全的异步队列,提供awaitable的put/get接口、内置同步机制及哨兵控制;list无等待能力、无相应方法且非线程安全。

asyncio.Queue 为什么不能直接用 list 替代
因为 asyncio.Queue 是线程安全且协程友好的,内部自带 put() 和 get() 的 awaitable 接口,能自动挂起协程直到队列非空/未满;而普通 list 没有等待能力,直接 pop(0) 或 append() 会阻塞事件循环,导致其他任务饿死。
常见错误现象:RuntimeWarning: coroutine 'Queue.get' was never awaited——忘了加 await;或者用 queue.append(item) 后调 await queue.get(),结果报 AttributeError,因为 list 没有 get 方法。
-
asyncio.Queue初始化时可设maxsize,超限时put()会 await 直到有空间 - 它不继承自
collections.deque,没有extend()、index()等方法 - 不支持索引访问(
queue[0]报错),只提供get()、put()、qsize()、empty()、full()
如何正确启动生产者和消费者协程并避免提前退出
关键在控制生命周期:消费者不能因队列暂时为空就结束,生产者也不能发完就立即 return,否则消费者可能永远等不到 “结束信号”。常用做法是让生产者发完后放入一个哨兵值(如 None),消费者收到后 break。
容易踩的坑:asyncio.run(main()) 中的 main() 如果只 create_task 了生产者,没显式 await 消费者,程序可能秒退——因为主协程结束,整个事件循环就关了。
- 用
asyncio.gather()并发运行多个生产者/消费者,确保全部完成才退出 - 消费者要用
while True:+item = await queue.get()循环,别写成for item in queue:(语法错误) - 记得调
queue.task_done(),否则join()会永远卡住(虽非必须,但配合join()控制结束时很关键)
asyncio.Queue 在高并发下的性能注意事项
asyncio.Queue 内部用 asyncio.Lock 和 asyncio.Event 实现同步,轻量但非零开销。当每秒吞吐量超过 10 万条消息时,队列本身可能成为瓶颈,尤其在大量 put()/get() 频繁切换时。
典型表现:CPU 使用率不高,但整体吞吐上不去,qsize() 长期接近 maxsize,说明消费者处理不过来或生产者太猛。
- 避免在
get()后立刻put()回同一队列(比如做中间转换),这会放大锁竞争 - 如果只是转发,考虑用
asyncio.create_task()直接 spawn 处理协程,绕过队列中转 -
maxsize=0(无界)看似省心,但内存失控风险大;建议设合理上限,配合queue.full()做降级(如丢弃低优先级任务)
完整可运行示例:带超时与异常传播的最小闭环
下面这段代码能直接复制运行,模拟 2 个生产者、3 个消费者,1 秒后自动停止:
import asyncio
import random
<p>async def producer(queue, name):
for i in range(5):
await asyncio.sleep(random.uniform(0.1, 0.3))
item = f"{name}-item-{i}"
await queue.put(item)
print(f"Produced: {item}")
await queue.put(None) # 哨兵</p><p>async def consumer(queue, name):
while True:
item = await queue.get()
if item is None:
queue.task_done()
break
print(f"Consumed by {name}: {item}")
await asyncio.sleep(random.uniform(0.2, 0.5))
queue.task_done()</p><p>async def main():
queue = asyncio.Queue(maxsize=10)
producers = [asyncio.create_task(producer(queue, f"P{i}")) for i in range(2)]
consumers = [asyncio.create_task(consumer(queue, f"C{i}")) for i in range(3)]
await asyncio.gather(*producers)
await queue.join() # 等待所有已入队任务被处理完
for c in consumers:
c.cancel()</p><p>asyncio.run(main())
</p>
注意 queue.join() 和 task_done() 的配对关系——漏掉任何一个,join() 就不会返回。这是最常被忽略的同步点。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











