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

Python asyncio 取消与同步:Timeout、Lock、Queue 和清理

asyncio 的并发模型建立在协作式调度之上:事件循环在一个线程中运行任务;某个任务执行到 await 时主动挂起,事件循环才有机会运行其他任务。因此,取消、超时、锁和队列并不是四组互不相关的 API,而是围绕同一个问题协作:

当任务停止等待、停止执行或被迫退出时,谁负责释放资源、交接状态和唤醒其他任务?

Python 3.14 的 asyncio 提供了协程、Task、TaskGroup、超时上下文管理器、同步原语和异步队列等高级 API。asyncio 的同步原语和队列专门面向异步任务,并不用于操作系统线程之间的同步。(docs.python.org)

本文先建立取消模型,再分别分析 TimeoutLockQueue 和清理逻辑,最后把它们组合成一个可以运行的生产者—消费者程序。


一、先建立模型:取消不是“杀掉任务”

1. Task、协程和 await

async def 定义的是协程函数。调用协程函数时,得到的是协程对象;仅仅调用并不会开始执行:

async def work():
    print("running")

coro = work()

此时 work() 的函数体还没有运行。只有下面两种方式会驱动它:

await work()

或者:

task = asyncio.create_task(work())
await task

协程对象表示“可以执行的异步计算”,Task 表示“已经交给事件循环调度的协程”。Task 还提供了取消、查询状态和等待结果的能力。asyncio.create_task() 会把协程包装为 Task 并安排执行;创建后应保存 Task 的强引用,否则事件循环只保留弱引用,Task 可能在完成前被垃圾回收。(docs.python.org)

2. 取消请求的传播路径

调用:

task.cancel()

并不是立即终止 Task,而是向 Task 发出取消请求。Task 在下一次可以被中断的时机收到:

asyncio.CancelledError

最常见的时机是任务正在等待某个 Future、Task、锁、队列或 asyncio.sleep() 时。

可以把取消过程表示为:

sequenceDiagram
    participant C as 取消方
    participant T as 目标 Task
    participant A as 当前 await
    participant F as finally 清理

    C->>T: task.cancel()
    T->>A: 注入 CancelledError
    A-->>T: await 中断
    T->>F: 执行 finally
    F-->>T: 释放资源/恢复不变量
    T-->>C: CancelledError 或正常结束

关键点是:取消必须经过任务自身的执行路径。如果任务正在运行一段没有 await 的长时间同步代码,事件循环无法插入取消异常,其他任务也无法运行。

async def bad():
    # 这里没有 await,会持续占用事件循环线程
    for _ in range(10**10):
        pass

这段代码中,即使外部调用 bad_task.cancel(),取消也要等循环结束后才有机会生效。对于 CPU 密集型工作,应切换到线程、进程或解释器池;对于异步代码,则应避免在事件循环线程中执行长时间阻塞操作。事件循环在同一线程中一次只运行一个 Task;Task 只有在 await 时才会让出执行权。(docs.python.org)

3. CancelledError 为什么不能随便吞掉

asyncio.CancelledError 直接继承自 BaseException,而不是普通的 Exception。这使得很多只捕获 Exception 的代码不会意外吞掉取消请求。

正确的清理模式是:

async def worker():
    resource = await acquire_resource()
    try:
        await use_resource(resource)
    finally:
        await close_resource(resource)

如果确实需要捕获 CancelledError,清理完成后通常应继续抛出:

async def worker():
    try:
        await use_resource()
    except asyncio.CancelledError:
        await record_cancelled()
        raise

如果任务吞掉取消异常,TaskGroupasyncio.timeout() 这类依赖取消机制实现的组件可能出现错误行为。除非确实需要抑制取消,否则不要调用 uncancel();如果真正抑制了取消,则还需要通过 uncancel() 清除任务内部的取消状态。(docs.python.org)


二、Timeout:超时本质上是有范围的取消

1. asyncio.timeout() 的工作方式

Python 3.11 引入的 asyncio.timeout() 返回一个异步上下文管理器:

async def request():
    async with asyncio.timeout(2):
        return await call_remote_service()

如果上下文中的等待超过 2 秒,超时管理器会取消当前 Task。它在内部处理 CancelledError,离开上下文时将其转换为内置的 TimeoutError。因此,TimeoutError 应在上下文外捕获:

async def request():
    try:
        async with asyncio.timeout(2):
            return await call_remote_service()
    except TimeoutError:
        return None

下面的写法不能按预期捕获 TimeoutError

async def wrong():
    async with asyncio.timeout(2):
        try:
            await call_remote_service()
        except TimeoutError:
            print("这里通常捕获不到超时")

原因是:上下文管理器只有在离开内部代码、执行 __aexit__() 时,才把内部的 CancelledError 转换成 TimeoutError。(docs.python.org)

2. 超时保护的是当前作用域

考虑这段代码:

async def outer():
    try:
        async with asyncio.timeout(1):
            await inner()
    except TimeoutError:
        print("outer timeout")

这里的超时取消的是执行 outer() 的 Task。若 inner() 是当前 Task 中直接执行的协程,它会沿着调用栈收到取消。

但如果 inner() 自己创建了一个独立 Task:

async def outer():
    child = asyncio.create_task(inner())
    try:
        async with asyncio.timeout(1):
            await child
    except TimeoutError:
        print("outer timeout")

需要区分两个对象:

  • 当前 Task 被超时取消;
  • child 是否被取消,取决于等待关系和封装方式。

如果要让子任务在调用方超时时继续运行,可以使用 asyncio.shield()

async def outer():
    child = asyncio.create_task(inner())

    try:
        async with asyncio.timeout(1):
            return await asyncio.shield(child)
    except TimeoutError:
        print("调用方超时,但 child 仍可能继续运行")
        raise

shield() 只保护被等待的 Task,不保护当前调用方。调用方仍会收到取消或超时异常;被保护的 Task 则不会因为调用方取消而自动取消。使用 shield() 时仍必须保存 Task 引用。(docs.python.org)

这会产生一个生命周期问题:

请求超时
  ├── 请求 Task 结束
  └── child 继续运行
       ├── 必须有明确的所有者
       ├── 必须有异常处理
       └── 必须最终回收

如果没有后续管理,shield() 很容易把“超时”变成“后台泄漏”。

3. asyncio.timeout()wait_for() 的区别

两者都能实现超时,但控制边界不同。

asyncio.timeout()

async with asyncio.timeout(2):
    await operation_a()
    await operation_b()

它限制的是一个作用域,可以包含多个 await。超时通过取消当前 Task 实现。

asyncio.wait_for()

await asyncio.wait_for(operation(), timeout=2)

它针对单个 awaitable。超时时会取消被等待的 awaitable,并抛出 TimeoutErrorwait_for() 会等待被取消的 awaitable 实际完成取消过程,所以总耗时可能超过给定的 2 秒。(docs.python.org)

例如:

import asyncio


async def slow_cancel():
    try:
        await asyncio.sleep(10)
    finally:
        print("开始清理")
        await asyncio.sleep(1)
        print("清理完成")


async def main():
    started = asyncio.get_running_loop().time()

    try:
        await asyncio.wait_for(slow_cancel(), timeout=0.1)
    except TimeoutError:
        elapsed = asyncio.get_running_loop().time() - started
        print(f"收到 TimeoutError,实际耗时约 {elapsed:.1f} 秒")


asyncio.run(main())

预期输出类似:

开始清理
清理完成
收到 TimeoutError,实际耗时约 1.1 秒

timeout=0.1 表示“开始取消的期限”,不保证调用者在 0.1 秒内立即恢复。被取消任务的 finally 如果还要等待网络关闭、事务回滚或资源释放,调用方必须等这些清理动作完成。

4. 相对超时和绝对截止时间

asyncio.timeout(5) 使用相对秒数;当多个函数共享同一个请求截止时间时,绝对 deadline 更准确:

async def request():
    loop = asyncio.get_running_loop()
    deadline = loop.time() + 5

    async with asyncio.timeout_at(deadline):
        await authenticate()
        await load_profile()
        await write_audit_log()

如果每一层都重新设置 5 秒超时,三层调用可能累计等待 15 秒;如果所有层共享同一个 deadline,则剩余预算会自然传递。

asyncio.Timeout 使用事件循环的单调时钟。可以先创建没有 deadline 的上下文,再在拿到配置后调用 reschedule();还可以调用 expired() 查询是否已经超时。超时上下文可以安全嵌套。(docs.python.org)

async def dynamic_timeout():
    cm = None

    try:
        async with asyncio.timeout(None) as cm:
            timeout = await get_timeout_from_config()
            loop = asyncio.get_running_loop()
            cm.reschedule(loop.time() + timeout)

            await operation()
    except TimeoutError:
        print("operation 超时")
    finally:
        if cm is not None and cm.expired():
            print("确认是 deadline 触发的退出")

5. Timeout 的边界:它不能中断同步阻塞函数

async def ineffective():
    async with asyncio.timeout(1):
        time.sleep(10)

这段代码不会在 1 秒时及时返回,因为 time.sleep() 阻塞了事件循环线程。超时回调本身也没有机会运行。

应将阻塞工作移出事件循环:

async def effective():
    async with asyncio.timeout(1):
        await asyncio.to_thread(time.sleep, 10)

但这里也要注意:取消等待线程的 Task,不等于强制杀死已经运行的线程。线程中的阻塞函数可能继续运行,直到自身结束。超时只保证异步调用方停止等待,不能为任意同步代码提供强制终止语义。


三、Lock:保护临界区,而不是保护整个任务

1. asyncio.Lock 的状态

asyncio.Lock 是异步任务之间的互斥锁。它只有两个核心状态:

UNLOCKED ── acquire 成功 ──> LOCKED
LOCKED   ── release       ──> UNLOCKED

当多个协程等待同一把锁时,官方语义保证获取锁具有公平性:通常先开始等待的协程先获得锁。锁本身不是线程安全对象,不应拿来同步不同 OS 线程。(docs.python.org)

推荐使用:

lock = asyncio.Lock()

async with lock:
    shared_state.update()

其语义等价于:

await lock.acquire()
try:
    shared_state.update()
finally:
    lock.release()

finally 是关键。只要任务已经成功获得锁,无论临界区正常返回、抛出异常还是被取消,都必须释放锁。

2. 取消发生在获取锁之前

async def update(lock, state):
    async with lock:
        state["count"] += 1

如果任务在 await lock.acquire() 期间被取消,它尚未拥有锁,因此不应该调用 release()async with lock 的上下文管理器能够区分“获取成功”和“获取失败”,避免错误释放。

手写代码时要保持同样的结构:

async def update(lock, state):
    acquired = False

    try:
        await lock.acquire()
        acquired = True
        state["count"] += 1
    finally:
        if acquired:
            lock.release()

以下代码存在严重风险:

async def broken(lock, state):
    try:
        await lock.acquire()
        state["count"] += 1
    finally:
        lock.release()

如果任务在 acquire() 过程中被取消,finally 仍会执行,而此时当前任务可能并未获得锁,最终可能触发 RuntimeError,或掩盖原始取消路径。

3. 取消发生在临界区内部

async def update(lock, state):
    async with lock:
        state["count"] += 1
        await asyncio.sleep(1)
        state["last_updated"] = True

如果任务在 sleep() 处被取消,async with lock 会退出并释放锁,但共享状态可能只更新了一半:

count 已增加
last_updated 尚未设置
任务被取消
锁已释放

因此,锁只保证互斥,不保证业务操作的原子性。要么把状态更新设计为可恢复的阶段,要么在取消时回滚:

async def update(lock, state):
    async with lock:
        old_count = state["count"]

        try:
            state["count"] += 1
            await persist(state)
            state["last_updated"] = True
        except asyncio.CancelledError:
            state["count"] = old_count
            raise

这里的回滚本身也可能需要异步操作。如果回滚动作必须完成,应明确其取消策略;不能简单地假设 finally 中的每个 await 都一定能执行到底。

4. 给锁操作设置超时

asyncio.Lock.acquire() 没有 timeout 参数。同步原语的方法通常不直接接收超时,应使用 asyncio.timeout()asyncio.wait_for() 包装等待。(docs.python.org)

async def try_update(lock, state):
    try:
        async with asyncio.timeout(0.2):
            async with lock:
                state["count"] += 1
                await persist(state)
    except TimeoutError:
        print("获取锁或执行临界区超过 200ms")
        return False

    return True

这个例子的超时范围包括两部分:

  1. 等待其他任务释放锁;
  2. 获得锁后执行 persist()

如果只想限制获取锁的时间,而不限制临界区:

async def try_acquire_only(lock, state):
    try:
        await asyncio.wait_for(lock.acquire(), timeout=0.2)
    except TimeoutError:
        return False

    try:
        state["count"] += 1
        await persist(state)
    finally:
        lock.release()

    return True

两种写法表达的是不同的业务约束,不能混用。

5. 不要在锁内等待不受控的外部操作

async with lock:
    await call_remote_service()
    shared_state["value"] = 1

远程服务变慢时,锁会长时间占用,其他任务都在锁外排队。更合理的结构通常是:

result = await call_remote_service()

async with lock:
    shared_state["value"] = result

前提是远程调用不依赖锁保护的状态,且结果仍然可以安全写入。如果读取、计算、写入必须组成一致操作,则需要保留锁,但应为整个临界区设置明确的 deadline。


四、Queue:用容量把背压传播回生产者

1. asyncio.Queue 的基本不变量

asyncio.Queue 是 FIFO 队列:

生产者 ── put ──> [ item item item ] ── get ──> 消费者

maxsize <= 0 时,队列容量视为无限;当 maxsize > 0 时,队列达到容量后,await queue.put(item) 会阻塞,直到消费者取走项目。异步队列的 qsize() 可以直接返回当前大小,但队列本身不是线程安全的。(docs.python.org)

容量有限时,队列形成背压:

消费者变慢
    ↓
队列逐渐填满
    ↓
put() 挂起
    ↓
生产者变慢
    ↓
上游压力被限制

如果使用无限队列,生产速度持续超过消费速度时,内存可能不断增长:

queue = asyncio.Queue()  # 无限容量

maxsize 不是一个随便填写的性能参数,而是系统愿意暂存多少未完成工作的明确上限。

2. get()task_done()join() 的计数关系

队列内部维护一个未完成任务计数:

put(item)      → unfinished += 1
task_done()   → unfinished -= 1
join()        → 等待 unfinished == 0

因此,对每一次成功的 get(),消费者必须在工作完成后调用一次 task_done()

async def worker(queue):
    item = await queue.get()
    try:
        await process(item)
    finally:
        queue.task_done()

task_done() 表示的不只是“我拿到了项目”,而是“这个项目对应的工作已经完成”。如果少调用一次,queue.join() 会永久等待;如果多调用一次,会抛出 ValueError。(docs.python.org)

错误示例:

async def broken_worker(queue):
    item = await queue.get()
    await process(item)

    # 忘记 task_done()

即使消费者最终退出,join() 仍然认为该项目未完成。

3. 取消消费者时如何保证计数正确

取消可能发生在不同阶段:

阶段 A:正在等待 get()
    → 没有取到项目
    → 不调用 task_done()

阶段 B:已经 get() 成功
    → 必须保证最终调用 task_done()

阶段 C:处理已完成
    → 调用 task_done()

因此,task_done() 应放在成功 get() 后的 try/finally 中:

async def worker(queue):
    while True:
        item = await queue.get()
        try:
            await process(item)
        except asyncio.CancelledError:
            print(f"取消处理 item={item!r}")
            raise
        finally:
            queue.task_done()

即使 process(item) 因取消而中断,队列计数也会归还。不过这并不代表项目已经成功处理;它只表示该项目不再由当前消费者继续处理。生产系统必须进一步决定:重新入队、写入失败存储、记录为未完成,还是允许丢弃。

4. QueueShutDown:生产者和消费者的关闭协议

Python 3.13 引入了 asyncio.Queue.shutdown(),Python 3.14 继续提供这一 API。调用:

queue.shutdown()

后:

  • 队列不能再增长;
  • 新的 put() 会抛出 QueueShutDown
  • 已经阻塞在 put() 的生产者会被唤醒并收到 QueueShutDown
  • 默认的正常关闭模式仍允许消费者取出队列中已有的项目;
  • 队列清空后,后续 get() 会抛出 QueueShutDown。(docs.python.org)

这比传统的哨兵对象更适合表达“队列本身已经关闭”:

queue.shutdown()

try:
    item = await queue.get()
except asyncio.QueueShutDown:
    return

正常关闭的状态可以表示为:

stateDiagram-v2
    [*] --> Open
    Open --> ShuttingDown: shutdown()
    ShuttingDown --> Draining: 允许 get() 取出已有项目
    Draining --> Closed: 队列为空
    Closed --> [*]

    Open --> Closed: shutdown(immediate=True)

shutdown(immediate=True) 是强制关闭:

  • 队列中的项目会被直接清空;
  • 阻塞在 get() 的任务会被唤醒并收到 QueueShutDown
  • 未完成任务计数会根据被清空的项目调整;
  • join() 可能在项目尚未真正处理完成时解除阻塞。

因此,立即关闭会破坏通常的 join() 不变量:join() 返回不再意味着每个项目都完成了处理。只有在进程关闭、任务不可恢复或明确接受丢弃工作时,才应使用立即关闭。(docs.python.org)

5. 为队列操作设置超时

队列的 put()get() 也没有直接的超时参数,应使用 asyncio.timeout()wait_for()

async def put_with_timeout(queue, item):
    try:
        async with asyncio.timeout(0.5):
            await queue.put(item)
    except TimeoutError:
        print("队列持续满载,生产者放弃等待")
        return False

    return True

超时放弃入队意味着项目尚未进入队列,调用方必须决定如何处理:

put 超时
  ├── 丢弃
  ├── 返回上游重试
  ├── 写入磁盘/外部缓冲
  └── 触发系统降级

不能把 put() 超时简单记录成“消费者处理失败”,因为项目可能根本没有进入队列。


五、清理:资源释放、状态交接和取消传播

1. finally 解决资源释放,不解决业务补偿

下面的代码可以保证锁释放:

async def critical(lock):
    async with lock:
        await operation()

也可以保证队列计数归还:

async def consume(queue):
    item = await queue.get()
    try:
        await operation(item)
    finally:
        queue.task_done()

但它们不能自动保证:

  • 数据库事务提交或回滚;
  • 文件是否写完整;
  • 网络请求是否在对端取消;
  • 已取出的队列项目是否需要重新投递;
  • 外部系统是否已经观察到部分副作用。

清理至少包含三类动作:

类型 目标 示例
资源清理 释放本地资源 关闭连接、释放锁
并发清理 结束相关任务 取消子任务、等待其结束
业务清理 恢复业务状态 回滚、重试、记录未完成

finally 通常只能直接解决第一类,并为后两类提供执行入口。

2. 清理中的 await 也可能被取消

async def worker():
    try:
        await operation()
    finally:
        await close_connection()

如果 Task 收到取消,通常会进入 finally。但如果清理阶段还存在新的取消请求,close_connection() 本身也可能被取消。

对于必须完成的短清理动作,可以谨慎使用独立 Task 加 shield()

async def worker(connection):
    try:
        await operation(connection)
    finally:
        cleanup_task = asyncio.create_task(connection.close())
        try:
            await asyncio.shield(cleanup_task)
        except asyncio.CancelledError:
            # 当前任务仍然需要把清理任务收尾
            await cleanup_task
            raise

这段代码的语义是:

  1. 创建独立的关闭 Task;
  2. 当前任务等待它,但不让当前任务的取消自动取消关闭 Task;
  3. 当前任务仍保留取消状态;
  4. 等关闭任务完成后重新抛出取消异常。

shield() 不应被用来无限期掩盖取消。如果清理操作可能长时间阻塞,应给清理设置独立的有限期限,并记录“清理未完成”这一故障状态。

3. 使用 TaskGroup 管理相关任务

手工使用 create_task() 时,必须自己维护:

  • Task 引用;
  • 异常收集;
  • 取消剩余任务;
  • 等待所有任务结束;
  • 防止后台任务泄漏。

TaskGroup 提供了结构化的生命周期:

async def run_group():
    async with asyncio.TaskGroup() as tg:
        tg.create_task(job_a())
        tg.create_task(job_b())

退出 async with 时,TaskGroup 会等待其中任务完成。如果某个子任务抛出非取消异常,其他任务会被取消;所有任务结束后,异常会组合成 ExceptionGroupBaseExceptionGroup。这与 asyncio.gather() 不同:gather() 默认在一个 awaitable 抛出异常后立即传播该异常,但不会自动取消其他仍在运行的 awaitable。(docs.python.org)

因此,TaskGroup 的清理路径通常是:

子任务失败
  ↓
TaskGroup 取消兄弟任务
  ↓
兄弟任务执行 finally
  ↓
TaskGroup 等待全部结束
  ↓
抛出 ExceptionGroup

如果清理代码吞掉 CancelledError,这条路径就可能无法正常收敛。


六、完整算例:带超时、锁、背压和正常关闭的消费者池

下面的程序展示四个组件如何配合:

  • Queue(maxsize=2) 提供有限容量;
  • 生产者在队列满时背压;
  • 每个消费者使用 asyncio.timeout() 限制处理时间;
  • 共享统计数据由 asyncio.Lock 保护;
  • task_done() 放在 finally 中;
  • queue.shutdown() 负责通知消费者退出;
  • TaskGroup 负责统一等待和传播异常。
import asyncio
from dataclasses import dataclass


@dataclass
class Stats:
    completed: int = 0
    timed_out: int = 0


async def process(item: int) -> None:
    # item=3 故意处理较慢,用于演示超时。
    delay = 0.5 if item == 3 else 0.05
    await asyncio.sleep(delay)


async def producer(queue: asyncio.Queue[int]) -> None:
    for item in range(1, 7):
        try:
            async with asyncio.timeout(0.4):
                await queue.put(item)
                print(f"producer: put {item}, qsize={queue.qsize()}")
        except TimeoutError:
            print(f"producer: put {item} timeout")
            # 这里表示项目没有进入队列。
            # 真实系统可以选择重试、落盘或丢弃。


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

        try:
            try:
                async with asyncio.timeout(0.2):
                    await process(item)
            except TimeoutError:
                async with stats_lock:
                    stats.timed_out += 1
                print(f"{name}: item {item} timed out")
            else:
                async with stats_lock:
                    stats.completed += 1
                print(f"{name}: item {item} completed")
        except asyncio.CancelledError:
            print(f"{name}: cancelled while handling item {item}")
            raise
        finally:
            # 无论成功、超时、异常还是取消,
            # 这个已经 get() 成功的项目都必须结束队列计数。
            queue.task_done()


async def main() -> None:
    queue: asyncio.Queue[int] = asyncio.Queue(maxsize=2)
    stats = Stats()
    stats_lock = asyncio.Lock()

    async with asyncio.TaskGroup() as tg:
        tg.create_task(worker("worker-1", queue, stats, stats_lock))
        tg.create_task(worker("worker-2", queue, stats, stats_lock))

        producer_task = tg.create_task(producer(queue))

        # 等待生产者完成,确保不会再有新的 put()。
        await producer_task

        # 正常关闭:允许消费者继续处理队列中的已有项目。
        queue.shutdown()

        # 等待所有已入队项目都调用 task_done()。
        await queue.join()

    print(
        f"completed={stats.completed}, "
        f"timed_out={stats.timed_out}"
    )


if __name__ == "__main__":
    asyncio.run(main())

运行前提

需要 Python 3.13 或更高版本,因为 Queue.shutdown()QueueShutDown 从 Python 3.13 开始提供。本文按 Python 3.14 语义编写。

关键执行过程

假设两个消费者都开始处理:

初始:
queue=[],unfinished=0

生产者 put 1:
queue=[1],unfinished=1

生产者 put 2:
queue=[1, 2],unfinished=2

消费者 get 1:
queue=[2],unfinished=2

消费者 get 2:
queue=[],unfinished=2

注意,get() 不会减少 unfinished。只有消费者调用 task_done() 才会减少计数:

item 1 完成:
unfinished=1

item 2 完成:
unfinished=0
queue.join() 返回

如果项目 3 处理超时:

process(3) 收到取消
    ↓
asyncio.timeout() 转换为 TimeoutError
    ↓
统计 timed_out += 1
    ↓
finally 调用 task_done()

这里“超时”不等于“队列计数仍然挂起”。从队列生命周期看,项目 3 已经完成了本次消费尝试;从业务生命周期看,它是否需要重试仍由应用决定。

为什么先等待生产者,再关闭队列

如果生产者还可能执行 put(),就提前调用:

queue.shutdown()

生产者会收到 QueueShutDown。因此正常关闭顺序应是:

停止接受新生产任务
    ↓
等待生产者结束
    ↓
queue.shutdown()
    ↓
消费者排空已有项目
    ↓
queue.join()
    ↓
消费者因 QueueShutDown 退出

如果生产者和消费者中任意一个发生不可恢复异常,TaskGroup 会取消其他任务。此时已有项目是否需要重新投递,不能由 TaskGroup 自动决定,必须由业务层设计。


七、失败路径与诊断方法

1. 任务被取消但程序迟迟不退出

常见原因包括:

async def cleanup():
    while True:
        await retry_forever()

如果清理逻辑无限重试,取消后任务可能始终无法结束。清理必须有自己的上限:

async def cleanup():
    try:
        async with asyncio.timeout(2):
            await close_resource()
    except TimeoutError:
        print("cleanup deadline exceeded")

注意,清理超时本身应记录为故障,而不是静默忽略。

2. queue.join() 永久阻塞

按以下顺序诊断:

  1. 是否每次成功 get() 都对应一次 task_done()
  2. 是否某个消费者在 get() 后、进入 try/finally 前被取消?
  3. 是否调用了 task_done() 多次?
  4. 是否使用了 shutdown(immediate=True),导致对 join() 的语义理解错误?
  5. 是否有生产者还在持续 put(),使得队列不断增加?

最安全的消费者结构仍然是:

item = await queue.get()
try:
    await handle(item)
finally:
    queue.task_done()

3. Lock 永久不释放

通常意味着任务在获得锁后退出,却没有走到释放路径:

await lock.acquire()
await operation()
lock.release()

如果 operation() 抛出异常或被取消,release() 不会执行。应改为:

await lock.acquire()
try:
    await operation()
finally:
    lock.release()

或者直接:

async with lock:
    await operation()

如果锁等待时间过长,可以给 acquire() 设置 timeout,并记录锁持有者、等待时间和临界区名称。不要通过强行调用 release() 来“修复”未知的锁状态,因为这可能破坏互斥关系。

4. 出现 Task exception was never retrieved

这通常表示通过 create_task() 创建了后台任务,却没有:

  • 保存引用;
  • await 它;
  • 或读取它的异常。

例如:

asyncio.create_task(background_job())

如果 background_job() 失败,异常可能直到 Task 被回收时才被事件循环记录。对于有明确生命周期的并发工作,优先使用 TaskGroup;真正的后台任务则必须有集合、完成回调、日志和关闭流程。官方文档也明确指出,TaskGroup 会保持任务引用、等待任务并传播异常。(docs.python.org)

5. 调试模式

开发阶段可以启用 asyncio debug mode:

PYTHONASYNCIODEBUG=1 python app.py

或:

asyncio.run(main(), debug=True)

调试模式可以帮助发现错误线程调用非线程安全 API、过慢回调以及事件循环选择器耗时过长等问题。还可以提高 asyncio 日志级别并启用 ResourceWarning。(docs.python.org)


八、几个容易混淆的边界

asyncio.Lock 不是线程锁

下面两种并发不是同一层次:

asyncio.Task 之间       → asyncio.Lock
OS thread 之间          → threading.Lock
进程之间                → multiprocessing 或外部协调机制

不要在一个线程中创建 asyncio.Lock,然后期待它能保护另一个线程访问的数据。

asyncio.Queue 不是跨线程队列

asyncio.Queue 适用于同一个事件循环中的异步任务。跨线程传递数据应使用线程安全队列,或通过 loop.call_soon_threadsafe()run_coroutine_threadsafe() 等桥接 API 安排事件循环工作。事件循环对象和大多数 asyncio 对象都不是线程安全的。(docs.python.org)

Timeout 不等于强制终止

Timeout 的实际语义是:

到达 deadline
    ↓
向相关 Task 发出取消
    ↓
等待取消路径完成
    ↓
调用方收到 TimeoutError

它不是操作系统级的强制杀死。被调用代码必须在 await 处让出控制权,并正确处理清理。

取消不等于失败

取消通常表示:

  • 调用方不再需要结果;
  • 上层任务失败后正在级联退出;
  • 服务正在关闭;
  • 超时上下文正在结束。

它与业务异常不同。一般不要把 CancelledError 转换成普通业务错误,也不要把取消记录成普通失败后继续重试。取消的正确传播通常比“尽量完成当前工作”更重要。


九、将四个组件统一起来

可以用三个不变量理解整个主题。

不变量一:拥有资源,就必须负责释放

成功 acquire Lock
    → 最终 release Lock

成功打开连接
    → 最终关闭连接

成功 get Queue item
    → 最终 task_done()

不变量二:取消必须能到达可等待点

任务有 await
    → 取消有机会注入

任务长期同步阻塞
    → 取消和其他任务都被延迟

不变量三:关闭必须有协议

停止生产
    → 关闭队列
    → 排空或立即丢弃
    → 等待消费者结束
    → 释放外部资源

如果系统能明确这三个不变量,TimeoutLockQueue 和清理就不再是零散的 API 记忆,而是一套可推导的生命周期设计:

  • Timeout 决定取消边界;
  • Lock 决定共享状态的互斥边界;
  • Queue 决定工作交接和背压边界;
  • finally、TaskGroup 和关闭协议决定故障路径能否收敛。

系列导航与关联阅读

官方资料

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