跳到主要内容

线程、进程、任务与 asyncio

本节目标

查询 Python 线程、进程、Future、TaskGroup、取消、背压和关闭边界。

并发设计的难点不在于“同时启动”,而在于明确工作由谁拥有、失败如何传播、等待何时有界,以及退出时哪些资源必须被观察和回收。本章以 threadingfree-threaded Pythonqueuemultiprocessingconcurrent.futuresasyncio协程与任务异步同步原语异步队列的 Python 3.14 合同为准。异步语法和协议见异步语法与协议,外部程序边界见操作系统、命令行参数与子进程,网络资源的生命周期见网络与常用互联网协议

并发模型与选择方法

先描述负载、隔离、共享状态和取消合同,再选择工具。线程共享进程地址空间,适合多个阻塞 I/O 工作,但任一线程破坏共享对象都可能影响整个进程;进程隔离地址空间,能利用多个 CPU 核心执行 Python 代码,却引入启动、序列化、IPC 和回收成本;asyncio 让一个事件循环协作调度任务,适合大量支持异步协议的 I/O,不会自动把阻塞函数变成非阻塞。Executor 则为同步调用提供统一的提交和结果接口,底层仍可能是线程或进程。

模型不能只按“快不快”决定。CPU 密集工作要先测量任务粒度、序列化量和部署平台;短小工作放进进程池可能比串行更慢。共享可变状态越多,锁顺序和失败恢复越复杂;跨进程消息虽需复制或序列化,却能形成更清楚的所有权边界。asyncio 任务只在 await 等显式让出点切换,但跨让出点的不变量同样需要设计。

无论选择哪种方案,都要写出终止清单:谁发出停止信号,谁停止接收新工作,谁排空或放弃已接收工作,谁等待 worker,谁观察异常,以及超时后采取什么有界升级动作。daemon、垃圾回收和解释器退出都不能替代正常清理协议。

线程与 CPython GIL 构建边界

默认 CPython 构建中的全局解释器锁限制同一时刻执行 Python 字节码的线程数量,因此普通 CPU 密集 Python 代码通常不会仅因增加线程就获得多核并行;阻塞 I/O 和某些释放 GIL 的扩展操作仍可与其他线程重叠。这是 CPython 实现和构建配置的边界,不是 Python 语言承诺。

更重要的是,GIL 不会把跨越多条操作的共享状态不变量自动变成原子事务。读取、判断再写入,或同时更新两个容器,即使每个单独方法当前受到内部保护,也可能在步骤之间被其他线程观察。应以 Lock 或更高层消息传递保护完整不变量,不依赖“这段代码在带 GIL 的版本碰巧没出错”。

free-threaded 构建允许禁用 GIL,但不免除数据竞争责任;内置容器会使用内部锁维持与默认构建相近的基本线程安全,却不承诺用户的复合操作原子。运行时还可能通过 -X gilPYTHON_GIL,或因导入未声明支持 free threading 的 C 扩展而启用 GIL。性能和正确性测试都应记录解释器构建、sys._is_gil_enabled() 的实际结果与扩展兼容性,不能从版本号推断运行状态。

LockConditionEvent

Lock 保护互斥临界区;优先使用 with lock:,让异常路径也释放锁,并建立全局一致的多锁获取顺序。原始 Lock 不记录所有者,任何线程都可释放一个已锁定的锁;需要同一线程递归获取时才使用 RLock。无论哪种锁,都不要在持锁时执行无界 I/O、等待另一个线程,或调用可能回调未知代码的接口。

Condition 把锁与“状态可能已变化”的通知关联起来。等待者必须在持有关联锁时检查业务谓词;Condition.wait() 会释放关联锁,并在返回前重新取得它。返回不证明谓词为真:通知可能被其他等待者先消费,状态也可能再次变化,因此要用 while 重新检查谓词,或直接使用 wait_for()notify()notify_all() 不会立即释放锁,被唤醒者要等通知者退出临界区后才能继续。

Event 是共享布尔标志:set() 唤醒等待者,clear() 重置标志,但它既不携带队列项目,也不确认每个接收者已处理通知。它适合表示“允许开始”或“请求停止”,不适合计数和一次一消费者的工作分发。线程没有安全的强制停止 API;正常退出应由非 daemon 线程定期观察事件,完成清理后由拥有者 join()

Queue 与生产者消费者

queue.Queue 在多个线程间提供同步的 put()get();设置有限 maxsize 可以让生产者在消费者落后时阻塞,从而把容量变成显式背压。qsize()empty()full() 都只是瞬时近似,检查之后状态可能立刻变化,不能用它们代替带 timeout 的实际操作。

任务完成跟踪与项目出队是两件事。每次成功的 get() 最终恰好对应一次 task_done(),通常把处理逻辑放在 try/finally 中;调用次数过多会抛出 ValueError,遗漏则让 join() 永久等待。join() 只等待未完成任务计数归零,不负责关闭工作线程,也不传播 worker 异常;线程异常、停止信号和 join() 仍需单独拥有和观察。

Python 3.14 的 shutdown(immediate=False) 禁止继续 put(),并让已阻塞的生产者以 ShutDown 失败;默认模式允许消费者排空现有项目,队列为空后 get() 才抛出 ShutDownshutdown(immediate=True) 会排空队列并可能在工作并未执行时解除等待,因此 shutdown(immediate=True) 会打破通常的 join() 完成不变量。只有业务明确选择“放弃排队工作”时才能使用立即关闭,并另行记录哪些工作未完成。

进程启动方式与序列化

multiprocessingspawn 会启动全新解释器,forkserver 通过专用服务器派生子进程。forkserver 服务器通常保持单线程,除非系统库或预加载导入以副作用启动线程;只有在这一前提下,服务器使用 fork() 才一般安全。fork 则复制当前进程状态。可用方式和默认值随平台而异;Python 3.14 已不在任何平台默认使用 fork,支持所需设施的 POSIX 平台通常默认 forkserver,Windows 与 macOS 默认 spawn。不要让正确性依赖隐式默认值:启动方式应通过 context 显式选择并传给所用 API,库应接受调用者提供的 context,而不是擅自调用一次性的全局 set_start_method()

使用 spawnforkserver 时,子解释器需要重新导入主模块。目标可调用对象、参数和返回结果必须满足 pickle 与可导入边界;局部函数、lambda、打开的文件句柄和许多运行时对象不能直接跨越这一边界。即使某对象能够 pickle,也不能把来自不可信来源的 pickle 当数据格式加载,因为反序列化可以执行任意代码。

创建进程、池或 manager 的入口放在 if __name__ == "__main__": 保护下,避免子进程导入模块时递归创建新进程。冻结可执行文件、交互解释器和平台资源继承还有额外限制;测试至少覆盖实际部署平台和所选启动方式,而不是只在默认 fork 环境偶然通过。

进程池、IPC 与清理

进程池适合把足够大的独立工作分配给固定数量的 worker。任务太小会被序列化、排队和进程切换成本吞没;任务太大又会降低负载均衡和取消响应。map() 的 chunk、worker 数、内存峰值以及 maxtasksperchild 都应由测量决定。提交后要消费 AsyncResult 或 Future,不能让子进程异常只留在无人读取的结果对象里。

池的关闭方式代表不同业务语义。正常收尾使用 close()join(),异常放弃工作才使用 terminate()multiprocessing.Pool 的上下文管理器在退出时调用 terminate(),不能把它描述成保证排空全部任务的优雅关闭。强制终止持有锁、队列或 pipe 的进程可能破坏共享资源,之后继续复用这些 IPC 对象并不安全。

QueuePipe、共享内存和 manager proxy 都是协议边界,不是普通本地对象。经 pipe/queue 发送的对象会被 pickle,单条消息和总量都应受限,只有可信对端才能反序列化。IPC 端点也有独立的 close() 与后台线程回收责任:不用的 pipe 端及时关闭,进程队列写入结束后按所有权调用 close()、必要时 join_thread(),并在消费队列数据后再等待生产进程,避免 feeder thread 与 join() 相互等待。

ExecutorFuture 与结果

Executor.submit() 把可调用对象封装为 Future;调用者通过 done()exception()result()as_completed() 观察完成,而不是从提交顺序推断执行顺序。Future.result() 会返回值或重新抛出工作函数的异常;timeout 只限制调用者等待,不等于工作已停止。cancel() 只可能取消尚未开始运行的工作,返回 False 时仍需决定如何等待或隔离其副作用。

executor 也是需要关闭的资源。shutdown(wait=True) 会拒绝新提交并等待已提交工作完成;cancel_futures=True 只取消尚未开始的 future,正在运行的仍会结束。with Executor(...) 等价于以等待方式退出,适合有界工作;若工作函数可无限阻塞,退出同样可能无限等待,必须在工作协议自身设置 timeout 和停止机制。

工作函数在容量有限的同一 executor 中等待另一个 Future,可能形成死锁;单 worker 中嵌套提交并等待是最小例子。ProcessPoolExecutor 还要求 __main__ 可导入,调用对象、参数和返回值可 pickle,并禁止从其 worker 内调用同一 executor 或 Future 方法。线程池适合阻塞 I/O,进程池适合经测量值得支付进程边界成本的 CPU 工作;两者都不能替调用者定义幂等、重试和副作用语义。

事件循环、协程与 Task

调用协程函数只创建协程对象;它必须被 await、包装为 Task,或交给结构化并发设施才会执行。通常以一次 asyncio.run(main()) 管理顶层事件循环,不在已经运行的循环内嵌套调用,也不把同步阻塞函数直接放进 loop 线程。语言层的 await、异步迭代和异步上下文协议已在第 13 章定义,本节关注调度与所有权。

asyncio.create_task() 会立即把协程安排进当前 loop。事件循环只保存 Task 的弱引用,所以调用者必须保存强引用并观察结果或异常;仅仅创建后丢弃变量,任务可能在完成前消失,失败也可能变成“Task exception was never retrieved”诊断。长期后台任务应保存在集合中,并用 done callback 只做移除等明确动作;应用关闭时仍要等待或取消并 await 它们。

任务取消是在下一个可取消等待点注入 CancelledError 的协作协议,不保证请求发出后立即停止。没有 await 的长 CPU 循环会阻塞整个事件循环;需要频繁 CPU 运算时应切分、转移到合适 executor,或改用进程。跨线程触碰 loop 对象还必须使用对应的线程安全入口,不能因为对象是 Python 对象就任意调用。

TaskGroup、取消与 timeout

TaskGroup 为相关子任务建立词法所有权:在 async with 内创建任务,退出上下文时统一观察。TaskGroup 退出上下文前会等待组内所有任务结束;非取消失败会取消其余任务,并在清理后组合为 ExceptionGroup(适当时为 BaseExceptionGroup)。KeyboardInterruptSystemExit 有特殊的重新抛出规则,不能当普通异常组处理。

取消代码使用 try/finally 释放资源;若捕获取消来补充日志或回滚,清理完成后通常必须重新抛出 CancelledError。它直接继承 BaseException,宽泛捕获 Exception 不会吞掉它。结构化并发和 timeout 都用取消实现,吞掉取消或错误使用 uncancel() 会破坏父子传播和剩余取消计数,只有确实要消除取消状态的高级协议才应这样做。

asyncio.timeout() 在上下文内部以取消实现,并在上下文外转换为 TimeoutError,因此 TimeoutError 要在 async with 外捕获。wait_for() 则取消所等待对象并等待它完成取消,所以总耗时可能超过名义 timeout;要避免取消某个共享任务,可按合同使用 shield(),但仍必须保存并最终观察该任务。

concurrency_report.py
import asyncio


async def square(index: int, value: int, results: list[int | None]) -> None:
results[index] = value * value


async def wait_until_cancelled() -> None:
await asyncio.Event().wait()


async def main() -> tuple[list[int], int, bool, bool, int]:
results: list[int | None] = [None, None, None]
async with asyncio.TaskGroup() as group:
for index, value in enumerate([1, 2, 3]):
group.create_task(square(index, value, results))

ordered_results = [value for value in results if value is not None]
cancelled_task = asyncio.create_task(wait_until_cancelled())
cancelled_task.cancel()
try:
await cancelled_task
except asyncio.CancelledError:
cancelled = True
else:
cancelled = False

current = asyncio.current_task()
pending = sum(
task is not current and not task.done()
for task in asyncio.all_tasks()
)
return ordered_results, 3, ordered_results == [1, 4, 9], cancelled, pending


results, task_count, ordered, cancelled, pending = asyncio.run(main())
print(f"results={results}")
print(f"task-count={task_count}")
print(f"ordered={ordered}")
print(f"cancelled={cancelled}")
print(f"pending={pending}")
results=[1, 4, 9]
task-count=3
ordered=True
cancelled=True
pending=0

代表脚本只演示 asyncio 调度:TaskGroup 的三个任务写入预先分配的位置,因此结果不依赖完成先后;另一个任务等待本地 Event,在尚未执行用户工作前被显式取消并观察。最后排除当前任务后检查 all_tasks(),证明没有未完成 child。脚本不创建线程或进程,不访问网络,不读取墙钟,也不使用 sleep 或计时轮询。

异步同步、背压与关闭

asyncio.LockEventConditionSemaphoreBarrier 协调同一事件循环内的任务;它们不是线程安全对象。与线程版一样,锁保护的是跨多个步骤的不变量,condition 等待者要循环检查谓词,event 只是共享标志。任何持锁代码都应尽快到达可控出口;在锁内等待可能反向依赖本锁的任务会造成协作式死锁。

asyncio.Queue 不是线程安全队列,也没有操作自身的 timeout 参数;需要有界等待时用 asyncio.timeout() 包围 put()get()join()。有界 maxsizeput() 在容量耗尽时形成背压;消费者每次成功 get() 后仍须恰好一次 task_done()。Python 3.14 的 shutdown() 可停止增长并正常排空;immediate=True 同样会让 join() 在工作未执行时解除,不能把它当完成证明。

关闭顺序通常是停止生产、关闭入口、排空已接收工作、通知消费者结束,再等待所有 task。若选择放弃,则取消任务并逐一 await,区分预期的取消与真正失败;退出前等待、取消并观察每一个已创建任务。信号处理、服务器和传输对象还有各自的 close()/wait_closed() 合同,不能只看到 pending=0 就推断外部资源已关闭。测试应使用内存事件、固定输入和显式状态断言,不用公网、DNS、真实 sleep 或调度先后制造“恰好成功”的时序。