Python 并发与 Asyncio:任务归属、等待和取消
并发程序可以在等待一个操作时推进另一个操作。例如,请求正在等待响应,程序可以先处理另一份已经收到的数据。并发指多项工作在重叠的时间段内推进,可以交错执行,也可以并行执行;并行指工作在同一时刻执行,例如不同 CPU 核心同时计算。重叠等待可以提高吞吐量,但不会自动加快一段计算。
使用并发之前,先明确每份工作由谁启动、谁接收结果、谁在失败时负责停止它。后面的示例把一批工作放在同一个生命周期边界内管理。
选择协程、线程还是进程
Python 的协程与 task 文档区分了协程对象和被调度的 task:调用 async def 函数只创建协程对象,直接 await 会等待它,创建 task 才能安排它与其他工作并发推进。连续等待两个协程仍然可以是顺序执行。
线程与进程执行器提供统一的提交与取结果接口。选择时先看工作主要在等什么,以及数据要怎样交接:
线程文档对 CPython 的条件很重要:同一解释器启用 GIL 时,多个线程不能同时执行 Python 字节码;释放 GIL 的扩展可能并行计算。Python 3.13 起还有可禁用 GIL 的 free-threaded 构建,但只有运行时确实禁用了 GIL,同一解释器中的多个线程才能同时执行 Python 字节码;运行选项或不兼容的 C 扩展可能重新启用 GIL。实际选择仍要看运行时、库和工作负载。
Python 3.14 还提供 InterpreterPoolExecutor:每个 worker 线程都有独立的解释器和 GIL,因此可利用多个核心。解释器之间不能直接共享可变对象,提交的函数、参数和返回值也要通过 pickle 序列化传递。
阻塞调用与协作调度
一个事件循环同一时刻运行一个 task。等待尚未完成的操作时,task 暂停,循环才有机会推进其他工作。await 本身不是保证让出执行权的标点:如果被等待的操作立即完成,就可能继续执行。
把同步文件读取、time.sleep 或长时间计算直接放进协程,会占住事件循环。给函数加上 async def 不会改变其内部调用。asyncio.sleep 会暂停当前 task;asyncio.to_thread 可把阻塞函数移到线程。CPU 密集循环可交给进程,也可拆成较小的计算批次,并在批次之间用 await asyncio.sleep(0) 主动让出执行权;仅拆分计算或增加 task 数量不会解决阻塞。
调用外部程序时,参数、管道和子进程清理仍要遵守子进程边界。异步调度只解决等待如何重叠,不替你处理子进程的完整生命周期。
Task 的归属、生命周期与取消
把一组相关工作放进 TaskGroup,就是给它们一个明确的生命周期边界:退出上下文前等待全部任务;某个任务抛出非取消异常时,取消其他任务,等它们结束后再抛出异常组;KeyboardInterrupt 和 SystemExit 则在收拢任务后重新抛出原异常。子任务要通过组的 create_task 注册,普通 asyncio.create_task 创建的额外后台任务不会自动归组。TaskGroup 和 asyncio.timeout 均从 Python 3.11 开始提供。
普通 create_task 需要调用方保存引用、接收异常并安排退出时的等待。默认 gather 会传播第一个异常,但不会因此取消其他工作;因此不能把“已经收到异常”当成“整组已经停下”。
取消请求会在下一个可处理的时机引发 CancelledError。用 finally 释放资源;如果显式捕获了取消异常,清理后通常应继续抛出。吞掉它可能破坏任务组或超时的控制流程。cancel() 发出请求后,仍要等待任务完成;已完成的写入也不会因此撤销。多 Agent 协作把类似的归属问题扩展到了文件、消息和多个执行环境。
线程的停止边界又不同:执行器的 Future.cancel() 无法取消已经运行的调用,result(timeout=...) 只限制取结果时的等待。因此,取消对 to_thread 的等待也不能保证底层线程停止。阻塞库需要自己的超时或协作停止机制。
共享状态也会发生竞争
即使只有一个事件循环,读取与写回之间有一次暂停,其他 task 也可能插入操作。asyncio.Lock可以让同一循环中的 task 独占一段操作;它不能用来同步操作系统线程,线程应使用 threading 的同步工具。
把下面的完整程序保存为 race.py,用 python3 race.py 运行。它故意在读取计数后暂停:
import asyncio
async def count(use_lock):
value = 0
lock = asyncio.Lock()
async def update():
nonlocal value
before = value
await asyncio.sleep(0)
value = before + 1
async def worker():
if use_lock:
async with lock:
await update()
else:
await update()
async with asyncio.TaskGroup() as group:
group.create_task(worker())
group.create_task(worker())
return value
async def main():
print(f"unlocked={await count(False)}")
print(f"locked={await count(True)}")
asyncio.run(main())
在 Python 3.13.15、默认任务工厂下得到:
unlocked=1
locked=2
无锁时两个 worker 都先读到 0,随后各自写回 1,丢失了一次更新。有锁时第二个 worker 在第一个写完后读到 1,最终写回 2。锁要覆盖整个“读取—修改—写回”,只锁最后一次赋值不能保护前面的读取。
这里把暂停放在锁内,是为了展示竞争。实际工作中尽量缩短临界区;可以先计算互不共享的结果,再交给一个任务汇总。Event 用来通知条件已发生,Semaphore 用来限制同时进入某段工作的数量。给每个输入都创建 task,再让它们等信号量,只限制正在执行的工作,不限制已经分配的 task 数量。
有界队列、背压与超时
FIFO 顺序和底层实现见队列。并发流水线还需要决定生产者过快时怎么办。asyncio.Queue的正数 maxsize 限制等待中的条目,队列满时 await put 会等待空位。这种让下游容量减慢上游生产的机制叫背压;maxsize=0 则不限制队列长度。它同样不能跨线程直接共享。
队列容量、worker 数和结果存储是三个不同的上限。下面的批次最多有 2 项在队列里、2 项在 worker 手中,即 4 项已接收的工作;生产者可能还持有 1 项等待入队,共 5 项处在交接或处理阶段。停止标记不算工作项。输入由 range 逐项产生,但结果字典随输入数增长;无限流应把结果交给也有容量控制的下游,条目大小也需要限制。
asyncio.timeout在期限到期时取消当前任务,并把这次超时引发的取消转成上下文外的 TimeoutError;外部取消仍是 CancelledError。把单项期限放在 worker 中,只计算取出条目后的处理时间;把整批期限放在任务组外,还覆盖入队和排空等待。事件循环阻塞或清理耗时都会让实际返回时间越过期限,超时不等于强制终止。
一个有界批次:正常完成、超时与失败
下面用异步睡眠模拟等待,随后计算 1 到 6 的平方;不访问网络、不启动外部进程。保存为 bounded.py,用 python3 bounded.py 运行(Python 3.11 或更新版本)。
import asyncio
WORKERS = 2
CAPACITY = 2
STOP = object()
async def operation(number, mode):
delay = 0.01
if mode in {"failure", "batch-timeout"}:
delay = 0.2
if mode == "failure" and number == 1:
delay = 0.01
if mode == "item-timeout" and number == 3:
delay = 0.2
await asyncio.sleep(delay)
if mode == "failure" and number == 1:
raise ValueError("invalid item 1")
return number * number
async def run_batch(mode):
queue = asyncio.Queue(maxsize=CAPACITY)
results = {}
tasks = []
closed = 0
status = "ok"
item_limit = 0.5 if mode == "batch-timeout" else 0.05
batch_limit = 0.02 if mode == "batch-timeout" else 2.0
async def produce():
for number in range(1, 7):
await queue.put(number)
await queue.join()
for _ in range(WORKERS):
await queue.put(STOP)
async def consume():
nonlocal closed
try:
while True:
number = await queue.get()
try:
if number is STOP:
return
try:
async with asyncio.timeout(item_limit) as item_deadline:
value = await operation(number, mode)
except TimeoutError:
if not item_deadline.expired():
raise
results[number] = "timeout"
else:
results[number] = value
finally:
queue.task_done()
finally:
closed += 1
try:
async with asyncio.timeout(batch_limit):
try:
async with asyncio.TaskGroup() as group:
tasks.append(group.create_task(produce()))
for _ in range(WORKERS):
tasks.append(group.create_task(consume()))
except* ValueError as errors:
status = "failed"
print(f"{mode}: {sorted(str(e) for e in errors.exceptions)}")
except TimeoutError:
status = "batch-timeout"
print(f"{mode}: status={status}, workers_closed={closed}, "
f"tasks_done={all(task.done() for task in tasks)}")
if status == "ok":
ordered = sorted(results.items())
total = sum(value for value in results.values()
if isinstance(value, int))
print(f"results={ordered}, total={total}")
async def main():
for mode in ("normal", "item-timeout", "failure", "batch-timeout"):
await run_batch(mode)
asyncio.run(main())
在 Python 3.13.15 上运行得到,结果按输入编号排序:
normal: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 9), (4, 16), (5, 25), (6, 36)], total=91
item-timeout: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 'timeout'), (4, 16), (5, 25), (6, 36)], total=82
failure: ['invalid item 1']
failure: status=failed, workers_closed=2, tasks_done=True
batch-timeout: status=batch-timeout, workers_closed=2, tasks_done=True
正常合计为 1 + 4 + 9 + 16 + 25 + 36 = 91;第 3 项超时后为 91 - 9 = 82。status=ok 表示按本例策略处理完批次,单项超时仍明确留在结果中,不能把它当作所有操作成功。
run_batch 打印状态与结果,返回值为 None。worker 捕获 TimeoutError 后,用超时上下文的 expired() 检查是否真是本项期限到期;若操作自身在期限到期前抛出同名异常,就重新抛出,由任务组取消其他工作并向调用方传播。
生产者先等 queue.join(),再为每个 worker 放入一个停止标记。每次成功 get 都对应一次 task_done,包括停止标记、超时和异常路径;get 尚未成功时不会多减计数。这个计数用于排空协调,不证明业务成功。
若 worker 失败,任务组会取消还在 join 或 put 中等待的生产者,以及其他 worker。失败路径不再等待队列排空:未取出的条目随本次本地队列一起丢弃。except* ValueError 处理 ValueError 及其子类,作为本例的整批失败;未匹配的异常仍通过异常组向调用方传播,外部取消也会继续传播。整批期限到期时,同一组任务也会先结束再打印超时状态。
每次输出的 workers_closed=2 和 tasks_done=True 表明两个 worker 都经过退出清理,三个子任务均已结束。结果汇总在组结束之后;各 worker 写入不同编号,写字典的这段操作中没有 await,所以本例不需要额外的锁。若改成共享计数并在读写之间等待,就需要前面展示的同步。
接入真实服务时,把 operation 换成可取消的异步调用,保留这组任务的归属与两层期限。这里没有重试:要增加重试,先确定总期限、尝试次数和副作用的幂等性;任务组的取消不会回滚远端已经完成的操作。