asyncio.queue更适合协程场景,因其put()/get()是awaitable且不阻塞事件循环;而queue.queue为同步阻塞,误用会导致协程卡死。

asyncio.Queue 为什么比普通 queue.Queue 更适合协程场景
因为 asyncio.Queue 的 put() 和 get() 都是 awaitable,不会阻塞事件循环;而 queue.Queue 的 put()/get() 是同步阻塞调用,一旦用在 async def 函数里,会直接卡死整个协程调度。常见错误现象是:消费者任务看似启动了,但永远收不到数据,或程序卡在某个 q.put(item) 不动——那基本是因为误用了线程安全的 queue.Queue。
使用场景明确:所有生产者、消费者都必须是协程(async def),且运行在同一个 asyncio 事件循环中。
性能影响:默认无最大长度(maxsize=0),写入永不阻塞;设了 maxsize 后,put() 在满时会 await,自动实现背压;这点比手动 sleep 或信号量更自然。
如何正确启动多个生产者 + 多个消费者并共用一个队列
关键不是“怎么创建”,而是“怎么协调退出”。常见错误是消费者不知道何时停止:没有收到结束信号,又没检测到队列空,结果 await q.get() 永远挂起。
- 不要依赖
q.empty()判断退出——它只是快照,不可靠;也不要在生产者结束后立即break - 推荐做法:生产者全部完成时,向队列放入
None(或其他 sentinel 值)作为结束标记,每个消费者收到一次就退出 - 或者用
asyncio.create_task()启动所有消费者,再用asyncio.gather(*producers)等待所有生产者结束,最后调用q.join()确保所有已入队任务被处理完
示例片段:
import asyncio
<p>async def producer(q: asyncio.Queue, name: str):
for i in range(3):
await q.put(f"{name}-item-{i}")
await asyncio.sleep(0.1)</p><h1>生产完毕,发结束信号(可选)</h1><pre class="brush:python;toolbar:false;">await q.put(None)async def consumer(q: asyncio.Queue, name: str): while True: item = await q.get() if item is None: q.task_done() break print(f"{name} got {item}") q.task_done()
async def main(): q = asyncio.Queue() 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(3)] await asyncio.gather(*producers) await q.join() # 等待所有已入队任务被 task_done()
q.task_done() 和 q.join() 必须成对出现,否则 await q.join() 永不返回
这是最容易漏掉的一步。q.join() 内部依赖未完成任务计数,而这个计数只在 q.get() 时加一、q.task_done() 时减一。如果消费者拿到 item 后忘记调用 q.task_done(),q.join() 就会一直等待。
SkillSub Pro - Python 题解与代码注释双功能技能功能概述SkillSub Pro - Python 题解与代码注释双功能技能是一项面向实际任务的技能,主要用于SkillSub Pro 是一个 Python 题解生成与代码注释的 双功能合体技能 ,专为学生、算法学习者和开发者设计;✅ 一个技能,两种用途 :;核心要点📝 题解模式 :输入题目/题号,自动生成完整 Python 题解(含详细注释、解题思路、复杂度分析);💬 注释模式 :输入 Python 代码,自动添加详细中。它将相关步骤、
注意点:
-
q.task_done()必须在处理完 item 后调用,不能放在try/except外侧——异常时也得确保调用,否则计数失准 - 每个
get()对应且仅对应一次task_done();重复调用会引发ValueError: task_done() called too many times - 不用
q.join()也能跑,但无法可靠判断“所有任务已处理完毕”,尤其在测试或资源清理阶段容易出问题
当需要限制并发消费数量时,别用 Queue.maxsize 控制,改用 asyncio.Semaphore
asyncio.Queue(maxsize=N) 控制的是队列容量(缓冲区大小),不是同时处理的任务数。想限制最多 3 个消费者在运行中,应该用 asyncio.Semaphore(3) 包裹消费逻辑,而不是把 Queue 设成 maxsize=3。
原因很直接:Queue 满了只会让生产者 await,不影响消费者数量;而 Semaphore 才真正控制协程进入临界区的并发度。
示例:
sem = asyncio.Semaphore(2)
<p>async def limited_consumer(q, name):
while True:
item = await q.get()
if item is None:
q.task_done()
break
async with sem: # 最多 2 个协程同时执行下面的逻辑
await asyncio.sleep(0.5) # 模拟耗时处理
print(f"{name} processed {item}")
q.task_done()</p>
实际项目里,Queue 负责解耦和缓冲,Semaphore 负责资源节流,两者职责分明。混用或误用 maxsize 是调试时最耗时间的坑之一。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










