news 2026/9/7 1:37:41

Python asyncio 队列(Queue/PriorityQueue/LifoQueue)完全指南:从生产者消费者模型到优雅关闭

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Python asyncio 队列(Queue/PriorityQueue/LifoQueue)完全指南:从生产者消费者模型到优雅关闭

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 三种队列类(QueuePriorityQueueLifoQueue)的完整 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_gettersasyncio.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个元素时返回Truemaxsize=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 起删除了Queueloop参数,事件循环通过_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.PriorityQueueQueue的子类,按优先级顺序取出元素(数值最小优先)。条目通常为(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)。

PriorityQueueLifoQueueQueue共享同一套阻塞/唤醒、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 挂入_gettersput_nowait()入队成功后唤醒队首等待者(queues.py)。

官方测试覆盖了若干竞态细节:

  • test_get_cancelled_race/test_put_cancelled_race:取消等待者不会破坏队列顺序,后续元素不丢失(test_queues.py);
  • test_cancelled_put_silence_value_error_exceptionput()被取消时若 future 已被get_nowait()移除,重复移除引发的ValueError会被吞掉,协程仍以CancelledError正常结束(test_queues.py)。

非阻塞方法对应的异常

  • QueueEmpty:对空队列调用get_nowait()时抛出;
  • QueueFull:队列已达maxsize时调用put_nowait()抛出。

两个异常类与队列类共同定义并导出(__all__中列有QueueFullQueueEmptyQueueShutDown,见 queues.py)。官方测试对它们逐一验证,例如test_nonblocking_get_exceptiontest_nonblocking_put_exception(test_queues.py)。

四、join() + task_done():让主协程等待全部工作完成

这是 asyncio 队列生产—消费模型中最关键的协作机制。

工作原理解析

  • task_done():由消费者协程调用,表示一个此前通过get()取出的工作项已被完整处理。每完成一个取出的元素,都应调用一次task_done()
  • join():阻塞当前协程,直到队列中所有元素都已被取出且处理完毕。

底层机制(结合 queues.py 源码):

  1. 内部维护计数器self._unfinished_tasks与一个locks.Event类型的self._finished
  2. 每次put_nowait()入队成功,_unfinished_tasks += 1_finished.clear()
  3. 每次task_done()将计数器减 1,当减到 0 时调用_finished.set()
  4. 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 源码):

  1. 队列不再增长:之后所有put()/put_nowait()调用都抛QueueShutDown
  2. 阻塞中的 putter 被唤醒:当前阻塞在put()上的协程会被解除阻塞,并在原等待处抛出QueueShutDown
  3. 阻塞中的 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被复用于QueueLifoQueuePriorityQueue三类,覆盖了:

  • test_shutdown_empty:空队列 shutdown 后,get/put全部抛QueueShutDownjoin()正常完成;
  • test_shutdown_nonempty:非空队列默认关闭后,put抛异常,但已入队的"data"仍可get()取出;配合一次task_done()join()正常完成,多出的task_done()ValueError
  • test_shutdown_immediate:立即关闭后队列被排空,join()立即完成;
  • test_shutdown_immediate_with_unfinished:有一条已取出但未task_done()的任务时immediate=Truejoin()不会立即返回,需补一次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参数后队列能自动感知运行循环的原因。

七、实践要点速查

  1. 容量控制asyncio.Queue(maxsize=N)可用于天然背压(backpressure)——当put()挂起等待时,生产协程自动放慢节奏,避免内存无限增长;maxsize<=0为无限队列,注意生产过快时无上限风险。
  2. 超时操作:队列方法没有 timeout,统一用asyncio.wait_for(queue.get(), timeout=...)包裹;超时会取消内部等待的 future,且 asyncio 会正确清理_getters/_putters中的残留项(有测试保障)。
  3. 每 get 必有 task_done:使用join()时必须保证每个get()取出的任务最终都对应一次task_done(),建议把task_done()放进finally或与处理逻辑紧邻,防止异常路径导致join()永久阻塞;task_done()多调会抛ValueError
  4. 优先级与 LIFO 场景:任务带紧迫程度用PriorityQueue+(priority, data)元组;需要"最近任务优先"(如栈式回溯、深度优先遍历)用LifoQueue;二者与Queue完全同构,可无缝替换。
  5. 优雅停机:Python 3.13+ 用shutdown()代替手动 cancel worker 的收尾套路;若任务可能丢弃,务必先评估immediate=Truejoin()语义的影响。
  6. 不要把 asyncio 队列当线程队列用:它们绑定单一事件循环且非线程安全;跨线程传数据请用标准库queue,跨进程请用multiprocessing或消息中间件。
  7. 单元测试参考:上述所有行为均有官方测试背书,路径 Lib/test/test_asyncio/test_queues.py,编写自己的队列消费者逻辑时可将其作为语义规范文档来读。

【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/7 1:36:07

Linux---基本指令

前言1.windows系统中&#xff0c;标识文件唯一性&#xff0c;是通过路径标识的。桌面其实也是一个文件夹&#xff0c;只不过用图形化界面的方式显示出来&#xff0c;方便操作。2.无论是windows还是linux&#xff0c;一旦登录&#xff0c;就会处在一个默认的路径下。那windows默…

作者头像 李华
网站建设 2026/9/7 1:35:16

MicroPython高频采样优化:DMA+乒乓缓冲彻底解放CPU

先说一个我自己的真实经历。去年我在 ESP32-S3 上跑 MicroPython 做一个小型环境监测节点&#xff0c;需要同时采集 4 路模拟信号&#xff0c;频率还不能太低&#xff0c;同时还要驱动一块小 OLED 做实时显示。一开始图省事&#xff0c;直接在 while True 里调 read_u16() …

作者头像 李华
网站建设 2026/9/7 1:35:10

PHC硬件时钟操作与时钟调整:PTP同步的核心实战指南

做PTP时间同步的兄弟应该都有这种经历&#xff1a;ptp4l日志刷得飞起&#xff0c;master offset看着也是几十纳秒&#xff0c;可一到验收&#xff0c;抓包软件打出来的时间戳还是稀巴烂。问题往往不在协议跑没跑通&#xff0c;而在你根本没和网卡里那只表——PHC&#xff08;PT…

作者头像 李华
网站建设 2026/9/7 1:34:12

CPH插件实战指南:用VS Code高效刷LeetCode,从安装到提交一步到位

1. CPH 是什么&#xff0c;为什么刷 LeetCode 需要它1.1 很多刷题党都遇到过的低效场景先聊一个很常见的场景。很多同学刷 LeetCode 时&#xff0c;习惯性地打开浏览器&#xff0c;进入 LeetCode 题目页面&#xff0c;读完题后在网页右侧的内嵌编辑器里写代码&#xff0c;然后点…

作者头像 李华
网站建设 2026/9/7 1:33:24

HL7 V3 Schema实战解析:消息校验、代码生成与避坑指南

简介&#xff1a;HL7 V3 Schema是医疗信息化领域实现标准化数据交换的关键资源&#xff0c;面向医疗软件开发者、系统集成工程师以及从事HL7标准实施的技术人员。压缩包共48个文件&#xff0c;以xsd模式定义文件为主&#xff0c;辅以dtd、xml、doc、vsd、xls等文档与图形说明&a…

作者头像 李华
网站建设 2026/9/7 1:32:19

STM32H725ZGT6深度解析:550MHz Cortex-M7高性能MCU实战指南

STM32H725ZGT6 这颗料&#xff0c;我第一次拿到手的时候其实没太当回事——毕竟 H7 系列已经出了好几年&#xff0c;H743、H750 这些老朋友大家都熟。但真正点开数据手册&#xff0c;看到主频 550MHz 那一栏的时候&#xff0c;我还是愣了一下&#xff1a;ST 居然把一颗 Cortex-…

作者头像 李华