Chapter 55: Concurrency Abstraction
并发抽象要解决的主问题是:当一段 Python 程序同时面对等待、计算、共享状态、取消、异常和资源回收时,读者如何判断该用线程、进程、事件循环、执行器或同步原语,并能解释这些选择落在 CPython runtime 的哪个边界上。
本章的贯穿材料是一类常见服务:批量读取外部资源,解析结果,更新统计信息,再把失败、超时和关闭路径收束到一个可维护的调度边界。这个服务看起来只是“同时做多件事”,实际涉及五个对象层级:Python callable、运行中的 thread 或 process、event loop 中的 Task / Future、共享数据上的同步原语,以及最终负责生命周期的 executor 或 context manager。
Python 的并发模型需要先区分 concurrency 和 parallelism。Concurrency 表示程序能管理多个进行中的任务,让等待和执行交错推进;parallelism 表示多个执行单元在同一时刻真正占用多个 CPU core 执行。CPython 默认构建中,Global Interpreter Lock(GIL)让同一个 interpreter 内同一时刻只有一个线程执行 Python bytecode;这让 threading 很适合 I/O 等待场景,让 multiprocessing 更适合 CPU-bound 计算场景。Python 3.13 起存在 free-threaded builds 可以关闭 GIL,但官方文档仍把它作为特殊构建边界处理;本章默认讨论常见 CPython 构建,并在版本敏感处标出边界。
读完本章后,读者应能按固定顺序判断一个并发设计:先定位任务是在等待 I/O、执行 Python CPU 代码,还是调用会释放 GIL 的 native code;再看数据是否共享、是否需要跨进程传递、结果和异常如何返回;最后检查取消、backpressure、shutdown 和资源回收是否形成闭环。
下面的简化代码给出本章反复使用的材料。它暂时不选择具体并发方案,只把工程问题拆成 I/O、CPU、共享统计和关闭四个接口。
from dataclasses import dataclass
@dataclass
class ParseResult:
source: str
item_count: int
def blocking_read(source: str) -> bytes:
"""代表文件、网络或数据库读取。真实实现会阻塞调用线程。"""
raise NotImplementedError
def parse_payload(payload: bytes) -> ParseResult:
"""代表 CPU-bound 解析。真实实现会创建对象并执行 Python 代码。"""
raise NotImplementedError
def record_metric(result: ParseResult) -> None:
"""代表共享统计写入。并发调用时需要明确同步边界。"""
raise NotImplementedError
这段代码的关键点是职责边界。blocking_read() 的成本主要来自等待外部资源,线程或事件循环都可以把等待隐藏起来;parse_payload() 的成本主要来自 Python 代码执行,默认 CPython 下多个线程会受 GIL 约束;record_metric() 改写共享状态,需要同步原语或单线程收敛;关闭路径需要保证未完成任务、异常和底层资源都有明确归属。
55.1 threading
threading 把 OS thread 包装成 Python 对象,让一个进程内出现多个独立的控制流。Thread 对象负责把 callable、参数、线程名、daemon 状态、启动和 join 生命周期合并成一个 Python 层接口;底层执行仍依赖操作系统线程调度。官方文档把 threading 定位为基于低层 _thread 模块的高层接口,并说明它适合 I/O-bound 工作;CPython 默认构建中,GIL 让多个线程共享一个 Python bytecode 执行闸口,CPU-bound Python 代码通常无法通过线程获得多核并行收益。Python 3.14 threading 文档 对这些边界有明确说明。
在贯穿材料中,blocking_read() 是 threading 的典型目标。一个线程在系统调用、socket read、磁盘 I/O 或外部库等待时,当前线程处于等待状态,其他线程可以继续推进。多个线程共享同一进程内存,所以把读取结果放入内存队列、更新缓存或复用连接对象都很直接;同一特性也带来数据竞争,需要用 Lock、Condition、Event 或 Queue 形成同步边界。
下面的例子把多个 source 交给 worker thread。示例的目的在于定位 Thread、Queue、Lock 和线程生命周期之间的关系。
import threading
from queue import Queue
metrics = {"ok": 0, "failed": 0}
metrics_lock = threading.Lock()
thread_state = threading.local()
def worker(jobs: Queue[str], results: Queue[ParseResult]) -> None:
thread_state.read_count = 0
while True:
source = jobs.get()
try:
if source is None:
return
payload = blocking_read(source)
thread_state.read_count += 1
result = parse_payload(payload)
with metrics_lock:
metrics["ok"] += 1
results.put(result)
except Exception:
with metrics_lock:
metrics["failed"] += 1
raise
finally:
jobs.task_done()
这段代码里,Queue 负责跨线程传递任务,Lock 负责保护共享字典,threading.local() 负责给每个线程一份独立属性字典。thread_state.read_count 在每个 worker 中同名,但它绑定到当前线程的 thread-local storage;另一个线程读取到的是自己那份状态。这个模型适合保存连接、请求上下文、临时统计等线程私有数据。
Thread 生命周期要从创建、启动、运行、终止和 join 五个状态理解。创建 Thread(target=worker, ...) 只得到 Python 对象;调用 start() 才请求 OS 创建新的控制流,并让该控制流执行 run();目标函数返回或抛出未捕获异常后,线程结束;其他线程调用 join() 等待其结束。daemon thread 让进程退出时能直接结束剩余后台线程,适合非关键后台活动;涉及文件、锁、事务或网络连接的工作应使用非 daemon 线程和显式停止信号,让资源释放路径可分析。
GIL 的约束要落在执行对象上理解。GIL 限制的是同一个 CPython interpreter 内多个线程执行 Python bytecode 的并行度;它不等于禁止线程并发,也不等于所有 native code 都串行。I/O 等待、某些 C 扩展、压缩库、数值库可能在内部释放 GIL,让其他线程推进。判断 threading 是否合适,应先看主体成本是否来自外部等待或释放 GIL 的 native 工作,再看共享内存是否带来同步复杂度。
Python 3.14 的 Thread 构造器包含 context 参数,用于控制 contextvars.Context 在线程启动时如何传递;free-threaded build 下默认行为也有差异。这个细节说明线程抽象不仅包含 OS thread,还包含 Python runtime 状态的继承策略。对库代码而言,线程启动点需要把“可调用对象、参数、上下文、停止信号、异常处理和 join 责任”放在同一处管理,防止线程成为失去 owner 的后台控制流。
55.2 multiprocessing
multiprocessing 把并发边界从 thread 提升到 process。一个 child process 拥有自己的 Python interpreter、地址空间、模块导入状态、文件描述符集合和 GIL;父进程与子进程通过 pipe、queue、socket、shared memory 或 manager 交换数据。官方文档把它描述为使用类似 threading 的 API 创建进程,并通过 subprocess 绕开 GIL;同时它标出 Queue 中的对象会被序列化,进程目标和参数通常需要可 pickle。Python 3.14 multiprocessing 文档 对这些约束有集中说明。
在贯穿材料中,parse_payload() 更接近 multiprocessing 的目标。解析函数如果主要执行 Python bytecode,多个进程可以让多个 CPU core 同时工作,因为每个进程有自己的 GIL。代价也很明确:启动进程需要创建或复用 worker;参数和返回值需要跨进程序列化;大对象复制会增加内存和 IPC 成本;子进程崩溃需要被父进程观察;关闭时要处理 close()、terminate()、join() 和未取出的队列数据。
下面的例子把 CPU-bound 解析放到子进程中。它强调输入输出必须能越过进程边界。
from multiprocessing import Process, Queue
def parse_worker(payloads: Queue[tuple[str, bytes]], results: Queue[ParseResult]) -> None:
while True:
item = payloads.get()
if item is None:
return
source, payload = item
results.put(parse_payload(payload))
if __name__ == "__main__":
payloads: Queue[tuple[str, bytes]] = Queue()
results: Queue[ParseResult] = Queue()
process = Process(target=parse_worker, args=(payloads, results))
process.start()
payloads.put(("source-a", b"..."))
payloads.put(None)
process.join()
这里的 Queue 表面上和线程队列相似,runtime 边界完全不同。线程队列传递的是同一进程内对象引用;进程队列传递的是序列化后的对象内容。父进程把 (source, payload) 放入队列,子进程拿到的是反序列化后的新对象。对象身份、打开的连接、锁、lambda、闭包、本地定义函数和大型内存图都可能在这个边界出问题。
进程启动方式会改变可观察行为。spawn 会启动新的 Python 解释器并重新导入主模块,所以目标函数必须位于可导入模块顶层,并且需要使用 if __name__ == "__main__": 保护启动代码;fork 会复制父进程当前地址空间,在存在多线程或持有锁时会继承复杂状态;forkserver 通过专门服务进程创建子进程,减少从任意父进程 fork 的风险。不同平台默认启动方式不同,库代码应让调用方传入 multiprocessing context,并在文档中说明自己依赖的 start method。
共享内存只解决“大块 bytes 或数组复制成本”问题,不会自动解决对象语义问题。multiprocessing.shared_memory 可以让多个进程访问同一块内存区域,但 Python 对象图、引用计数、锁状态和生命周期仍需要额外协议管理。把共享内存用于批量数据时,通常要同时设计 owner、offset、length、释放时机和异常路径;否则进程隔离带来的清晰边界会被手工共享状态重新打散。
资源回收是 multiprocessing 的核心成本之一。父进程应清楚哪些 worker 由自己创建,哪些队列或 pipe 由自己关闭,异常时是否要终止子进程,子进程退出码如何映射到业务错误。对长期服务而言,进程池还要处理 worker 泄漏、内存增长、initializer 失败和 broken pool。进程并行解决 GIL 约束,同时把错误传播、对象传输和资源回收提升为设计问题。
55.3 asyncio
asyncio 把并发边界放到单线程 event loop。它使用 async def 定义 coroutine function,调用后得到 coroutine object;event loop 把 coroutine 包装成 Task,并在 await、I/O readiness、timer 或 Future 完成时切换执行。官方文档把 asyncio 定位为使用 async / await 编写并发代码的库,适合 I/O-bound 和高层网络代码,并提供 high-level API 与面向框架作者的 low-level event loop、transport 和 protocol API。Python 3.14 asyncio 文档 给出了这层结构。
在贯穿材料中,blocking_read() 需要先改成真正的 async I/O,才能直接放进 event loop。如果它仍是普通阻塞函数,放进 coroutine 内调用会占用 event loop 所在线程,让其他 Task 无法推进。asyncio 的并发收益来自“在等待点主动交还控制权”,所以每个耗时操作都必须表现为 awaitable,或者被移交到线程池、进程池等外部执行边界。
下面的示例用 async 形式表达读取和调度。它的目的在于区分 coroutine、Task、Future 和 event loop 的角色。
import asyncio
async def async_read(source: str) -> bytes:
await asyncio.sleep(0.01)
return source.encode()
async def parse_one(source: str) -> ParseResult:
payload = await async_read(source)
return parse_payload(payload)
async def parse_many(sources: list[str]) -> list[ParseResult]:
tasks = [asyncio.create_task(parse_one(source)) for source in sources]
return await asyncio.gather(*tasks)
parse_one() 调用后只创建 coroutine object,代码尚未执行。create_task() 把 coroutine 注册到当前 event loop,得到 Task;Task 负责驱动 coroutine 执行,保存结果或异常,并在 await 的对象完成后继续推进。gather() 聚合多个 awaitable,调用方等待聚合结果。整个过程通常发生在一个 OS thread 内,多个 Task 通过 await 协作切换。
Future 在 asyncio 中表示“稍后完成的结果占位符”。Task 是 Future 的子类,额外负责驱动 coroutine。低层 I/O 回调、timer、transport/protocol 或线程池回调完成时,会把结果写入 Future;等待它的 Task 被 event loop 放回 ready 队列。这个结构让回调式 I/O 能接入 async / await 语法。
transport / protocol 是 asyncio 的低层抽象。transport 负责实际 I/O 通道,例如 socket 写入、缓冲和关闭;protocol 负责在连接建立、数据到达、连接丢失时接收回调。高层 streams API 把这层包装为 StreamReader / StreamWriter,业务代码一般使用 streams;网络框架或协议库会直接面对 transport/protocol,因为它们需要控制缓冲、连接生命周期和协议状态机。
取消是 asyncio 设计里的普通控制流。调用 task.cancel() 会安排在 coroutine 的下一个可取消等待点注入 CancelledError;被取消的 coroutine 需要在 finally 或 async context manager 中释放资源。取消不会自动停止已经提交到线程池或进程池的正在运行 callable,也不会自动撤销外部系统中的请求。工程判断要把“Python Task 已取消”和“底层工作已停止”分开验证。
asyncio 的单线程性质让共享状态写入更容易分析,因为同一个 event loop 中的普通 Python 代码片段在两个 await 之间不会被其他 Task 抢占。边界在于任何 await 都可能让其他 Task 运行;一个对象在 await 前后可能被别的 Task 修改。设计 async 代码时,要把临界状态改写压缩在无 await 的短段落里,或者使用 asyncio.Lock、asyncio.Queue、Semaphore 等异步同步原语。
55.4 concurrent.futures
concurrent.futures 把“把 callable 提交给 worker,并在稍后取回结果”抽象成统一接口。Executor.submit() 接收 callable 和参数,返回 concurrent.futures.Future;调用方可以 result()、exception()、cancel()、done() 或注册 callback。具体 worker 可以是 thread、process,Python 3.14 文档中还包含 InterpreterPoolExecutor,但本章只把它作为版本边界提示,详细 runtime 隔离留到后续 subinterpreters 章节。Python 3.14 concurrent.futures 文档 说明这些 executor 实现共享同一抽象接口。
在贯穿材料中,concurrent.futures 适合把调度边界稳定下来。读取阶段可以使用 ThreadPoolExecutor,解析阶段可以使用 ProcessPoolExecutor,调用方只面对 Future。这种封装让“怎么执行”从“怎么提交、等待、取消、传播异常、关闭 worker”中分离出来。
下面的例子把同一个函数分别提交给不同 executor。示例的重点是接口统一,成本模型不同。
from concurrent.futures import Executor, ThreadPoolExecutor, ProcessPoolExecutor, as_completed
def run_parse_pipeline(sources: list[str]) -> list[ParseResult]:
results: list[ParseResult] = []
with ThreadPoolExecutor(max_workers=8) as readers:
read_futures = {
readers.submit(blocking_read, source): source
for source in sources
}
payloads: list[tuple[str, bytes]] = []
for future in as_completed(read_futures):
source = read_futures[future]
payloads.append((source, future.result()))
with ProcessPoolExecutor(max_workers=4) as parsers:
parse_futures = [
parsers.submit(parse_payload, payload)
for _, payload in payloads
]
for future in as_completed(parse_futures):
results.append(future.result())
return results
future.result() 是异常传播点。worker 内 callable 正常返回时,result() 返回该值;callable 抛出异常时,result() 在调用线程重新抛出对应异常。这个设计让异常从 worker 边界回到 owner 边界,调用方可以集中处理失败。对于 ProcessPoolExecutor,异常对象本身也要跨进程序列化;worker 进程异常退出时,pool 可能进入 broken 状态,后续提交失败。
cancel() 的语义依赖 Future 状态。尚未开始运行的任务可以取消;已经运行的线程 callable 通常继续执行;进程池中的任务也受队列、worker 状态和底层进程生命周期影响。Future 的取消表示“这个异步结果不再正常完成”,不保证底层副作用撤销。凡是 callable 会写文件、发请求、提交事务或改写外部状态,都要让 callable 自己支持幂等、超时或显式停止信号。
Executor.map() 提供批量提交和按输入顺序取结果的接口。Python 3.14 文档中 map() 包含 buffersize 参数,用于限制已提交但尚未 yield 的任务数量;这把 backpressure 引入 executor 层。没有容量控制时,调用方可能一次性把大量任务、参数对象和 Future 放入内存,并把下游 worker 队列推满。批量接口要同时考虑顺序、容量、超时和异常传播。
Executor 的核心价值是形成 shutdown 边界。with ThreadPoolExecutor(...) 或 with ProcessPoolExecutor(...) 会在离开作用域时触发关闭;默认等待已提交任务结束。显式 shutdown(cancel_futures=True) 可以取消尚未开始的 pending futures。库代码如果创建 executor,应说明 owner 是谁;长期服务如果复用 executor,应说明何时关闭、如何处理 broken 状态、是否允许调用方传入外部 executor。
55.5 synchronization primitives
同步原语解决的是“多个控制流如何围绕共享状态建立顺序关系”。它们把可见状态、等待条件和唤醒责任写进 runtime 对象,是并发程序的顺序协议。Python 标准库里,threading、multiprocessing 和 asyncio 都提供各自语境下的同步工具;名字相似,阻塞方式和适用边界不同。线程原语阻塞 OS thread,进程原语跨 process 协调,async 原语让 Task 挂起并交还 event loop。
Lock 表达互斥。一个控制流进入临界区前 acquire,离开时 release;同一时间只有一个控制流能持有它。它适合保护共享字典、计数器、缓存结构或不可重入外部资源。临界区应短,并且里面不应执行无界等待、网络请求或调用不可控用户回调;否则持锁时间会放大成系统吞吐瓶颈。
RLock 表达同一线程可重入的互斥。它记录 owner 和递归计数,同一线程多次 acquire 后需要对应次数 release。它适合对象方法之间互相调用且共享同一内部状态锁的场景。使用 RLock 前要确认重入来自真实调用链需求;如果只是为了让死锁现象消失,代码仍可能隐藏了锁顺序混乱。
Semaphore 表达容量。它维护一个计数器,acquire 消耗一个许可,release 归还一个许可。容量为 1 时接近互斥锁;容量大于 1 时适合限制并发连接数、并发请求数、打开文件数或外部 API 配额。BoundedSemaphore 会在 release 超过初始容量时报错,适合暴露许可归还错误。
Condition 表达“状态谓词变化后的通知”。它通常和一个 lock 配合使用,等待方在持锁状态下检查条件,条件不满足时 wait;通知方修改共享状态后 notify。正确使用 Condition 的重点是循环检查谓词,因为唤醒只表示“需要重新检查”,不表示条件必定成立。
import threading
from collections import deque
items: deque[bytes] = deque()
condition = threading.Condition()
closed = False
def producer(payload: bytes) -> None:
with condition:
items.append(payload)
condition.notify()
def consumer() -> bytes | None:
with condition:
while not items and not closed:
condition.wait()
if items:
return items.popleft()
return None
这段代码里,等待谓词是 items 非空或 closed 为真。notify() 只负责叫醒等待方;等待方醒来后仍要在锁内重新检查谓词。这个模式把“队列状态”和“关闭信号”放进同一个同步边界,读者可以直接分析正常消费、空队列等待和关闭退出三条路径。
Event 表达一次性或阶段性信号。一个控制流调用 set() 后,等待它的控制流继续执行;调用 clear() 可把它恢复为未触发。它适合停止信号、初始化完成信号、配置加载完成信号。Event 不携带队列数据,也不记录通知次数;如果需要传递多个工作项,应使用 Queue。
Queue 表达跨控制流通信。线程版 queue.Queue 内部包含锁和条件变量,支持 put()、get()、task_done() 和 join();它把共享列表、锁、通知和容量控制合并成一个对象。asyncio.Queue 使用 awaitable 形式挂起 Task;multiprocessing.Queue 通过序列化跨进程传递对象。选择 Queue 时,要先确认通信双方处在线程、进程还是 event loop 语境,再判断对象能否共享或需要序列化。
同步原语的检查顺序是:先写出共享状态是谁拥有,再写出哪些控制流会读写它;接着判断要表达互斥、容量、状态变化通知、单次信号还是消息传递;最后检查等待路径是否有超时、关闭信号和异常释放。没有这条顺序时,代码很容易把 Lock、Event 和 Queue 混用成难以停止的后台系统。
55.6 process isolation
进程隔离把内存、GIL、崩溃影响和权限边界拆开。每个进程有独立地址空间和独立 interpreter 状态;一个子进程的 Python heap、全局变量、模块缓存和 GIL 与父进程分离。这个边界让 CPU-bound Python 代码获得多核并行机会,也让崩溃影响更容易隔离:worker 进程崩溃时,父进程可以观察退出码、记录失败、重建 worker 或让 pool 进入 broken 状态。
隔离的直接后果是对象引用无法直接跨进程传递。父进程中的 list、dict、class instance、函数对象和文件对象都属于父进程地址空间;子进程拿到的要么是序列化后的值,要么是继承或复制来的 OS 句柄,要么是共享内存、manager proxy 等专门协议对象。把线程模型中的“共享对象加锁”照搬到进程模型,会在对象身份、延迟和一致性上产生错误判断。
序列化成本来自三个层次。第一层是 CPU 成本,pickle 需要遍历对象图并重建对象;第二层是内存成本,父子进程分别持有对象副本;第三层是语义成本,反序列化依赖可导入类路径、版本兼容和可信边界。对小任务来说,这些成本可能超过并行收益;对大批量 CPU 任务来说,把输入切块、减少返回数据、复用 worker 和使用共享内存能降低边界成本。
IPC(Inter-Process Communication)是进程隔离后的通信协议。Pipe 适合两个端点之间的消息;Queue 适合多 producer / consumer;socket 适合跨机器或服务化边界;shared memory 适合大块连续数据;manager proxy 适合把对象操作代理到 server process。每种 IPC 都要回答三个问题:数据所有权在哪里,写入和读取的顺序如何定义,异常或关闭时未消费数据如何处理。
权限边界也可以借助进程隔离表达。一个 worker 进程可以用更少权限、更受限的环境变量、更窄的文件描述符集合运行;父进程只通过受控消息协议交给它任务。Python 标准库不会自动建立安全沙箱,但进程边界提供了比线程更清晰的系统级隔离点。涉及不可信输入、可能崩溃的 native 扩展或高风险解析器时,进程隔离通常比线程隔离更可审计。
从 runtime 角度看,process isolation 的判断顺序是:先确认并行收益是否足以覆盖启动和序列化成本;再列出跨边界对象是否可 pickle、是否过大、是否依赖打开句柄;接着设计 IPC 的消息格式和关闭协议;最后定义 worker 崩溃、父进程取消和资源清理的行为。这个顺序能防止“为了绕开 GIL 使用进程”变成一组不可追踪的后台进程。
55.7 executor abstraction
executor abstraction 把任务提交、worker 生命周期、结果获取、异常传播和 shutdown 包装成一个调度边界。调用方不直接管理每个 thread 或 process,而是把 callable 提交给 executor;executor 负责把 callable 分配给 worker,返回 Future 表达异步结果。这个抽象让业务代码可以围绕 Future 编排,而把 worker 数量、启动方式和复用策略压缩到 executor 配置中。
在贯穿材料中,executor 是把多种并发模型组合起来的关键位置。读取阶段的 ThreadPoolExecutor、解析阶段的 ProcessPoolExecutor、event loop 的 run_in_executor() 都能围绕同一个概念建模:提交一个函数,获得一个未来结果,在 owner 边界等待它,最后关闭执行资源。代码结构因此从“到处创建线程或进程”转成“少数几个受控调度器”。
下面的图展示 executor 边界内外的角色。图中不描述具体线程池实现,只描述调用方能依赖的最小模型。
这条路径的关键点是 Future 同时承载结果和状态。调用方提交任务后,不再直接知道哪一个 worker 执行它;调用方通过 Future 观察 pending、running、finished、cancelled 和 exception。异常传播也在 Future 上完成,业务错误从 worker 回到调用方时保留原始语义,executor 自身错误则通过 BrokenExecutor、BrokenThreadPool 或 BrokenProcessPool 这类状态暴露。
worker 生命周期要由 owner 明确。函数内部临时创建 executor 并用 with 包住,owner 就是当前函数;服务启动时创建全局 executor,owner 就是服务生命周期管理器;库函数接收外部 executor,owner 就是调用方。owner 不清晰会产生两类问题:一类是 executor 过早关闭,导致后续 submit 失败;另一类是 executor 长期存活,导致线程、进程、队列和引用保留超过业务生命周期。
executor 的队列是 backpressure 的第一道边界。无限制 submit 会让内存中堆积大量 Future、参数对象和待执行任务。同步代码可以用 Semaphore 包住 submit 数量,或者使用 Python 3.14 Executor.map(buffersize=...) 这类容量控制;async 代码可以用 asyncio.Semaphore 或有界 Queue 限制进入 executor 的任务数。容量设计应靠近提交点,因为提交点最清楚输入规模和下游吞吐。
executor shutdown 是异常路径的一部分。正常路径下,with 离开时等待已提交任务;错误路径下,调用方要决定 pending 任务是否取消、running 任务是否等待、进程池是否强制终止、未取结果如何记录。shutdown() 之后继续 submit 是调用方错误;worker 初始化失败或进程异常退出导致 executor broken 时,应把它作为调度层故障处理;把所有失败都归因于业务 callable 会丢失 executor 状态信息。
55.8 async abstraction
async abstraction 把“正在进行的工作”表示成 awaitable,并把调度推进交给 event loop。它让调用方用顺序代码表达并发等待,但也要求每个等待点、取消点、缓冲边界和阻塞调用都可见。写 async 代码时,真正需要管理的是 coroutine object、Task、Future、event loop、I/O transport、同步原语和外部 executor 之间的组合关系。
取消需要作为接口契约设计。一个 async 函数如果持有连接、锁、临时文件或事务,它必须在 finally、async with 或 cleanup callback 中释放资源;否则 CancelledError 到达等待点后会中断正常路径。取消还会向下游 awaitable 传播,例如等待 gather()、TaskGroup 或单个 Task 时,外层取消会影响内层任务。工程上应说明取消后的状态:结果是否丢弃,部分写入是否提交,外部请求是否仍可能完成。
backpressure 在 async 系统中通常由 Queue(maxsize=...)、Semaphore、stream drain、连接池容量或 executor 提交窗口表达。没有 backpressure 时,producer 可以远快于 consumer 创建 Task、填满内存、堆积 socket buffer 或压垮下游服务。一个稳定的 async pipeline 应让每个 producer 在容量耗尽时 await,让等待本身成为调度信号。
import asyncio
async def bounded_parse(sources: list[str], limit: int) -> list[ParseResult]:
semaphore = asyncio.Semaphore(limit)
async def one(source: str) -> ParseResult:
async with semaphore:
payload = await async_read(source)
return parse_payload(payload)
return await asyncio.gather(*(one(source) for source in sources))
这段代码中,Semaphore 把并发读取数量限制在 limit。注意 parse_payload() 仍在 event loop 线程内执行;如果它是真正的 CPU-bound Python 代码,应通过进程池或专门 worker 转移出去。这个例子证明 async backpressure 解决的是任务进入速度和等待规模,CPU 计算边界仍要单独判断。
任务泄漏来自失去 owner 的 Task。调用 asyncio.create_task() 后,如果没有保存 Task、没有等待它、没有在关闭时取消它,任务仍可能继续运行;它抛出的异常可能只在 event loop 的异常处理器中出现。稳定做法是使用 TaskGroup 或明确的 task registry,把创建、等待、取消和异常收集放到同一个作用域。Task owner 应能回答:谁创建它,谁等待它,谁取消它,异常由谁观察。
上下文传播影响日志、trace id、request id、locale、权限标记和事务上下文。contextvars 让 async 代码能把上下文绑定到当前执行上下文;Task 创建时会捕获当前 context。跨线程或 executor 时,传播策略依赖具体 API:asyncio.to_thread() 会传播当前 context,直接使用某些 executor 接口时需要显式复制或传参。判断上下文问题时,应把 context 当成调度输入的一部分,而非普通全局变量。
阻塞调用进入 event loop 的后果是整个 loop 停顿。普通 time.sleep()、同步网络请求、同步数据库驱动、CPU-heavy 循环和长时间持有 GIL 的解析函数都会让其他 Task 无法推进。修复顺序是:优先使用真正 async 的库;无法替换时,用 asyncio.to_thread() 或 loop.run_in_executor() 移交阻塞 I/O;CPU-bound Python 工作移交进程池;需要限制数量时,在移交前放置 async semaphore 或 queue。
本章的最终判断顺序可以收束为一条路径:先判断任务成本来自等待、Python CPU、native code 还是跨系统通信;再判断状态共享是否必要、对象是否能跨边界传递;接着选择 thread、process、event loop 或 executor;最后把同步、取消、backpressure、异常传播和 shutdown 写成显式 owner 关系。并发抽象的价值不在于让代码“同时运行”,而在于让同时进行的工作仍有可解释的对象边界和生命周期边界。
最小自检任务
阅读下面的代码,判断每个函数在哪个执行边界上运行,结果和异常通过什么对象返回,取消请求能影响哪些部分,以及这段设计中最需要补强的边界是什么。
import asyncio
from concurrent.futures import ProcessPoolExecutor
async def handle_sources(sources: list[str]) -> list[ParseResult]:
loop = asyncio.get_running_loop()
payloads = await asyncio.gather(*(
asyncio.to_thread(blocking_read, source)
for source in sources
))
with ProcessPoolExecutor(max_workers=4) as pool:
futures = [
loop.run_in_executor(pool, parse_payload, payload)
for payload in payloads
]
return await asyncio.gather(*futures)
答案要点
handle_sources() 是 coroutine function,调用后得到 coroutine object;它被 event loop 包装成 Task 后才会执行。asyncio.to_thread(blocking_read, source) 把同步阻塞读取移交到线程池,返回 event loop 可等待的 coroutine;真正的 blocking_read() 在 worker thread 中运行,结果回到等待它的 Task。asyncio.gather() 聚合这些 awaitable,任何一个读取失败都会让异常从 await 点传播到 handle_sources()。
ProcessPoolExecutor 创建进程池,loop.run_in_executor(pool, parse_payload, payload) 把 CPU-bound 解析提交给进程池,并返回 asyncio.Future 风格的可等待对象;底层还有 concurrent.futures.Future 承载进程池任务状态。parse_payload() 在 child process 中运行,参数和结果需要可 pickle。解析函数抛出的异常会通过 Future 在 await asyncio.gather(*futures) 处重新抛出;进程异常退出可能让 executor 进入 broken 状态。
取消 handle_sources() 所在 Task 时,event loop 会让当前 await 点收到 CancelledError。尚未开始或尚未提交的工作可以停在 Python 调度层;已经进入 to_thread() 的 blocking_read() 通常继续在 worker thread 执行;已经进入进程池的 parse_payload() 也可能继续运行。取消请求不等于底层阻塞 I/O 或 CPU 计算已经停止,代码需要为 blocking read 设置超时,为进程池任务设计输入规模和关闭策略。
这段设计最需要补强 backpressure 和 owner 边界。所有 source 会一次性创建读取任务,读取结果会一次性保存在内存中,然后一次性提交给进程池。更稳定的设计应使用有界队列或 semaphore 控制读取并发,用流式 pipeline 减少 payload 堆积,并把 executor 的 owner 提升到服务生命周期或明确的调用作用域。若 parse_payload() 返回对象较大,还要检查跨进程序列化成本和内存副本。
本章知识点总结
- 并发边界:并发设计要先定位等待、CPU、共享状态、结果返回和关闭责任分别落在哪个 runtime 对象上。
- GIL 约束:默认 CPython 中 GIL 限制同一 interpreter 内多个线程并行执行 Python bytecode,I/O 等待和部分 native code 仍能通过线程并发推进。
- threading:
Thread封装 OS thread 的启动、运行、异常和 join 生命周期,适合 I/O-bound 工作和共享内存协作。 - 线程同步:
Lock、Condition、Event和Queue把共享状态的互斥、通知、信号和消息传递写成可检查对象。 - thread-local:
threading.local()为每个线程提供独立属性字典,适合保存线程私有连接、上下文或临时统计。 - multiprocessing:进程模型通过独立 interpreter、地址空间和 GIL 获得 CPU 并行机会,同时引入启动、pickle、IPC 和回收成本。
- 进程隔离:process isolation 分离内存、崩溃影响和权限边界,跨进程通信必须设计消息格式、所有权和关闭协议。
- asyncio:event loop 通过 Task 和 Future 驱动 coroutine,在
await等待点协作切换,适合 async I/O 和高层网络并发。 - Future 语义:
Future是异步结果占位符,保存完成、取消、异常和回调状态;concurrent.futures.Future与asyncio.Future属于不同调度语境。 - executor:
Executor统一 callable 提交、worker 分配、结果获取、异常传播和 shutdown,是线程池、进程池和 event loop 桥接的调度边界。 - 取消边界:取消 Task 或 Future 主要改变 Python 调度层状态,已经运行的线程函数、进程函数或外部副作用需要额外停止协议。
- backpressure:有界 Queue、Semaphore、stream drain、executor buffersize 或提交窗口能把下游容量反馈到上游生产速度。
- 任务 owner:每个 Thread、Process、Task、Future 和 Executor 都应有明确 owner,负责等待、取消、异常观察和资源释放。
- 判断顺序:先判断成本来源,再判断共享和传输边界,接着选择并发模型,最后检查同步、取消、backpressure、异常和 shutdown。