Python 基础体系 · 第 58/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 队列与背压:queue、asyncio.Queue、容量和关闭协议

队列不是“把对象临时放进去”的容器,而是并发组件之间的交接协议。它至少同时回答四个问题:

  1. 谁负责生产任务,谁负责消费任务?
  2. 队列满时,生产者是等待、丢弃,还是失败?
  3. 消费者取出任务后,什么时候才算真正完成?
  4. 系统停止时,如何让生产者、消费者和等待方都结束,而不是互相等待?

Python 3.14 中,线程场景主要使用 queue.Queue,异步场景主要使用 asyncio.Queue。二者都支持 FIFO、容量限制和任务完成跟踪;但它们分别服务于线程调度和事件循环,不能因为接口名称相似就混用。queue 模块面向多生产者、多消费者线程,并通过锁提供同步语义;asyncio.Queue 专门用于 async/await 代码,且不是线程安全对象。(docs.python.org)


一、队列解决的不是“存储”,而是并发交接

设生产速率为:

λp=单位时间生产的任务数时间\lambda_p = \frac{\text{单位时间生产的任务数}}{\text{时间}}

设消费者处理速率为:

μc=单位时间完成的任务数时间\mu_c = \frac{\text{单位时间完成的任务数}}{\text{时间}}

如果一段时间内满足:

λp>μc\lambda_p > \mu_c

那么未完成任务数会持续增加。无界队列只能把问题从“生产者现在被阻塞”推迟为“进程最终耗尽内存”;它没有消除过载,只是失去了过载边界。

有容量限制的队列则把系统状态限制在:

0Q(t)K0 \leq Q(t) \leq K

其中:

  • Q(t)Q(t) 是时刻 tt 队列中的任务数;
  • KK 是队列容量;
  • Q(t)=KQ(t)=K 且生产者继续提交时,生产者必须等待、失败或采取丢弃策略。

这就是背压:下游处理能力不足时,通过队列容量把压力传播回上游,使上游不能无限制地产生待处理任务。

背压并不等于“永远阻塞”。它是一种明确的容量策略:

  • 阻塞式背压:队列满时等待空位;
  • 超时式背压:等待一段时间,超时后失败;
  • 拒绝式背压:立即返回“队列已满”;
  • 丢弃式背压:丢弃最旧、最新或低优先级任务;
  • 降级式背压:不再生成昂贵任务,改用缓存结果或简化处理。

队列的 maxsize 只负责提供一个容量边界,具体的过载策略由调用方式和业务代码决定。


二、queue.Queue:线程之间的同步队列

2.1 基本模型

queue.Queue(maxsize=0) 是 FIFO 队列。maxsize <= 0 表示不限制容量;正数表示队列最多保存多少个尚未取出的项目。队列满时,阻塞式 put() 会等待,直到消费者取走项目;非阻塞式 put_nowait() 则立即抛出 queue.Full。(docs.python.org)

import queue

q = queue.Queue(maxsize=2)

q.put("a")
q.put("b")

print(q.qsize())  # 常见输出:2

try:
    q.put_nowait("c")
except queue.Full:
    print("queue is full")

这里的 maxsize=2 约束的是队列中尚未被消费者取走的项目数,不是所有已提交但仍在处理的任务数。

例如:

生产者 put A
生产者 put B
消费者 get A
消费者开始处理 A
队列中只剩 B

此时队列大小是 1,但系统仍有两个尚未完成的任务:正在处理的 A 和队列中的 B。因此,qsize() 不能直接等价于“系统总负载”。

queue.Queue 内部使用锁和条件变量协调竞争线程,但文档明确指出,qsize()empty()full() 都不适合用来做跨线程的先检查后操作。比如:

if not q.empty():
    item = q.get()

检查完成后,另一个消费者可能已经取走了项目,后续 get() 仍然可能阻塞。官方文档将 qsize() 定义为近似大小,并明确说明它不能保证紧随其后的 get()put() 一定不会阻塞。(docs.python.org)

正确的思路是直接执行操作,并通过阻塞、超时或异常处理竞争结果:

try:
    item = q.get(timeout=1.0)
except queue.Empty:
    print("one second内没有任务")

2.2 put() 的三种行为

q.put(item)                    # 队列满时一直等待
q.put(item, timeout=2.0)       # 最多等待两秒
q.put_nowait(item)             # 不等待,满时抛出 queue.Full

put(item, block=False, timeout=...) 中的 timeout 会被忽略,因为调用方已经明确要求不阻塞。相同规则也适用于 get()。(docs.python.org)

选择哪种方式,取决于任务是否允许延迟:

  • 请求入口通常不能无限等待,应设置超时并返回过载响应;
  • 日志、遥测等可丢失数据可以使用非阻塞提交;
  • 必须处理的任务可以阻塞,但必须让上游也具备取消或超时能力;
  • 后台批处理可以等待,但要考虑程序退出时如何收尾。

2.3 task_done()join():容量和完成度是两套状态

队列中有两个容易混淆的概念:

  • 队列大小:还有多少项目没有被 get() 取出;
  • 未完成任务数:还有多少项目没有完成处理。

每次 put() 会增加未完成任务数;消费者调用 task_done() 后才会减少该计数;join() 会等待未完成任务数归零。(docs.python.org)

因此,下面的代码是正确的基本结构:

def worker(q: queue.Queue[str]) -> None:
    while True:
        try:
            item = q.get()
        except queue.ShutDown:
            return

        try:
            process(item)
        finally:
            q.task_done()

task_done() 必须放在 finally 中。否则,process(item) 抛出异常时,任务计数不会归零,主线程的 q.join() 可能永久等待。

另一方面,不能在没有对应 get() 的情况下多调用 task_done()

q.task_done()  # 如果没有尚未完成的 put,对应调用会抛出 ValueError

完整线程示例:

from __future__ import annotations

import queue
import threading
import time


def process(item: int) -> None:
    time.sleep(0.05)
    print(f"processed {item}")


def worker(name: str, q: queue.Queue[int]) -> None:
    while True:
        try:
            item = q.get()
        except queue.ShutDown:
            print(f"{name} stopped")
            return

        try:
            process(item)
        finally:
            q.task_done()


def main() -> None:
    q: queue.Queue[int] = queue.Queue(maxsize=3)

    threads = [
        threading.Thread(target=worker, args=(f"worker-{i}", q))
        for i in range(2)
    ]

    for thread in threads:
        thread.start()

    for item in range(10):
        q.put(item)
        print(f"submitted {item}")

    # 等待所有已提交任务处理完成。
    q.join()

    # 不再接受新任务,并让消费者在队列排空后退出。
    q.shutdown()

    for thread in threads:
        thread.join()

    print("all stopped")


if __name__ == "__main__":
    main()

运行时,生产者最多只能领先消费者三个排队项目。当两个工作线程都在处理任务且队列已满时,q.put(item) 会暂停;某个消费者执行 get() 后,生产者才获得新的空位。这条等待链就是背压的实际表现。


三、Python 3.13+ 的队列关闭协议

Python 3.14 中,queue.Queueasyncio.Queue 都提供了 shutdown(immediate=False),以及对应的 ShutDownQueueShutDown 异常。这些 API 在 Python 3.13 引入,因此不能用于要求 Python 3.12 或更早版本的程序。(docs.python.org)

3.1 正常关闭:停止增长,但处理剩余任务

调用:

q.shutdown()

之后:

  1. 队列不再接受新项目;
  2. 未来的 put() 会抛出关闭异常;
  3. 已经阻塞在 put() 上的生产者会被唤醒,并抛出关闭异常;
  4. 队列中已经存在的项目仍可通过 get() 取出;
  5. 队列排空后,后续 get() 会抛出关闭异常;
  6. 如果每个已取出的项目都调用了 task_done()join() 仍按正常完成协议解除阻塞。(docs.python.org)

可以把正常关闭理解为:

OPEN
  │ shutdown()
  ▼
CLOSING:禁止 put,允许 get 已有项目
  │ 队列排空
  ▼
CLOSED:get/put 都失败

生产者应把关闭异常视为“工作通道已结束”,而不是系统错误:

try:
    q.put(item)
except queue.ShutDown:
    # 生产者停止生成或把任务转移到其他持久化通道
    return

3.2 立即关闭:放弃队列中的工作

调用:

q.shutdown(immediate=True)

会立即清空队列,并减少相应的未完成任务计数;阻塞中的 get() 会被唤醒并抛出关闭异常。此时 join() 可能解除阻塞,即使被清空的任务实际上从未处理完成。因此,立即关闭会破坏通常的“不完成就不能 join 返回”不变量,必须谨慎使用。(docs.python.org)

这适合:

  • 进程即将被强制终止;
  • 队列中的任务已经过期;
  • 上游事务已回滚,旧任务全部失效;
  • 继续处理任务的代价高于丢弃任务。

这不适合:

  • 订单、账务、消息投递等必须确认结果的工作;
  • 仅仅因为某个消费者失败,就把所有未处理任务静默清空;
  • 依赖 join() 作为“所有任务都成功完成”的场景。

join() 只表示未完成计数归零。使用 shutdown(immediate=True) 后,它可能只表示“队列中的任务被清理了”,不再表示“任务被业务处理完成”。


四、asyncio.Queue:事件循环内的协作式队列

4.1 与 queue.Queue 的根本区别

asyncio.Queue 面向协程,不是线程安全对象;它的 put()get() 是协程操作,需要 await。它没有 timeout 参数,超时需要用 asyncio.wait_for()asyncio.timeout() 包装。(docs.python.org)

import asyncio

q = asyncio.Queue(maxsize=2)

await q.put("item")
item = await q.get()

不能把它当作线程间队列:

# 错误思路:其他线程直接操作 asyncio.Queue
q.put_nowait(item)

即使某次调用看似成功,也不能因此获得跨线程同步保证。需要在线程和事件循环之间传递数据时,应使用明确的线程到事件循环桥接机制,例如在线程中通过事件循环调度回调,而不是把异步队列当作普通线程安全容器。

4.2 异步背压

async def producer(q: asyncio.Queue[int]) -> None:
    for item in range(10):
        await q.put(item)
        print("submitted", item)

当队列容量已满时,await q.put(item) 会暂停当前协程,但不会阻塞整个事件循环。事件循环可以继续运行其他协程,包括消费者。

这与线程队列的区别在于:

queue.Queue.put()
    阻塞当前线程

asyncio.Queue.put()
    挂起当前任务,让事件循环运行其他任务

不过,“不会阻塞事件循环”不等于“没有代价”。如果生产者所在的请求协程一直等待队列空位,请求本身仍然会变慢。因此,异步背压依然需要超时、取消和上游限流。

使用 Python 3.11+ 的 asyncio.timeout()

async def submit_with_timeout(
    q: asyncio.Queue[int],
    item: int,
    timeout: float,
) -> bool:
    try:
        async with asyncio.timeout(timeout):
            await q.put(item)
    except TimeoutError:
        return False
    return True

asyncio.timeout() 在截止时间到达时取消当前任务,并在上下文管理器外部把取消转换为 TimeoutError;因此 TimeoutError 应在 async with 外捕获。(docs.python.org)

4.3 异步消费者的正确清理

async def worker(name: str, q: asyncio.Queue[int]) -> None:
    while True:
        try:
            item = await q.get()
        except asyncio.QueueShutDown:
            print(f"{name} stopped")
            return

        try:
            await process(item)
        finally:
            q.task_done()

这里有两个独立的退出路径:

  1. 队列正常关闭且排空,get() 抛出 QueueShutDown
  2. 工作任务被取消,正在执行的 await process(item) 收到 CancelledError

task_done() 必须覆盖第二种路径,否则消费者已经取出项目,却没有减少未完成计数,等待 join() 的协程会卡住。

取消是协作式的:任务被取消时,CancelledError 会在下一次合适的挂起点抛出。清理逻辑应放在 finally 中;如果显式捕获 CancelledError,清理完成后通常应继续抛出,而不是吞掉取消。TaskGroupasyncio.timeout() 等结构化并发组件内部也依赖取消机制,吞掉取消可能破坏它们的行为。(docs.python.org)

一个可运行的异步端到端示例:

from __future__ import annotations

import asyncio


async def process(item: int) -> None:
    await asyncio.sleep(0.05)
    print(f"processed {item}")


async def worker(name: str, q: asyncio.Queue[int]) -> None:
    while True:
        try:
            item = await q.get()
        except asyncio.QueueShutDown:
            print(f"{name} stopped")
            return

        try:
            await process(item)
        finally:
            q.task_done()


async def producer(q: asyncio.Queue[int]) -> None:
    for item in range(10):
        try:
            async with asyncio.timeout(1.0):
                await q.put(item)
        except TimeoutError:
            print(f"drop {item}: queue remained full")
        except asyncio.QueueShutDown:
            print("producer stopped: queue is closed")
            return


async def main() -> None:
    q: asyncio.Queue[int] = asyncio.Queue(maxsize=3)

    async with asyncio.TaskGroup() as tg:
        for i in range(2):
            tg.create_task(worker(f"worker-{i}", q))

        await producer(q)

        # 禁止继续 put,但允许消费者处理队列中已有的项目。
        q.shutdown()

        # 等待每个已提交且最终被 get 的项目完成。
        await q.join()

        # worker 会因为队列关闭且已排空而自行退出。


asyncio.run(main())

这个生命周期的关键顺序是:

启动消费者
    ↓
生产者提交任务
    ↓
生产者结束
    ↓
queue.shutdown()
    ↓
消费者继续 get 已有任务
    ↓
每个任务 task_done()
    ↓
queue.join() 返回
    ↓
队列排空,get() 抛 QueueShutDown
    ↓
消费者退出
    ↓
TaskGroup 等待所有消费者结束

TaskGroup 会保存任务、等待任务结束,并在某个任务以非取消异常失败时取消同组其他任务;相较于直接使用 gather(),它提供更强的任务生命周期约束。(docs.python.org)


五、task_done() 应该表示什么

task_done() 不是“我已经调用过 get()”的通知,而是:

这个项目对应的业务处理已经结束,无论结果是成功、失败还是明确记录为放弃。

因此下面的写法通常是错误的:

item = await q.get()
q.task_done()

await process(item)

它会让 join() 误以为任务已经完成,而实际处理还没有开始。

更合理的写法是:

item = await q.get()
try:
    await process(item)
finally:
    q.task_done()

如果处理失败但任务已经被记录到错误存储中,仍可以在 finally 中调用 task_done(),因为“队列交接生命周期”已经结束。是否重试,是业务协议;是否调用 task_done(),是队列协议。

重试也必须明确计数模型。例如:

async def worker(q: asyncio.Queue[tuple[int, int]]) -> None:
    while True:
        try:
            item, attempt = await q.get()
        except asyncio.QueueShutDown:
            return

        try:
            try:
                await process(item)
            except Exception:
                if attempt < 3:
                    try:
                        await q.put((item, attempt + 1))
                    except asyncio.QueueShutDown:
                        record_lost_retry(item)
                else:
                    record_failed(item)
        finally:
            q.task_done()

这里每次重新 put() 都会增加一次未完成任务计数,而当前取出的任务在 finally 中减少一次。只有在“重新提交成功”时,未完成任务总量才保持不变;如果重试提交失败,就必须记录任务丢失或转移到其他可靠通道。


六、容量应该如何推导,而不是随意填写

容量不是越大越好。设:

  • 生产者平均速率为 rpr_p
  • 消费者平均速率为 rcr_c
  • 允许生产者被背压的最长时间为 TT
  • 希望吸收的突发任务量为 BB

如果长期平均满足 rprcr_p \leq r_c,队列容量主要用来吸收突发。一个粗略下界是:

KBK \geq B

如果生产者在某段时间内以 rpr_p 生产,而消费者只能以 rcr_c 消费,持续 TT 秒,则这段时间增加的排队量近似为:

Kmax(0,rprc)TK \geq \max(0, r_p-r_c)T

例如:

  • 生产者每秒生成 100 个任务;
  • 消费者每秒处理 80 个任务;
  • 允许生产者最多被持续阻塞 5 秒。

则需要吸收的积压约为:

(10080)×5=100(100-80)\times 5=100

容量至少应接近 100,实际还要为处理耗时抖动、任务大小差异和调度误差预留余量。

但如果 rpr_p 长期大于 rcr_c,任何有限容量最终都会被填满。容量只能决定“多久后进入背压”,不能使系统获得超过消费者能力的吞吐量。

队列延迟还可用 Little 定律理解:

L=λWL=\lambda W

其中:

  • LL 是系统内平均任务数;
  • λ\lambda 是稳定吞吐率;
  • WW 是任务从进入到完成的平均时间。

在吞吐率不变时,队列中积压的任务越多,平均等待时间通常越长。增大容量可以减少生产者立即失败的概率,却可能让任务在队列里等待更久。因此,容量上限同时是:

  • 内存保护;
  • 延迟保护;
  • 故障扩散边界;
  • 上游节流触发点。

七、无界队列和 SimpleQueue 的边界

queue.Queue()asyncio.Queue() 默认 maxsize=0,表示无界容量。这个默认值适合简单示例,却不应自动成为生产配置。

无界队列的失败路径通常是:

下游变慢
  ↓
队列持续增长
  ↓
对象、缓存、日志和上下文持续占用内存
  ↓
垃圾回收压力增加
  ↓
延迟扩大
  ↓
进程被 OOM 或系统级联失败

queue.SimpleQueue 是无界 FIFO 队列,不提供任务跟踪功能,因此没有 task_done()join() 这套完成协议。它适合只需要线程安全交接、不需要等待“所有任务处理完”的场景。其 put() 不会因容量满而阻塞,但仍可能因底层内存分配失败。(docs.python.org)

不要把 SimpleQueue 当作带背压的队列。它能提供线程安全交接,但没有容量边界,也没有任务完成计数。


八、哨兵值与 shutdown() 的区别

在 Python 3.13 之前,常见做法是向队列放入哨兵值:

STOP = object()

for _ in workers:
    q.put(STOP)

消费者收到哨兵后退出:

item = q.get()
try:
    if item is STOP:
        return
    process(item)
finally:
    q.task_done()

哨兵方案仍然有价值,尤其是在:

  • 需要兼容 Python 3.12 及更早版本;
  • 不同队列实现没有统一关闭 API;
  • 关闭信号本身需要作为业务消息持久化;
  • 需要区分多种结束原因。

但它有几个协议成本:

  1. 多个消费者通常需要多个哨兵;
  2. 哨兵必须与普通数据可靠区分;
  3. 生产者可能在关闭过程中继续提交普通任务;
  4. 队列满时,投递哨兵本身也可能阻塞;
  5. 哨兵不是队列状态,无法自动唤醒所有已经阻塞在 put() 上的生产者。

shutdown() 是队列状态转换,不是一个特殊数据项。它可以禁止队列继续增长,并唤醒阻塞的生产者;正常关闭时还保留队列中已有项目。(docs.python.org)

因此,在 Python 3.14 的新代码中:

  • 队列本身负责生命周期关闭时,优先使用 shutdown()
  • 关闭信号属于业务数据时,可以继续使用哨兵;
  • 需要兼容旧版本时,使用哨兵或封装一个兼容层;
  • 不要同时使用两套关闭协议而不定义优先级。

九、get()、超时、取消和关闭异常的关系

消费者等待任务时,可能发生三种不同事件:

9.1 队列获得任务

item = await q.get()

返回项目后,消费者对该项目负有调用 task_done() 的责任。

9.2 队列被关闭

如果队列已关闭且没有可取项目,get() 会抛出 asyncio.QueueShutDown;线程队列对应 queue.ShutDown。这表示“不会再有任务”,通常是正常退出路径。(docs.python.org)

9.3 当前任务被取消或等待超时

try:
    async with asyncio.timeout(2):
        item = await q.get()
except TimeoutError:
    ...

这表示“这次等待超过了时限”,不等于“队列已经关闭”。超时后队列可能仍然开放,消费者可以重试、转移或退出。

如果外部直接取消消费者任务:

worker_task.cancel()

则消费者可能在 await q.get()await process(item) 处收到 CancelledError。它应该执行清理并传播取消,而不是把取消当成队列关闭来吞掉。取消和关闭是两个不同的控制面:

QueueShutDown  → 工作源结束,消费者正常退出
CancelledError → 当前协程被要求停止,应执行清理
TimeoutError   → 某次等待超过期限,可由业务决定重试或降级

把三者混为一谈,会导致典型错误:

  • 把超时当作永久关闭,消费者过早退出;
  • 把关闭异常当作系统故障,产生无意义报警;
  • 吞掉取消异常,导致上层 TaskGroup 或超时上下文无法正确结束。

十、关闭顺序中的竞态

考虑下面的异步流程:

await producer(q)
q.shutdown()
await q.join()

它成立的前提是:producer(q) 返回时,所有可能向该队列提交任务的生产者都已经结束。

如果还有其他生产者并发运行:

生产者 A 结束
主协程调用 shutdown()
生产者 B 尚未结束,准备 put()
生产者 B 收到 QueueShutDown

这可能是正确行为,也可能意味着任务被意外拒绝。关键不是异常本身,而是系统是否在设计上允许生产者 B 继续产生任务。

因此,关闭前需要先完成“生产者集合”的生命周期同步:

停止接受新输入
    ↓
等待所有生产者结束
    ↓
shutdown()
    ↓
消费者排空队列
    ↓
join()
    ↓
等待消费者退出

如果先调用 shutdown(),再等待生产者结束,那么正在生产的任务可能被拒绝。对于不可丢失任务,生产者必须在收到关闭异常后把任务写入持久化存储、重试队列或失败记录。

线程代码中也存在相同竞态。q.join() 只能等待当前已计入未完成计数的任务完成,不能自动阻止其他线程继续 put()。因此,“join() 返回”不等于“系统已经停止接受新任务”。


十一、queue.Queueasyncio.Queue 的选择

需求 选择
多个线程之间交接任务 queue.Queue
线程工作池中的 FIFO/LIFO/优先级调度 queue.QueueLifoQueuePriorityQueue
只需要简单线程安全 FIFO,不需要 join() queue.SimpleQueue
同一事件循环中的协程交接任务 asyncio.Queue
需要异步等待容量空位 asyncio.Queue(maxsize=...)
进程之间传递对象 multiprocessing.Queue,而不是前两者
网络服务中的跨进程、持久化、可确认消息 外部消息系统或数据库协议

Python 并发工具的选择取决于任务是 CPU 密集还是 I/O 密集,以及需要事件驱动的协作式并发还是线程、进程等抢占式并发模型。(docs.python.org)

尤其要注意,asyncio.Queue 不是“异步版的跨线程队列”。如果系统同时包含线程和事件循环,应在边界处明确转换:

线程池
  ↓ 线程安全队列
事件循环
  ↓ asyncio.Queue
异步消费者

每个边界都需要自己的同步协议,不能只因为两边都叫 Queue 就直接共享。


十二、如何诊断队列系统的失败

12.1 join() 永久等待

优先检查:

item = q.get()
try:
    ...
finally:
    q.task_done()

常见原因包括:

  • 某条异常路径没有执行 task_done()
  • 消费者取出项目后被取消,清理逻辑不完整;
  • 任务被重新提交,但旧任务和新任务的计数没有正确配对;
  • 调用了 shutdown(immediate=True),却错误地把 join() 当作业务成功确认;
  • 仍有生产者提交任务,导致未完成计数不断增加。

12.2 生产者似乎“卡死”

如果 maxsize 为正,检查:

  • 消费者是否已启动;
  • 消费者是否在处理某个慢任务;
  • 消费者是否因异常退出;
  • 消费者是否等待另一个永远不会发生的事件;
  • 是否在持有锁时调用了可能阻塞的 put()get()

线程队列通过锁协调竞争线程,但不设计为同一线程内可重入的同步结构。把队列操作放进不恰当的锁范围,可能形成锁顺序反转或自锁。(docs.python.org)

12.3 异步任务不退出

检查是否满足:

生产者结束
→ queue.shutdown()
→ 消费者得到 QueueShutDown
→ 消费者退出
→ TaskGroup 或 gather 完成

如果消费者写成:

while True:
    item = await q.get()
    ...

却没有关闭队列、哨兵或取消任务,它在队列排空后会永久等待,这是正常语义,不是 asyncio.Queue 的故障。

12.4 队列长度下降但延迟继续上升

这通常说明只观测了排队任务数,没有观测:

  • 正在处理的任务数;
  • 单任务处理耗时;
  • 任务等待时间;
  • 失败和重试次数;
  • 生产者等待容量的时间;
  • 队列关闭或拒绝次数。

队列大小是一个瞬时指标,不是完整的端到端延迟指标。对于容量和背压诊断,至少要记录任务进入队列的时间、开始处理的时间和完成时间。


十三、一个完整的状态模型

可以把一个带容量的队列抽象成以下状态:

stateDiagram-v2
    [*] --> Open
    Open --> Full: put until qsize == maxsize
    Full --> Open: get frees a slot
    Open --> Closing: shutdown()
    Full --> Closing: shutdown()
    Closing --> Closing: get existing item
    Closing --> Closed: queue becomes empty
    Closed --> Closed: get/put raises shutdown exception
    Open --> Aborted: shutdown(immediate=True)
    Full --> Aborted: shutdown(immediate=True)
    Aborted --> Closed: drain and wake waiters

关键点不在状态名称,而在允许的操作:

  • Open:可以 putget
  • Fullget 可以继续,阻塞式 put 等待;
  • Closing:禁止 put,允许提取已有任务;
  • Closed:不会再产生任务,getput 都以关闭异常结束;
  • Aborted:已有任务被丢弃,join() 的通常完成含义不再成立。

一个可靠的队列系统,必须让每条退出路径都能到达终态:

正常完成 → graceful shutdown → drain → join → consumers exit
取消退出 → cancel → finally cleanup → propagate cancellation
紧急终止 → immediate shutdown → explicitly account for discarded work

其中,shutdown() 解决的是队列是否还接受新任务join() 解决的是已计入的任务是否完成,任务取消解决的是执行者是否继续运行。三者不能互相替代。

队列真正提供的价值,是把生产、消费、容量、完成确认和关闭动作放进同一个可验证协议中。只要明确区分“排队”“取出”“处理完成”和“系统关闭”,queueasyncio.Queue 就不再只是 API 记忆题,而会成为可以推导、测试和诊断的并发基础设施。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。