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)
本文先建立取消模型,再分别分析 Timeout、Lock、Queue 和清理逻辑,最后把它们组合成一个可以运行的生产者—消费者程序。
一、先建立模型:取消不是“杀掉任务”
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
如果任务吞掉取消异常,TaskGroup 和 asyncio.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,并抛出 TimeoutError。wait_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
这个例子的超时范围包括两部分:
- 等待其他任务释放锁;
- 获得锁后执行
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
这段代码的语义是:
- 创建独立的关闭 Task;
- 当前任务等待它,但不让当前任务的取消自动取消关闭 Task;
- 当前任务仍保留取消状态;
- 等关闭任务完成后重新抛出取消异常。
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 会等待其中任务完成。如果某个子任务抛出非取消异常,其他任务会被取消;所有任务结束后,异常会组合成 ExceptionGroup 或 BaseExceptionGroup。这与 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() 永久阻塞
按以下顺序诊断:
- 是否每次成功
get()都对应一次task_done()? - 是否某个消费者在
get()后、进入try/finally前被取消? - 是否调用了
task_done()多次? - 是否使用了
shutdown(immediate=True),导致对join()的语义理解错误? - 是否有生产者还在持续
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
→ 取消有机会注入
任务长期同步阻塞
→ 取消和其他任务都被延迟
不变量三:关闭必须有协议
停止生产
→ 关闭队列
→ 排空或立即丢弃
→ 等待消费者结束
→ 释放外部资源
如果系统能明确这三个不变量,Timeout、Lock、Queue 和清理就不再是零散的 API 记忆,而是一套可推导的生命周期设计:
Timeout决定取消边界;Lock决定共享状态的互斥边界;Queue决定工作交接和背压边界;finally、TaskGroup 和关闭协议决定故障路径能否收敛。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python asyncio 网络流:TCP、Reader、Writer、TLS 与背压
- 下一篇:Python 线程:生命周期、锁、条件变量、竞态和死锁诊断
- 延伸:Python 异步任务:create_task、TaskGroup、异常组和结构化并发
- 延伸:Python 队列与背压:queue、asyncio.Queue、容量和关闭协议
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论