Python asyncio 队列(Queue/PriorityQueue/LifoQueue)完全指南:从生产者消费者模型到优雅关闭
【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython
本文以 CPython 官方文档 Doc/library/asyncio-queue.rst 为核心骨架,结合 asyncio 队列的源码实现(Lib/asyncio/queues.py)与官方测试(Lib/test/test_asyncio/test_queues.py)展开讲解。读完本文,你将掌握 asyncio 三种队列类(Queue、PriorityQueue、LifoQueue)的完整 API、阻塞与非阻塞操作、join()/task_done()任务协作机制、带超时的队列操作技巧,以及 Python 3.13 引入的shutdown()优雅关闭机制,并能在真实异步程序中落地生产者—消费者分发模型。
说明:本文所述 API 行为以当前仓库(开发版
PY_VERSION为 "3.16.0a0",见 Include/patchlevel.h)为准,其中Queue.shutdown()与QueueShutDown自 Python 3.13 起可用(见 Doc/whatsnew 相关说明与文档版本标注)。
一、asyncio 队列是什么:与线程安全队列模块的关系
asyncio 队列被设计为与标准库queue模块的类高度相似,但二者有本质区别:
- asyncio 队列不是线程安全的,它们专门用于 async/await 代码中,在单线程事件循环内协调协程任务;
- 线程安全版本
queue.Queue依赖锁机制在多个线程之间传递数据;而 asyncio 队列通过Future 等待/唤醒机制,让生产协程与消费协程在事件循环内协作,无需任何锁(源码中完全没有使用threading.Lock,而是内部维护_getters与_putters两个 future 双向队列,见 Lib/asyncio/queues.py)。
与线程安全 queue 相比的显著差异
| 对比维度 | asyncio 队列 | threadingqueue模块 |
|---|---|---|
| 适用场景 | async/await 单线程协程协作 | 多线程数据交换 |
| 线程安全 | 否(:ref:not thread safe`) | 是 |
qsize()可靠性 | 总是可知且准确 | 仅表示当时近似值 |
| 阻塞语义 | await put()/get()挂起协程 | put()/get()阻塞线程 |
| 方法超时参数 | 无 timeout 参数,需借助asyncio.wait_for() | 支持timeout=参数 |
关于qsize():文档明确指出,不同于线程版本队列,asyncio 队列的大小总是可知的,可直接调用qsize()返回。这是因为单线程的 asyncio 应用在调用qsize()与后续队列操作之间不会被其他执行流"插队打断"。源码中qsize()就是len(self._queue)(queues.py),而_queue在基类中是一个collections.deque。
关键注意点:asyncio 队列的方法都没有timeout参数。如果需要对队列操作设置超时,应使用asyncio.wait_for()函数包裹,例如:
# 最多等待 5 秒取一个元素,超时抛出 asyncio.TimeoutError item = await asyncio.wait_for(queue.get(), timeout=5.0)官方测试对此有直接验证:test_cancelled_getters_not_being_held_in_self_getters用asyncio.wait_for(queue.get(), 0.1)制造超时后,断言queue._getters中不再残留被取消的等待者,说明wait_for取消会正确清理内部等待队列(test_queues.py)。
二、三种队列类详解
1.Queue:标准 FIFO 队列
class asyncio.Queue(maxsize=0)构造参数maxsize:
maxsize小于或等于0(默认值0):队列容量无限;maxsize为大于0的整数:当队列达到maxsize时,await put()会阻塞,直到有元素被get()移除腾出空位。
maxsize不必是整数——从源码可见比较逻辑是数值比较(self.qsize() >= self._maxsize),官方测试test_float_maxsize就使用了maxsize=1.3验证其仍可正常工作(test_queues.py)。
FIFO 语义由内部数据结构保证:基类_init()创建collections.deque,_put()追加到右端、_get()从左侧弹出(popleft()),先进先出(queues.py)。
API 一览:
| 方法/属性 | 类型 | 行为 |
|---|---|---|
maxsize | 属性 | 队列允许存放的元素数量(即构造传入的maxsize) |
empty() | 同步 | 队列为空返回True |
full() | 同步 | 队列中已有maxsize个元素时返回True;maxsize=0时永远返回False |
qsize() | 同步 | 返回队列中元素个数 |
put(item) | 协程 | 放入一个元素;队列满时挂起等待空位 |
put_nowait(item) | 同步 | 非阻塞放入;无空位立即抛QueueFull |
get() | 协程 | 取出并返回一个元素;队列空时挂起等待 |
get_nowait() | 同步 | 非阻塞取出;队列空立即抛QueueEmpty |
join() | 协程 | 阻塞直到队列中所有元素都被取出并处理完毕 |
task_done() | 同步 | 标记一个已取出的工作项处理完成 |
shutdown(immediate=False) | 同步 | 将队列置入关闭模式(Python 3.13+) |
关于full():源码中当self._maxsize <= 0直接返回False,因此用默认maxsize=0初始化的队列full()永不返回True(queues.py)。
版本变更:Python 3.10 起删除了Queue的loop参数,事件循环通过_get_loop()自动绑定(类继承自mixins._LoopBoundMixin,会校验队列与当前运行事件循环一致,否则抛RuntimeError,见 Lib/asyncio/mixins.py)。该队列非线程安全,跨事件循环或跨线程使用会造成不可预期行为。
类型标注支持:Queue实现了__class_getitem__(通过GenericAlias),支持asyncio.Queue[int]形式的泛型下标;官方测试test_generic_alias验证了这一点(test_queues.py)。
2.PriorityQueue:按优先级取出的变体
asyncio.PriorityQueue是Queue的子类,按优先级顺序取出元素(数值最小优先)。条目通常为(priority_number, data)形式的元组。
其内部实现覆写了基类的三个钩子方法(queues.py):
def _init(self, maxsize): self._queue = [] # 用列表替代 deque def _put(self, item, heappush=heapq.heappush): heappush(self._queue, item) # 小顶堆入堆 def _get(self, heappop=heapq.heappop): return heappop(self._queue) # 弹出堆顶(最小元素)从源码结构可以推断:PriorityQueue底层是 Python 标准库heapq实现的小顶堆,因此插入与取出的时间复杂度均为 O(log n)。官方测试PriorityQueueTests.test_order验证了"1、3、2"入队后按"1、2、3"出队(test_queues.py)。
使用注意:由于堆需要比较条目,
PriorityQueue中所有元素必须彼此可比较。若放入(priority, data)元组且两个元素的priority相同,Python 会继续比较第二项data——若data类型不支持比较,会抛TypeError。实践中可给优先级元组加第三项序号或用包装类规避。
3.LifoQueue:后进先出(栈)变体
asyncio.LifoQueue同样是Queue子类,最先取出最近加入的元素(后进先出)。其实现同样覆写三个钩子:底层用列表,_put追加到尾部、_get从尾部pop()(queues.py)。官方测试验证"1、3、2"入队后按"2、3、1"出队(test_queues.py)。
PriorityQueue、LifoQueue与Queue共享同一套阻塞/唤醒、join()/task_done()、shutdown()逻辑——测试代码通过_QueueJoinTestMixin和_QueueShutdownTestMixin将同一批测试跑在三种类上,间接证明了这一点(test_queues.py)。
三、阻塞与非阻塞 API 的完整语义
put / put_nowait:入队
await queue.put(item) # 队列满则挂起,直到空位出现或队列被 shutdown queue.put_nowait(item) # 队列满立即抛 QueueFull,不阻塞put()的实现是一个while self.full()循环:一旦满,就创建一个 future 追加到_putters中挂起等待;当某个get()取走元素时会通过_wakeup_next(self._putters)唤醒队首等待者(queues.py)。若在等待期间协程被取消,源码会清理_putters中已取消的 future;若出现多个等待者同时被取消的竞态,还会礼貌地唤醒排队中的下一个(_wakeup_next)。
get / get_nowait:出队
item = await queue.get() # 队列空则挂起,直到有元素入队或队列被 shutdown item = queue.get_nowait() # 队列空立即抛 QueueEmpty,不阻塞get()与put()对称:空队列时创建 future 挂入_getters,put_nowait()入队成功后唤醒队首等待者(queues.py)。
官方测试覆盖了若干竞态细节:
test_get_cancelled_race/test_put_cancelled_race:取消等待者不会破坏队列顺序,后续元素不丢失(test_queues.py);test_cancelled_put_silence_value_error_exception:put()被取消时若 future 已被get_nowait()移除,重复移除引发的ValueError会被吞掉,协程仍以CancelledError正常结束(test_queues.py)。
非阻塞方法对应的异常
QueueEmpty:对空队列调用get_nowait()时抛出;QueueFull:队列已达maxsize时调用put_nowait()抛出。
两个异常类与队列类共同定义并导出(__all__中列有QueueFull、QueueEmpty、QueueShutDown,见 queues.py)。官方测试对它们逐一验证,例如test_nonblocking_get_exception与test_nonblocking_put_exception(test_queues.py)。
四、join() + task_done():让主协程等待全部工作完成
这是 asyncio 队列生产—消费模型中最关键的协作机制。
工作原理解析
task_done():由消费者协程调用,表示一个此前通过get()取出的工作项已被完整处理。每完成一个取出的元素,都应调用一次task_done()。join():阻塞当前协程,直到队列中所有元素都已被取出且处理完毕。
底层机制(结合 queues.py 源码):
- 内部维护计数器
self._unfinished_tasks与一个locks.Event类型的self._finished; - 每次
put_nowait()入队成功,_unfinished_tasks += 1并_finished.clear(); - 每次
task_done()将计数器减 1,当减到 0 时调用_finished.set(); join()在计数器大于 0 时await self._finished.wait(),一旦归零立即解除阻塞。
容易踩的坑:
task_done()调用次数多于已入队元素数时,会抛ValueError(源码检查_unfinished_tasks <= 0,测试test_task_done_underflow验证,见 test_queues.py);- 只
get()不task_done(),join()永远不会返回; task_done()与get()的数量必须一一对应。
官方经典示例:多 Worker 分发任务
官方文档给出的完整示例演示了"一个生产者入队 + 三个并发 worker 消费 + 主协程 join 等待全部完成"的标准模式(文档 Examples 节),可直接运行:
import asyncio import random import time async def worker(name, queue): while True: # 从队列取出一个 "work item"。 sleep_for = await queue.get() # 模拟处理:睡眠 sleep_for 秒。 await asyncio.sleep(sleep_for) # 通知队列:该 "work item" 已处理完毕。 queue.task_done() print(f'{name} has slept for {sleep_for:.2f} seconds') async def main(): # 创建用于存放 "workload" 的队列。 queue = asyncio.Queue() # 生成 20 个随机睡眠时长并入队。 total_sleep_time = 0 for _ in range(20): sleep_for = random.uniform(0.05, 1.0) total_sleep_time += sleep_for queue.put_nowait(sleep_for) # 创建三个 worker 任务并发消费队列。 tasks = [] for i in range(3): task = asyncio.create_task(worker(f'worker-{i}', queue)) tasks.append(task) # 等待队列被完全处理完毕。 started_at = time.monotonic() await queue.join() total_slept_for = time.monotonic() - started_at # 取消所有 worker 任务。 for task in tasks: task.cancel() # 等待 worker 任务全部取消完成。 await asyncio.gather(*tasks, return_exceptions=True) print('====') print(f'3 workers slept in parallel for {total_slept_for:.2f} seconds') print(f'total expected sleep time: {total_sleep_time:.2f} seconds') asyncio.run(main())该示例的核心手法:worker 是"死循环 +await queue.get()",主协程用await queue.join()判断全部消费完成,随后取消并gather回收 worker 任务(return_exceptions=True容忍CancelledError),这是收尾的标准姿势,避免 worker 协程泄漏。示例运行后输出形如:
worker-0 has slept for 0.32 seconds worker-1 has slept for 0.55 seconds ... ==== 3 workers slept in parallel for X.XX seconds total expected sleep time: Y.YY seconds其中总耗时远小于串行总睡眠时间,直观展示多任务并行收益。Queue.join()与task_done()的完整语义(包括多 worker 交替消费 100 个元素并正确累加的测试test_task_done)可在 test_queues.py 中对照查看。
五、shutdown() 优雅关闭机制(Python 3.13+)
从 Python 3.13 起,队列新增了shutdown(immediate=False)方法与QueueShutDown异常,用于解决一个常见难题:如何让阻塞在get()/put()上的协程在程序结束时被干净地唤醒。
基本规则
queue.shutdown(immediate=False) # 优雅排空模式(默认) queue.shutdown(immediate=True) # 立即终止模式调用shutdown()后,队列进入关闭模式,其核心规则(来自文档与 queues.py 源码):
- 队列不再增长:之后所有
put()/put_nowait()调用都抛QueueShutDown; - 阻塞中的 putter 被唤醒:当前阻塞在
put()上的协程会被解除阻塞,并在原等待处抛出QueueShutDown; - 阻塞中的 getter 被唤醒:所有阻塞的
get()协程被唤醒后,依据关闭模式决定抛出QueueShutDown还是继续取走残留元素。
immediate=False(默认):允许排空已装载任务
- 队列可以继续通过
get()正常取走已经入队的任务; - 只要对每个剩余任务调用一次
task_done(),处于 pending 状态的join()就能被正常解除阻塞; - 一旦队列被取空,后续
get()/get_nowait()一律抛QueueShutDown。
源码实现要点:shutdown(False)只设置_is_shutdown = True并唤醒全部 getter/putter;被唤醒的get()协程回到while self.empty()循环后,会因_is_shutdown为真而抛QueueShutDown——但若队列尚有元素,则正常执行get_nowait()取出任务(queues.py)。
immediate=True:立即终止
- 队列被立即排空(清空内部
_queue),未完成任务计数按被排空的任务数减少; - 若未完成任务数因此归零,阻塞中的
join()调用者被解除阻塞; - 阻塞中的
get()调用者被唤醒并抛QueueShutDown(队列已空)。
使用join()+immediate=True的警告
文档明确提醒:谨慎在immediate=True时使用join()——因为即使任务完全没有被处理(work 尚未执行),join()也会被解除阻塞,违背了 join 队列的常规不变式(invariant)。
源码确证了这一点:shutdown(immediate=True)直接循环self._get()排空队列并同步递减_unfinished_tasks,不做任何处理动作即触发_finished.set()(queues.py)。
三种场景的官方测试印证
test_queues.py 中的_QueueShutdownTestMixin被复用于Queue、LifoQueue、PriorityQueue三类,覆盖了:
test_shutdown_empty:空队列 shutdown 后,get/put全部抛QueueShutDown,join()正常完成;test_shutdown_nonempty:非空队列默认关闭后,put抛异常,但已入队的"data"仍可get()取出;配合一次task_done(),join()正常完成,多出的task_done()抛ValueError;test_shutdown_immediate:立即关闭后队列被排空,join()立即完成;test_shutdown_immediate_with_unfinished:有一条已取出但未task_done()的任务时immediate=True,join()不会立即返回,需补一次task_done()才解除——精确演示了"未完成任务数"如何参与计数。
一个结合 shutdown 的生产者—消费者收尾模板
import asyncio async def worker(name, q): while True: try: item = await q.get() except asyncio.QueueShutDown: print(f'{name}: queue shut down, exiting') return try: print(f'{name}: processing {item}') await asyncio.sleep(0.05) finally: q.task_done() async def main(): q = asyncio.Queue() for i in range(10): q.put_nowait(i) workers = [asyncio.create_task(worker(f'w{i}', q)) for i in range(3)] await q.join() # 等待 10 个任务全部处理完 q.shutdown() # 优雅关闭:唤醒仍阻塞在 q.get() 上的 worker await asyncio.gather(*workers) # worker 捕获 QueueShutDown 后自行退出 asyncio.run(main())此模式不需要手动cancel()每个 worker,而是借shutdown()让 worker 循环自然退出,逻辑更清晰、异常更可预期。
六、队列的运行时诊断与格式化输出
队列内置了调试友好的字符串表示,在日志与 REPL 中非常有用(源码_format(),queues.py):
q = asyncio.Queue() str(q) # '<Queue maxsize=0>' —— 不含内存地址 repr(q) # '<Queue at 0x7f... maxsize=0>' —— 含对象地址当队列处于不同状态时,输出会附加对应字段:
- 有元素:追加
_queue=[...]; - 有阻塞等待的 getter:追加
_getters[N]; - 有阻塞等待的 putter:追加
_putters[N]; - 有未完成任务:追加
tasks=N; - 已 shutdown:追加
shutdown。
例如官方测试test_format断言q._format() == 'maxsize=0 tasks=2',test_shutdown系列断言关闭后格式为maxsize=0 shutdown;_test_repr_or_str则验证了_getters[1]、_putters[1]、_queue=[1]等片段会出现在格式化字符串中(test_queues.py)。队列通过继承mixins._LoopBoundMixin与事件循环绑定(Lib/asyncio/mixins.py),这也是 3.10 移除loop参数后队列能自动感知运行循环的原因。
七、实践要点速查
- 容量控制:
asyncio.Queue(maxsize=N)可用于天然背压(backpressure)——当put()挂起等待时,生产协程自动放慢节奏,避免内存无限增长;maxsize<=0为无限队列,注意生产过快时无上限风险。 - 超时操作:队列方法没有 timeout,统一用
asyncio.wait_for(queue.get(), timeout=...)包裹;超时会取消内部等待的 future,且 asyncio 会正确清理_getters/_putters中的残留项(有测试保障)。 - 每 get 必有 task_done:使用
join()时必须保证每个get()取出的任务最终都对应一次task_done(),建议把task_done()放进finally或与处理逻辑紧邻,防止异常路径导致join()永久阻塞;task_done()多调会抛ValueError。 - 优先级与 LIFO 场景:任务带紧迫程度用
PriorityQueue+(priority, data)元组;需要"最近任务优先"(如栈式回溯、深度优先遍历)用LifoQueue;二者与Queue完全同构,可无缝替换。 - 优雅停机:Python 3.13+ 用
shutdown()代替手动 cancel worker 的收尾套路;若任务可能丢弃,务必先评估immediate=True对join()语义的影响。 - 不要把 asyncio 队列当线程队列用:它们绑定单一事件循环且非线程安全;跨线程传数据请用标准库
queue,跨进程请用multiprocessing或消息中间件。 - 单元测试参考:上述所有行为均有官方测试背书,路径 Lib/test/test_asyncio/test_queues.py,编写自己的队列消费者逻辑时可将其作为语义规范文档来读。
【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考