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

执行过程可以拆成四步:

  1. ThreadPoolExecutor 创建或复用 worker 线程。
  2. submit(add, 2, 3) 把函数、位置参数和关键字参数放入任务队列。
  3. 某个 worker 取出任务并执行 add(2, 3)
  4. worker 把返回值写入 Futurefuture.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() 使用进程的 terminatekill_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:

  1. first 先运行;
  2. second 仍处于等待状态;
  3. second.cancel() 有机会成功;
  4. 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=Truecancel_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 只执行叶子任务
主线程或调度层组合结果

例如先提交 ab,再由主线程调用它们的 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 任务调用 ExecutorFuture 方法可能导致死锁;这类嵌套调度不能套用线程池中的直觉。(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.Futureasyncio.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()

这个程序体现了几个重要事实:

  1. as_completed() 让已经完成的任务先被处理。
  2. future.result() 是观察任务异常的地方。
  3. stop_event.set() 只提出停止请求,任务必须主动检查。
  4. other.cancel() 只对尚未开始的任务有效。
  5. 离开 with 时仍会执行执行器关闭逻辑,因此协作式停止信号应尽早发出。
  6. StopRequestedValueError 分开处理,避免把主动停止误报成业务失败。

concurrent.futures 的核心并不是“自动让代码变快”,而是提供了一个可组合的执行协议:

Executor:在哪里运行
Future:一次运行的状态和结果是什么
result / wait / as_completed:如何观察完成
cancel / shutdown:如何放弃尚未开始的工作
Event / timeout / terminate:如何处理正在运行或失控的工作

只要把这些边界区分清楚,线程池和进程池就不再是两个“换一个类名”的 API,而是两种不同的数据传递、故障传播和生命周期模型。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。