asyncio.Queue是协程安全的队列,必须await put/get,不可用同步queue.Queue;需用sentinel或task_done/join管理结束;maxsize实现背压,避免OOM。

asyncio.Queue 是 asyncio 生产者-消费者模型的首选通信机制,它线程安全、协程友好,且天然支持 await,但直接照搬同步 Queue 的用法会卡死或漏数据。
为什么不能直接用 queue.Queue 替代 asyncio.Queue
同步 queue.Queue 的 get() 和 put() 是阻塞调用,放进 async def 里会阻塞整个事件循环;而 asyncio.Queue 的对应方法是协程,必须 await。常见错误是忘记 await,导致返回一个 coroutine 对象却没执行,后续 get() 拿不到东西,程序看似“卡住”。
-
q = asyncio.Queue()创建的是异步队列,不是同步队列的包装 -
await q.put(item)和await q.get()是必须的,写成q.put(item)会静默失败 -
q.qsize()在协程中可读,但不保证原子性,仅作调试参考,不可用于条件判断
asyncio.Queue 的关闭信号怎么传:用 sentinel 还是 task_done() + join()
没有内置“关闭”信号,需约定结束方式。最常用的是放一个哨兵值(如 None),消费者收到即退出;task_done() 和 join() 更适合等待所有已入队任务完成,而非终止消费。
- 多个生产者时,每个都应放入一个 sentinel,消费者按 sentinel 数量退出(例如 3 个生产者 → 收到 3 个
None后停) -
await q.join()会等待所有已get()出来的 item 被task_done()标记完成,和“停止消费”无关 - 避免用
q.empty()判断退出,因为可能刚判空、生产者就 put 了
如何防止消费者提前退出或死锁
典型死锁场景:消费者在 await q.get() 等待时,所有生产者已结束但没发 sentinel,消费者永远挂起。关键是要确保生产者显式结束,并统一管理生命周期。
立即学习“Python免费学习笔记(深入)”;
- 用
asyncio.create_task()启动生产者,再await asyncio.gather(*producers)等它们全部完成 - 生产者结束后,立即
await q.put(sentinel),不要依赖 finally 块——如果生产者协程被 cancel,finally 可能不执行 - 消费者内部用
try/except CancelledError清理资源,但不要吞掉异常,否则asyncio.wait_for等超时控制会失效
性能与边界注意点:maxsize、内存增长与背压
asyncio.Queue 默认无上限(maxsize=0),大量生产者快速 put 会导致内存暴涨甚至 OOM。启用 maxsize 后,put 会自动 await 等待空间,形成天然背压。
-
asyncio.Queue(maxsize=100)是合理起点,根据处理延迟和内存预算调整 -
q.full()和q.empty()返回布尔值,但状态瞬时变化,仅适合日志或监控,不能用于控制流 - 如果消费者处理慢,队列积压会拖慢生产者,这是设计意图,不是 bug;若需跳过旧数据,得自己实现环形缓冲或用
asyncio.Queue(maxsize=1)配合put_nowait()+try/except asyncio.QueueFull
真正难的不是写通逻辑,而是想清楚:谁负责发结束信号、什么时候发、发几个;以及是否允许队列无限增长——这决定了你是在做流式处理,还是在攒批处理。


















