Python 基础体系 · 第 57/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python concurrent.futures:线程池、进程池、Future 和取消
concurrent.futures 为“把可调用对象提交给执行器,并在稍后取得结果”提供了一套统一接口。它把任务提交、任务执行和结果获取分离开来:Executor 负责调度,Future 代表一次异步执行,线程池或进程池负责真正运行任务。Python 3.14 还提供了 InterpreterPoolExecutor,但本文重点放在线程池、进程池、Future 和取消机制上。(docs.python.org)
一、先建立模型:提交任务不等于执行任务
假设有一个函数:
def work(x):
return x * 2
直接调用时,调用线程会执行函数并立即得到结果:
result = work(10)
这里的时间关系是:
调用者 ──调用 work──> work 执行 ──返回──> 调用者继续执行
而使用执行器时:
future = executor.submit(work, 10)
调用者得到的不是 20,而是一个 Future:
调用者 ──submit──> Executor ──排队/调度──> worker 执行 work
│ │
└──────────────得到 Future <───────────┘
Future 是一次异步调用的句柄。它保存了这次调用的状态、结果或异常,但不代表调用一定已经完成。Executor.submit() 会返回由执行器创建的 Future,应用代码通常不应该直接实例化 Future。(docs.python.org)
这一区分是理解整个 API 的基础:
future = executor.submit(work, 10)
# 这里 work 可能尚未运行,也可能正在运行,也可能已经完成
print(future.done())
# result() 会等待,直到任务完成、抛出异常或被取消
print(future.result())
submit() 解决的是“把任务交给谁执行”,result() 解决的是“我现在需要这个任务的结果”。如果在提交后立刻调用 result(),并发带来的重叠执行空间就可能被立即阻塞抵消。
二、Executor 的统一接口
Executor 是抽象基类,不能直接用于执行任务。常用的具体实现包括:
| 执行器 | worker 类型 | 主要隔离边界 | 典型用途 |
|---|---|---|---|
ThreadPoolExecutor |
线程 | 共享进程地址空间 | I/O 等待、调用释放 GIL 的库 |
ProcessPoolExecutor |
进程 | 独立地址空间 | CPU 密集型 Python 计算、进程级隔离 |
InterpreterPoolExecutor |
线程,每个线程一个解释器 | 解释器状态隔离 | Python 3.14 中需要多解释器并行的场景 |
这几个执行器都支持类似的 submit()、map()、shutdown() 操作,但任务数据如何传递、异常如何传播、取消能取消到什么程度,都取决于具体执行器。(docs.python.org)
1. submit():提交单个任务
from concurrent.futures import ThreadPoolExecutor
def add(a, b):
return a + b
with ThreadPoolExecutor(max_workers=2) as executor:
future = executor.submit(add, 2, 3)
result = future.result()
print(result)
预期输出:
5
执行过程可以拆成四步:
ThreadPoolExecutor创建或复用 worker 线程。submit(add, 2, 3)把函数、位置参数和关键字参数放入任务队列。- 某个 worker 取出任务并执行
add(2, 3)。 - worker 把返回值写入
Future,future.result()读取该返回值。
调用者不需要自己创建线程,也不需要手动设计“结果队列 + 完成事件”。执行器和 Future 组合完成了这些工作。
2. with:资源生命周期
推荐使用上下文管理器:
with ThreadPoolExecutor(max_workers=4) as executor:
futures = [
executor.submit(work, value)
for value in range(10)
]
results = [future.result() for future in futures]
离开 with 块时,执行器会关闭,并等待任务完成,相当于调用 shutdown(wait=True)。因此,下面这段代码不会在离开代码块后立即丢弃尚未执行的任务:
with ThreadPoolExecutor(max_workers=2) as executor:
executor.submit(time_consuming_task)
如果任务本身可能无限运行,或者需要由外部信号强制终止,就不能简单假设 with 能快速退出;后文会讨论取消和进程终止的区别。(docs.python.org)
三、线程池:共享进程中的并发执行
1. 线程池的组件关系
ThreadPoolExecutor 运行在当前进程中。多个 worker 线程共享:
- Python 进程的堆;
- 模块级变量;
- 文件描述符;
- 网络连接;
- 锁、条件变量和其他线程同步对象。
因此线程之间传递普通 Python 对象不需要序列化:
from concurrent.futures import ThreadPoolExecutor
shared = []
def append_value(value):
shared.append(value)
return value
with ThreadPoolExecutor(max_workers=4) as executor:
futures = [executor.submit(append_value, i) for i in range(5)]
for future in futures:
future.result()
print(shared)
输出中的元素顺序不应作为业务保证。多个线程可能以不同顺序执行并完成,列表的单次 append() 通常不会表现为“半条数据”,但“先后顺序”和更复杂的复合操作仍然需要明确同步协议。
例如:
if key not in cache:
cache[key] = build_value()
这不是一个不可分割的业务操作。两个线程可能同时看到 key 不存在,然后都执行 build_value()。即使字典的单个操作在某些实现中不会损坏数据,也不能据此推出整个“检查后写入”流程是原子的。
2. 线程池适合什么任务
线程池的核心价值通常不是让 Python 字节码获得多核并行,而是让一个线程等待 I/O 时,其他线程继续运行。例如:
from concurrent.futures import ThreadPoolExecutor
from urllib.request import urlopen
def fetch(url):
with urlopen(url, timeout=5) as response:
return response.status, len(response.read())
urls = [
"https://www.python.org/",
"https://docs.python.org/3/",
]
with ThreadPoolExecutor(max_workers=4) as executor:
futures = [executor.submit(fetch, url) for url in urls]
for future in futures:
try:
status, size = future.result()
except Exception as exc:
print(f"请求失败:{exc!r}")
else:
print(status, size)
一个线程阻塞在网络读取上时,其他 worker 可以处理别的请求。这里的并发来自等待时间重叠,而不是必须依赖 CPU 并行。
在传统 CPython 执行模型中,多个线程执行纯 Python CPU 密集型字节码通常受到 GIL 约束;如果任务主要是 Python 层面的计算,线程池未必能带来预期的多核加速。另一方面,某些 C 扩展在执行计算时会释放 GIL,因此不能简单把“CPU 密集型”与“线程一定无效”等同起来。选择执行器时应看实际代码路径和基准测试,而不是只看函数名称。
3. max_workers 不是“同时运行任务数”的绝对保证
max_workers 表示池中最多使用多少个 worker。它不保证:
- 每个任务立刻开始;
- 所有任务同时运行;
- 线程数越大吞吐量越高;
- 线程数越大延迟越低。
例如:
with ThreadPoolExecutor(max_workers=2) as executor:
futures = [
executor.submit(time.sleep, 10)
for _ in range(100)
]
最多只有两个任务同时占用 worker,其余任务处于等待调度状态。若任务数量远大于 worker 数量,提交动作本身可能很快完成,但任务队列会持续增长。
Python 3.14 的 ThreadPoolExecutor 默认 worker 数沿用 3.13 的规则:min(32, (os.process_cpu_count() or 1) + 4)。这只是默认策略,不是针对具体应用负载的最优值。(docs.python.org)
4. 线程初始化失败
可以用 initializer 为每个 worker 建立线程局部资源:
import threading
from concurrent.futures import ThreadPoolExecutor
local = threading.local()
def init_worker(worker_name):
local.name = worker_name
def task(value):
return threading.current_thread().name, local.name, value * 2
with ThreadPoolExecutor(
max_workers=2,
thread_name_prefix="calc",
initializer=init_worker,
initargs=("worker-resource",),
) as executor:
future = executor.submit(task, 10)
print(future.result())
如果 initializer 抛出异常,线程池会进入不可正常工作的状态,已有待处理任务以及后续提交可能得到 BrokenThreadPool。初始化阶段应尽量完成确定性的本地准备,并将不可恢复错误显式暴露,而不是让池在半初始化状态下继续接收请求。(docs.python.org)
四、进程池:独立地址空间与序列化边界
1. 进程池为什么能绕过 GIL
ProcessPoolExecutor 使用独立进程执行任务。每个进程拥有独立的 Python 运行时和地址空间,因此一个进程中的 Python 线程锁和 GIL 不会直接阻止另一个进程运行。官方文档也明确指出,进程池通过使用 multiprocessing 绕过 GIL,但代价是任务和返回值必须可 pickle。(docs.python.org)
数据流不再是线程池中的“直接传引用”:
主进程
│
├─ pickle(函数、参数)
│
├─ 进程间通道 ──> worker 进程
│ │
│ ├─ unpickle
│ ├─ 执行函数
│ └─ pickle(返回值或异常)
│
└─ <────────────── 返回结果
因此,下面这些对象通常不能直接作为进程池任务边界:
- lambda;
- REPL 中定义的函数;
- 局部函数;
- 不可 pickle 的打开连接、锁或某些扩展对象。
任务函数应放在模块顶层,并确保主模块可被 worker 导入:
from concurrent.futures import ProcessPoolExecutor
def square(value):
return value * value
def main():
with ProcessPoolExecutor() as executor:
results = list(executor.map(square, range(5)))
print(results)
if __name__ == "__main__":
main()
if __name__ == "__main__": 不是装饰性写法。某些进程启动方式会重新导入主模块;如果创建进程池的代码位于模块顶层,子进程可能再次创建进程池,造成递归启动或启动失败。进程池要求 __main__ 模块可被 worker 进程导入,因此不能依赖交互式解释器中的临时定义。(docs.python.org)
2. 参数可 pickle 不代表传输成本可以忽略
下面代码逻辑上合法:
with ProcessPoolExecutor() as executor:
future = executor.submit(sum, range(10_000_000))
但 range、参数封装、进程间通信以及返回值处理都有成本。若每个任务只计算几百微秒,序列化和调度可能比计算本身更慢。
进程池更适合这样的结构:
每个任务的有效计算时间
远大于
序列化 + 进程间通信 + 调度开销
这也是 ProcessPoolExecutor.map() 提供 chunksize 的原因。对进程池而言,输入可以按块提交,减少大量小任务逐个跨进程调度的开销;对线程池和解释器池,chunksize 没有实际效果。(docs.python.org)
3. Python 3.14 的启动方式变化
在 Python 3.14 中,ProcessPoolExecutor 的默认进程启动方式已改为不再默认使用 fork。如果应用明确要求 fork,需要显式传入 multiprocessing context:
import multiprocessing
from concurrent.futures import ProcessPoolExecutor
def task(value):
return value * value
if __name__ == "__main__":
context = multiprocessing.get_context("fork")
with ProcessPoolExecutor(mp_context=context) as executor:
print(executor.submit(task, 6).result())
启动方式会影响全局状态继承、线程与 fork 的交互、初始化成本和可移植性。多线程程序中再 fork 还可能继承处于锁定状态的运行时资源,因此不能把启动方式当作纯性能开关。Python 3.14 的默认变化以及显式指定方式的要求见官方文档。(docs.python.org)
4. worker 回收与强制终止
max_tasks_per_child 可以限制一个进程最多执行多少个任务:
with ProcessPoolExecutor(max_tasks_per_child=1000) as executor:
...
任务数量达到上限后,worker 退出并由池创建新的 worker。这可用于控制长期运行进程中的资源积累,但会引入进程重建成本;该参数与 fork 启动方式不兼容,未显式提供 context 时会倾向使用 spawn。(docs.python.org)
Python 3.14 新增:
executor.terminate_workers()
executor.kill_workers()
二者都会立即处理存活的 worker,并内部执行资源关闭逻辑。terminate_workers() 使用进程的 terminate,kill_workers() 使用更强制的 kill。调用后不应继续向执行器提交任务。强制终止意味着正在执行的任务可能没有机会完成清理、提交事务或写出完整结果,因此它们是故障恢复工具,不是普通取消 API 的替代品。(docs.python.org)
五、Future 的状态机
一个 Future 可以抽象为以下状态:
PENDING ──开始执行──> RUNNING ──成功──> FINISHED
│ │
│ └─异常──> FINISHED
│
└─cancel() 成功──> CANCELLED
状态的关键区别是:
PENDING:已提交但尚未开始执行;RUNNING:worker 已经开始执行;CANCELLED:任务尚未开始,且取消成功;FINISHED:任务已经返回或抛出异常。
对应的查询方法:
future.cancelled() # 是否成功取消
future.running() # 是否正在执行
future.done() # 是否取消或已经完成
done() 为真并不等价于“任务成功”。任务可能是正常返回、抛出异常,或者被取消。判断成功必须调用 result(),或先调用 exception() 检查异常。(docs.python.org)
1. 成功结果与异常传播
from concurrent.futures import ThreadPoolExecutor
def divide(a, b):
return a / b
with ThreadPoolExecutor(max_workers=2) as executor:
good = executor.submit(divide, 10, 2)
bad = executor.submit(divide, 10, 0)
print(good.result())
try:
bad.result()
except ZeroDivisionError as exc:
print(f"任务失败:{exc!r}")
输出类似:
5.0
任务失败:ZeroDivisionError('division by zero')
异常不会在 submit() 调用处直接抛出,因为此时 worker 可能尚未执行函数。异常被记录在 Future 中,调用 result() 时重新抛出。exception() 则返回异常对象;如果任务成功完成,它返回 None。(docs.python.org)
因此,下面的写法会漏掉任务异常:
with ThreadPoolExecutor() as executor:
for value in range(10):
executor.submit(divide, value, 0)
print("主流程结束")
主线程没有调用任何 Future.result(),所以不会在这里观察到 worker 中发生的 ZeroDivisionError。任务失败并不自动等价于主流程失败;必须定义结果消费和错误汇报路径。
2. result(timeout) 的超时不是取消
from concurrent.futures import ThreadPoolExecutor, TimeoutError
import time
def slow_task():
time.sleep(5)
return "done"
with ThreadPoolExecutor() as executor:
future = executor.submit(slow_task)
try:
print(future.result(timeout=1))
except TimeoutError:
print("等待超时")
print("任务是否已取消:", future.cancel())
可能输出:
等待超时
任务是否已取消: False
result(timeout=1) 只表示调用者最多等待 1 秒。如果任务仍在运行,它抛出 TimeoutError,但不会自动停止任务。随后 cancel() 仍可能失败,因为任务已经进入 RUNNING 状态。
这两个动作的因果关系必须分开:
等待超时 ──只改变调用者的等待结果
取消成功 ──改变尚未开始任务的执行状态
六、取消:只能取消尚未开始的任务
1. Future.cancel()
from concurrent.futures import ThreadPoolExecutor
import time
def task(name):
print(f"{name} started")
time.sleep(1)
return name
with ThreadPoolExecutor(max_workers=1) as executor:
first = executor.submit(task, "first")
second = executor.submit(task, "second")
print("second cancel:", second.cancel())
print("first result:", first.result())
try:
print(second.result())
except Exception as exc:
print(type(exc).__name__)
由于只有一个 worker:
first先运行;second仍处于等待状态;second.cancel()有机会成功;second.result()抛出CancelledError。
预期输出类似:
second cancel: True
first started
first result: first
CancelledError
如果任务已经运行或已经完成,cancel() 返回 False;只有尚未开始执行的任务才可能被取消。(docs.python.org)
2. 取消竞态
以下代码不能保证 cancel() 一定成功:
future = executor.submit(task)
cancelled = future.cancel()
在 submit() 返回和 cancel() 执行之间,worker 可能已经取走任务并开始运行。于是存在竞态:
线程 A:submit() 返回 Future
线程 B:worker 取出任务
线程 B:Future 进入 RUNNING
线程 A:cancel() 返回 False
因此,取消是一个“尝试取消排队任务”的操作,不是对任务执行过程的强制控制。
3. shutdown(cancel_futures=True)
当需要关闭执行器并取消尚未开始的任务时,可以使用:
from concurrent.futures import ThreadPoolExecutor
import time
def task(index):
print(f"start {index}")
time.sleep(2)
return index
executor = ThreadPoolExecutor(max_workers=2)
futures = [
executor.submit(task, index)
for index in range(6)
]
executor.shutdown(wait=False, cancel_futures=True)
for index, future in enumerate(futures):
if future.cancelled():
print(index, "cancelled")
elif future.done():
print(index, "finished")
else:
print(index, "running or pending")
cancel_futures=True 只取消尚未开始的任务;已经完成或正在运行的任务不会被取消。wait=False 只让 shutdown() 尽快返回,并不意味着 Python 进程会立刻退出;程序仍会等待 pending futures 最终完成或被处理。(docs.python.org)
若 wait=True 且 cancel_futures=True,已开始的任务会完成后 shutdown() 才返回,剩余未开始的任务被取消。这适合“停止接收新任务,排空已运行任务,放弃排队任务”的关闭策略。
七、协作式停止:真正停止正在运行的任务
由于 Future.cancel() 不能停止正在运行的函数,长任务需要自己检查停止信号。线程中可使用 threading.Event:
from concurrent.futures import ThreadPoolExecutor
import threading
import time
stop_event = threading.Event()
def cancellable_task():
for step in range(100):
if stop_event.is_set():
return "stopped cooperatively"
# 模拟一小段工作
time.sleep(0.05)
return "finished"
with ThreadPoolExecutor(max_workers=1) as executor:
future = executor.submit(cancellable_task)
time.sleep(0.2)
stop_event.set()
print(future.result())
这里的停止路径是:
主线程 set()
↓
worker 下一次检查 Event
↓
函数主动返回
↓
Future 进入 FINISHED
这不是 Future 层面的取消,因为 future.cancelled() 最终仍为 False;任务是正常返回了一个“主动停止”的结果。也可以让函数抛出自定义异常,但调用方必须区分“外部请求停止”和“任务自身失败”。
循环中的检查间隔决定停止延迟。如果函数一次调用阻塞在不可中断的系统调用、第三方库或无限期 I/O 中,外部 Event 也不能凭空打断它。此时需要为 I/O 设置超时,或使用库本身提供的关闭接口。
进程池中的停止信号不能简单依赖线程共享对象。进程拥有独立地址空间,应使用 multiprocessing 提供的进程间通信机制,或者在无法协作停止时使用 terminate_workers() / kill_workers()。后两者可能跳过业务清理,必须配合幂等任务和外部一致性恢复。
八、按完成顺序消费:as_completed()
提交顺序和完成顺序通常不同:
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task(index, delay):
time.sleep(delay)
return index
jobs = [
(0, 0.3),
(1, 0.1),
(2, 0.2),
]
with ThreadPoolExecutor(max_workers=3) as executor:
future_to_index = {
executor.submit(task, index, delay): index
for index, delay in jobs
}
for future in as_completed(future_to_index):
index = future_to_index[future]
try:
result = future.result()
except Exception as exc:
print(index, "failed:", repr(exc))
else:
print(index, "done:", result)
预期完成顺序通常是:
1 done: 1
2 done: 2
0 done: 0
as_completed() 返回一个迭代器,每当某个 future 完成或被取消,就产生对应的 future。它适合:
- 谁先完成就先处理谁;
- 先完成的结果可以尽快发送给下游;
- 某个任务失败后尝试取消其他尚未开始的任务。
但要注意,as_completed() 只改变结果消费顺序,不改变任务已经被提交的事实。若希望限制在途任务数量,应控制提交量,而不是提交无限任务后只用 as_completed() 消费。
九、批量等待:wait()
wait() 返回两个集合:
done, not_done = wait(
futures,
timeout=2,
return_when=FIRST_COMPLETED,
)
其中:
done包含已经完成或已经取消的 future;not_done包含仍未完成的 future;FIRST_COMPLETED表示任意一个 future 完成或取消时返回;FIRST_EXCEPTION表示某个 future 因异常完成时返回;如果没有异常,则等价于等待全部完成;ALL_COMPLETED等待全部 future 完成或取消。(docs.python.org)
一个“首个成功结果”的实现不能直接使用 FIRST_EXCEPTION,因为它关注的是异常而不是成功:
from concurrent.futures import (
ThreadPoolExecutor,
as_completed,
)
def query(source):
# 模拟不同数据源
...
with ThreadPoolExecutor(max_workers=5) as executor:
futures = [executor.submit(query, source) for source in sources]
winner = None
try:
for future in as_completed(futures):
try:
winner = future.result()
except Exception:
continue
else:
break
finally:
for future in futures:
future.cancel()
这里的 cancel() 只能减少尚未开始的任务;已经运行的查询仍可能继续。若查询函数支持超时或停止信号,应把停止机制传入任务本身。
十、map():有序结果与背压
Executor.map() 适合“多个输入应用同一个函数”的场景:
from concurrent.futures import ThreadPoolExecutor
def square(value):
return value * value
with ThreadPoolExecutor(max_workers=3) as executor:
for result in executor.map(square, range(5)):
print(result)
输出顺序是:
0
1
4
9
16
即使后面的任务先完成,map() 迭代器仍按输入顺序产生结果。这个特性适合批处理,但可能造成队头阻塞:
任务 0 很慢
任务 1、2、3 已完成
调用 map 迭代器仍必须先等待任务 0
与之相比,as_completed() 按完成顺序返回。
Python 3.14 为 map() 增加了 buffersize。默认情况下,输入迭代器可能被立即收集并提交;设置 buffersize 后,可以限制“已经提交但结果尚未被迭代器取出”的任务数量。当缓冲区满时,输入迭代会暂停,直到调用方消费结果。这提供了一种提交端背压。(docs.python.org)
with ThreadPoolExecutor(max_workers=4) as executor:
for result in executor.map(
square,
range(1_000_000),
buffersize=100,
):
consume(result)
这里的 buffersize 不是 worker 数,也不是任务执行超时时间:
max_workers限制并发 worker 数;buffersize限制提交和消费之间的在途结果规模;timeout限制从map()调用开始计算的等待时间。
十一、回调:完成通知而不是结果替代品
可以通过 add_done_callback() 注册完成回调:
from concurrent.futures import ThreadPoolExecutor
def report(future):
try:
print("result:", future.result())
except Exception as exc:
print("error:", repr(exc))
with ThreadPoolExecutor(max_workers=2) as executor:
future = executor.submit(pow, 2, 10)
future.add_done_callback(report)
回调会在 future 被取消或执行结束时调用。如果 future 已经完成,添加回调时可能立即执行。回调接收 future 作为唯一参数;回调自身抛出的普通 Exception 会被记录并忽略,因此重要的错误不能只依赖回调抛出。(docs.python.org)
回调适合轻量通知,例如:
- 增加完成计数;
- 放入另一个队列;
- 记录任务状态;
- 触发非阻塞的后续调度。
不应在回调中执行长时间阻塞工作,也不应在回调里无条件等待另一个依赖当前执行器的 future,否则容易制造线程饥饿或死锁。
十二、线程池中的 Future 死锁
线程池最容易被忽略的故障是:worker 等待另一个也需要该线程池执行的 future。
1. 两个任务互相等待
from concurrent.futures import ThreadPoolExecutor
executor = ThreadPoolExecutor(max_workers=2)
def wait_on_a():
return future_a.result()
def wait_on_b():
return future_b.result()
future_a = executor.submit(wait_on_b)
future_b = executor.submit(wait_on_a)
两个 worker 分别执行 wait_on_b() 和 wait_on_a():
worker 1:等待 future_b
worker 2:等待 future_a
但两个 future 都必须等各自的函数返回,而两个函数都在等待对方,因此形成环路。
2. 单 worker 自我等待
from concurrent.futures import ThreadPoolExecutor
executor = ThreadPoolExecutor(max_workers=1)
def outer():
inner = executor.submit(pow, 2, 10)
return inner.result()
future = executor.submit(outer)
print(future.result())
状态变化是:
唯一 worker 执行 outer
↓
outer 提交 inner
↓
inner 排队,因为没有空闲 worker
↓
outer 等待 inner
↓
永远没有 worker 执行 inner
官方文档明确列出了这两类死锁示例。(docs.python.org)
避免方法不是简单地“把 max_workers 调大”。根本问题是任务图中存在同步等待环。更稳妥的结构是:
worker 只执行叶子任务
主线程或调度层组合结果
例如先提交 a 和 b,再由主线程调用它们的 result(),而不是让池中的任务相互等待。
十三、异常、取消和执行器损坏的区别
这几个状态必须分开处理:
1. 任务异常
任务函数运行了,但函数抛出异常:
try:
result = future.result()
except ValueError:
...
此时执行器通常仍然可以继续提交其他任务。
2. 任务取消
任务尚未开始,被成功取消:
try:
future.result()
except CancelledError:
...
取消不是业务异常,也不表示函数执行失败;函数可能根本没有被调用。
3. 执行器损坏
如果线程 initializer 失败,可能出现 BrokenThreadPool;如果进程 worker 异常退出,可能出现 BrokenProcessPool。这表示执行器本身无法可靠地继续承担任务,而不仅是某一个任务失败。(docs.python.org)
from concurrent.futures import (
ProcessPoolExecutor,
BrokenProcessPool,
)
try:
with ProcessPoolExecutor() as executor:
results = list(executor.map(task, inputs))
except BrokenProcessPool:
# 记录 worker 异常退出,丢弃旧 executor,重建前先检查原因
...
重建执行器之前,应先确认原因:
- worker 是否被操作系统杀死;
- 是否发生内存耗尽;
- 是否有 native 扩展崩溃;
- 是否违反 pickle 或启动方式要求;
- 是否在进程任务中调用了 executor/Future 方法。
进程池中,worker 任务调用 Executor 或 Future 方法可能导致死锁;这类嵌套调度不能套用线程池中的直觉。(docs.python.org)
十四、线程池和进程池的端到端对比
下面用同一个 CPU 任务比较调用方式,而不是预设某个固定性能结论:
from concurrent.futures import (
ThreadPoolExecutor,
ProcessPoolExecutor,
)
import math
import time
def count_primes(limit):
count = 0
for number in range(2, limit):
is_prime = True
for divisor in range(2, math.isqrt(number) + 1):
if number % divisor == 0:
is_prime = False
break
if is_prime:
count += 1
return count
def run(executor_cls, limits):
start = time.perf_counter()
with executor_cls(max_workers=4) as executor:
futures = [
executor.submit(count_primes, limit)
for limit in limits
]
results = [future.result() for future in futures]
elapsed = time.perf_counter() - start
return results, elapsed
if __name__ == "__main__":
limits = [20_000, 21_000, 22_000, 23_000]
print("threads:", run(ThreadPoolExecutor, limits))
print("processes:", run(ProcessPoolExecutor, limits))
这个例子的重点不是某台机器上的具体秒数,而是执行边界:
- 线程池共享内存,参数传递便宜,但纯 Python 字节码不一定获得多核并行;
- 进程池有独立地址空间,可以使用多个 CPU 核心,但需要序列化参数和返回值;
- 任务太小时,进程通信成本可能超过并行收益;
- 任务函数使用全局可变状态时,线程池和进程池的语义完全不同。
在真实系统中,应测量端到端耗时,包括:
任务排队时间
+ worker 执行时间
+ 数据序列化时间
+ 进程间传输时间
+ 结果消费时间
只测函数体内部时间,无法代表执行器方案的实际成本。
十五、与 asyncio 的关系:两个 Future 不能混用
concurrent.futures.Future 和 asyncio.Future 名字相同,但用途不同。前者属于线程池、进程池等同步并发执行器;后者属于 asyncio 事件循环和协程调度模型,不能把一个当成另一个直接 await。官方文档也特别提醒二者不要混淆。(docs.python.org)
在异步程序中,如果必须把阻塞函数放到线程池,可以使用 asyncio.to_thread(),或使用事件循环提供的执行器接口:
import asyncio
import time
def blocking_work(value):
time.sleep(1)
return value * 2
async def main():
result = await asyncio.to_thread(blocking_work, 21)
print(result)
asyncio.run(main())
这里的协作关系是:
asyncio 事件循环
│
└─把阻塞函数交给线程执行
│
└─返回可等待的 asyncio 对象
如果直接在协程中调用 blocking_work(21),事件循环线程会被 sleep() 阻塞,其他协程无法正常推进。线程池解决的是阻塞函数的执行位置,asyncio 解决的是异步任务的组织方式;它们可以组合,但抽象层次不同。
十六、生产边界:提交、等待和关闭必须形成闭环
一个可靠的执行流程至少应明确以下路径:
flowchart TD
A[创建 Executor] --> B[提交任务]
B --> C{Future 状态}
C -->|PENDING| D[等待或取消]
C -->|RUNNING| E[等待结果或协作式停止]
C -->|FINISHED| F[读取结果或异常]
C -->|CANCELLED| G[处理取消]
B --> H[任务提交失败]
H --> I[检查 BrokenExecutor 或参数错误]
F --> J[关闭 Executor]
G --> J
E --> J
D --> J
其中最容易出错的是“任务已经提交,但没有任何结果消费路径”。例如:
futures = [executor.submit(task, item) for item in items]
如果后续不调用 result()、wait()、as_completed() 或回调,应用就可能无法发现任务异常,也无法统计任务是否完成。
另一个常见错误是把执行器当作无限队列:
while True:
executor.submit(handle, next_item())
这会让待处理任务数量持续增长,最终表现为内存上涨、延迟扩大,甚至关闭时长不可控。map(buffersize=...)、显式生产者—消费者队列或限制批量提交,都可以让系统拥有明确的在途任务上限。
最后,关闭策略必须和业务语义一致:
- 优雅关闭:停止接收新任务,等待已提交任务完成;
- 放弃排队任务:
shutdown(cancel_futures=True); - 协作式停止运行任务:使用事件、超时或业务取消协议;
- 强制结束进程任务:
terminate_workers()或kill_workers(),并接受清理不完整的风险。
十七、一个完整的可取消批处理示例
下面的示例组合了线程池、Future、按完成顺序消费、异常处理和协作式停止:
from concurrent.futures import (
ThreadPoolExecutor,
as_completed,
)
import threading
import time
class StopRequested(Exception):
pass
def process_item(item, stop_event):
for _ in range(10):
if stop_event.is_set():
raise StopRequested(f"item {item} stopped")
time.sleep(0.05)
if item == 3:
raise ValueError("invalid item")
return item * 10
def main():
stop_event = threading.Event()
items = range(8)
with ThreadPoolExecutor(
max_workers=3,
thread_name_prefix="batch",
) as executor:
future_to_item = {
executor.submit(process_item, item, stop_event): item
for item in items
}
try:
for future in as_completed(future_to_item):
item = future_to_item[future]
try:
result = future.result()
except StopRequested as exc:
print(f"停止:{exc}")
except Exception as exc:
print(f"失败 item={item}: {exc!r}")
# 如果一个错误足以终止整个批次:
stop_event.set()
# 只能取消仍未开始的任务
for other in future_to_item:
if other is not future:
other.cancel()
else:
print(f"成功 item={item}: {result}")
finally:
# 防止任务在离开 with 后仍继续做长时间工作
stop_event.set()
if __name__ == "__main__":
main()
这个程序体现了几个重要事实:
as_completed()让已经完成的任务先被处理。future.result()是观察任务异常的地方。stop_event.set()只提出停止请求,任务必须主动检查。other.cancel()只对尚未开始的任务有效。- 离开
with时仍会执行执行器关闭逻辑,因此协作式停止信号应尽早发出。 StopRequested与ValueError分开处理,避免把主动停止误报成业务失败。
concurrent.futures 的核心并不是“自动让代码变快”,而是提供了一个可组合的执行协议:
Executor:在哪里运行
Future:一次运行的状态和结果是什么
result / wait / as_completed:如何观察完成
cancel / shutdown:如何放弃尚未开始的工作
Event / timeout / terminate:如何处理正在运行或失控的工作
只要把这些边界区分清楚,线程池和进程池就不再是两个“换一个类名”的 API,而是两种不同的数据传递、故障传播和生命周期模型。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python 多进程:启动方式、IPC、共享内存、Pool 与回收
- 下一篇:Python 队列与背压:queue、asyncio.Queue、容量和关闭协议
- 延伸:Python 线程:生命周期、锁、条件变量、竞态和死锁诊断
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论