完全指南:从生产者消费者模型到优雅关闭)
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(), timeout5.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(maxsize0)构造参数maxsizemaxsize小于或等于0默认值0队列容量无限maxsize为大于0的整数当队列达到maxsize时await put()会阻塞直到有元素被get()移除腾出空位。maxsize不必是整数——从源码可见比较逻辑是数值比较self.qsize() self._maxsize官方测试test_float_maxsize就使用了maxsize1.3验证其仍可正常工作test_queues.py。FIFO 语义由内部数据结构保证基类_init()创建collections.deque_put()追加到右端、_get()从左侧弹出popleft()先进先出queues.py。API 一览方法/属性类型行为maxsize属性队列允许存放的元素数量即构造传入的maxsizeempty()同步队列为空返回Truefull()同步队列中已有maxsize个元素时返回Truemaxsize0时永远返回Falseqsize()同步返回队列中元素个数put(item)协程放入一个元素队列满时挂起等待空位put_nowait(item)同步非阻塞放入无空位立即抛QueueFullget()协程取出并返回一个元素队列空时挂起等待get_nowait()同步非阻塞取出队列空立即抛QueueEmptyjoin()协程阻塞直到队列中所有元素都被取出并处理完毕task_done()同步标记一个已取出的工作项处理完成shutdown(immediateFalse)同步将队列置入关闭模式Python 3.13关于full()源码中当self._maxsize 0直接返回False因此用默认maxsize0初始化的队列full()永不返回Truequeues.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.pydef _init(self, maxsize): self._queue [] # 用列表替代 deque def _put(self, item, heappushheapq.heappush): heappush(self._queue, item) # 小顶堆入堆 def _get(self, heappopheapq.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 挂入_gettersput_nowait()入队成功后唤醒队首等待者queues.py。官方测试覆盖了若干竞态细节test_get_cancelled_race/test_put_cancelled_race取消等待者不会破坏队列顺序后续元素不丢失test_queues.pytest_cancelled_put_silence_value_error_exceptionput()被取消时若 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_exceptiontest_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(fworker-{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_exceptionsTrue) print() print(f3 workers slept in parallel for {total_slept_for:.2f} seconds) print(ftotal expected sleep time: {total_sleep_time:.2f} seconds) asyncio.run(main())该示例的核心手法worker 是死循环 await queue.get()主协程用await queue.join()判断全部消费完成随后取消并gather回收 worker 任务return_exceptionsTrue容忍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(immediateFalse)方法与QueueShutDown异常用于解决一个常见难题如何让阻塞在get()/put()上的协程在程序结束时被干净地唤醒。基本规则queue.shutdown(immediateFalse) # 优雅排空模式默认 queue.shutdown(immediateTrue) # 立即终止模式调用shutdown()后队列进入关闭模式其核心规则来自文档与 queues.py 源码队列不再增长之后所有put()/put_nowait()调用都抛QueueShutDown阻塞中的 putter 被唤醒当前阻塞在put()上的协程会被解除阻塞并在原等待处抛出QueueShutDown阻塞中的 getter 被唤醒所有阻塞的get()协程被唤醒后依据关闭模式决定抛出QueueShutDown还是继续取走残留元素。immediateFalse默认允许排空已装载任务队列可以继续通过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。immediateTrue立即终止队列被立即排空清空内部_queue未完成任务计数按被排空的任务数减少若未完成任务数因此归零阻塞中的join()调用者被解除阻塞阻塞中的get()调用者被唤醒并抛QueueShutDown队列已空。使用join()immediateTrue的警告文档明确提醒谨慎在immediateTrue时使用join()——因为即使任务完全没有被处理work 尚未执行join()也会被解除阻塞违背了 join 队列的常规不变式invariant。源码确证了这一点shutdown(immediateTrue)直接循环self._get()排空队列并同步递减_unfinished_tasks不做任何处理动作即触发_finished.set()queues.py。三种场景的官方测试印证test_queues.py 中的_QueueShutdownTestMixin被复用于Queue、LifoQueue、PriorityQueue三类覆盖了test_shutdown_empty空队列 shutdown 后get/put全部抛QueueShutDownjoin()正常完成test_shutdown_nonempty非空队列默认关闭后put抛异常但已入队的data仍可get()取出配合一次task_done()join()正常完成多出的task_done()抛ValueErrortest_shutdown_immediate立即关闭后队列被排空join()立即完成test_shutdown_immediate_with_unfinished有一条已取出但未task_done()的任务时immediateTruejoin()不会立即返回需补一次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(fw{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.pyq asyncio.Queue() str(q) # Queue maxsize0 —— 不含内存地址 repr(q) # Queue at 0x7f... maxsize0 —— 含对象地址当队列处于不同状态时输出会附加对应字段有元素追加_queue[...]有阻塞等待的 getter追加_getters[N]有阻塞等待的 putter追加_putters[N]有未完成任务追加tasksN已 shutdown追加shutdown。例如官方测试test_format断言q._format() maxsize0 tasks2test_shutdown系列断言关闭后格式为maxsize0 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(maxsizeN)可用于天然背压backpressure——当put()挂起等待时生产协程自动放慢节奏避免内存无限增长maxsize0为无限队列注意生产过快时无上限风险。超时操作队列方法没有 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 的收尾套路若任务可能丢弃务必先评估immediateTrue对join()语义的影响。不要把 asyncio 队列当线程队列用它们绑定单一事件循环且非线程安全跨线程传数据请用标准库queue跨进程请用multiprocessing或消息中间件。单元测试参考上述所有行为均有官方测试背书路径 Lib/test/test_asyncio/test_queues.py编写自己的队列消费者逻辑时可将其作为语义规范文档来读。【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考