asyncio.queue 更适合协程场景,因其 put()/get() 返回 awaitable、不阻塞事件循环,而 queue.queue 会阻塞线程导致 runtimeerror;需用 create_task 启动协程,配合 gather 和 join 控制生命周期,并确保每次 get 后调用 task_done。

asyncio.Queue 为什么比普通 queue.Queue 更适合协程场景
因为 queue.Queue 是线程安全的阻塞队列,所有 put() 和 get() 都会直接阻塞当前线程;而 asyncio.Queue 提供的是协程友好的异步等待接口——put() 和 get() 返回 awaitable,不会阻塞事件循环。如果你在 async def 函数里误用 queue.Queue,程序会卡死或抛出 RuntimeError: cannot be called from a running event loop。
关键区别在于:asyncio.Queue 的 put() 在队列满时 await,get() 在空时 await,整个过程不交出线程控制权,只让出协程调度权。
如何正确启动多个生产者和消费者协程
必须确保所有生产者先启动、全部数据 put 完再通知消费者结束,否则消费者可能提前退出。常见错误是没控制好“生产完成”信号,导致消费者在队列为空后就停了,但其实还有生产者没来得及 put。
- 用
asyncio.create_task()启动所有生产者和消费者协程,不要用await逐个执行 - 生产者任务结束后调用
queue.put(None)或单独发一个哨兵值(如queue.put(StopIteration)),消费者检测到后break - 更稳妥的做法是用
asyncio.gather()等待所有生产者完成,再await queue.join()确保所有已 put 的项都被task_done()标记
示例片段:
import asyncio
<p>async def producer(queue, name):
for i in range(3):
await queue.put(f"{name}-item-{i}")
print(f"Produced: {name}-item-{i}")
await asyncio.sleep(0.1)
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"{name} consumed {item}")
await asyncio.sleep(0.15)
queue.task_done()</p><p>async def main():
q = asyncio.Queue(maxsize=2)
producers = [asyncio.create_task(producer(q, f"p{i}")) for i in range(2)]
consumers = [asyncio.create_task(consumer(q, f"c{i}")) for i in range(2)]
await asyncio.gather(*producers)
await q.join() # 等待所有已入队项被处理完</p>
maxsize 设置为 0 或负数意味着什么
maxsize=0(默认值)表示“无上限”,但注意:这不等于“无限快”,它只是不限制数量,底层仍用 collections.deque 存储,内存耗尽时会抛 MemoryError。真正需要背压控制时,必须显式设一个正整数。
-
maxsize=1:最严格的背压,生产者每次put()都要等消费者get()后才能继续 -
maxsize=0:生产者不会因队列满而暂停,容易在突发流量下撑爆内存 - 负数会被当作 0 处理,没有特殊语义
如果你观察到生产者协程迟迟不 await put(),大概率是 maxsize 设得太小,且消费者处理太慢。
忘记调用 task_done() 会导致 join() 永远不返回
queue.join() 内部依赖每个 get() 后显式调用 task_done() 来递减未完成计数。漏掉一次,计数就卡住,join() 会永远挂起——这是最隐蔽也最常踩的坑。
- 每个
get()对应且仅对应一次task_done(),哪怕消费逻辑抛异常也要包在try/finally里 - 不能靠
put()触发计数变化;也不能在put()后调用task_done() - 如果消费者做了重试或丢弃,依然要调用
task_done(),因为“取走即负责”
正确写法示例:
try:
item = await queue.get()
await process(item)
finally:
queue.task_done()
实际用的时候,asyncio.Queue 的行为边界很清晰:它只管“排队 + 协程等待”,不负责序列化、持久化、跨进程或失败重试。这些都得你自己补全——比如加锁处理共享状态、用 asyncio.Lock 保护临界区、或者把失败 item 写进另一个 error queue。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











