本节目标:用统一接口管理线程池与进程池,掌握 Future、map/as_completed、异常传播、进程池守卫与共享数据,并理解两者的开销差异。
适用版本:Python 3.12+(实测 3.14.6)
12.2 concurrent.futures 与 multiprocessing
12.1 得出结论:CPU 密集该用进程,IO 密集该用线程。但直接手写 threading.Thread / multiprocessing.Process 并管理一堆 join,既啰嗦又容易出错。标准库的 concurrent.futures 提供了一套统一接口,把线程池和进程池抽象成同一种用法——换一行代码就能切换。
12.2.1 两种执行器,一套接口
| 执行器 | 并行单位 | 适合 | 启动成本 |
|---|---|---|---|
ThreadPoolExecutor | 线程 | IO 密集 | 极低 |
ProcessPoolExecutor | 进程 | CPU 密集 | 高(spawn + pickle 序列化) |
两者都提供 submit、map、shutdown,都支持 with 上下文管理器(退出时自动关闭并等待)。下面统一用 ThreadPoolExecutor 演示,把类名换成 ProcessPoolExecutor 即可迁移到进程池。
12.2.2 submit、Future 与异常传播
submit() 提交一个任务,立即返回一个 Future 对象——它是一张「提货单」,代表「将来会有结果」:
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
def fetch(name: str) -> str:
time.sleep(0.1)
return f"{name} done"
with ThreadPoolExecutor(max_workers=4) as ex:
fut = ex.submit(fetch, "A") # 提交后立刻返回 Future
print("Future 类型:", type(fut).__name__)
print("结果:", fut.result()) # result() 阻塞到完成
print("map 保序:", list(ex.map(fetch, ["A", "B", "C", "D"])))
futures = {ex.submit(fetch, n): n for n in ["A", "B", "C", "D"]}
for f in as_completed(futures): # 谁先完成谁先返回
print("as_completed:", f.result())
Future 类型: Future
结果: A done
map 保序: ['A done', 'B done', 'C done', 'D done']
as_completed: B done
as_completed: D done
as_completed: A done
as_completed: C done
map(fn, iterable):按输入顺序返回结果(上面恒为 A/B/C/D 顺序),内部自动并发。as_completed(futures):谁先完成先产出,顺序不定(每次运行可能不同),适合「先到先处理」。future.result():取值,未完成时阻塞。
关键点:任务里的异常不会立刻抛出,而是被存进 Future。下面这段代码主线程完全不会崩溃,直到你调用 result():
import time
from concurrent.futures import ThreadPoolExecutor
def boom(x: int) -> int:
raise ValueError(f"处理 {x} 失败")
with ThreadPoolExecutor(max_workers=2) as ex:
f = ex.submit(boom, 7)
print("提交后主线程未崩溃,异常类型:", type(f.exception()).__name__)
try:
f.result() # 异常在这里才抛出
except ValueError as e:
print("调用 result() 时抛出:", e)
with ThreadPoolExecutor(max_workers=1) as ex:
f = ex.submit(time.sleep, 1.0)
try:
f.result(timeout=0.2)
except TimeoutError:
print("超时: TimeoutError")
提交后主线程未崩溃,异常类型: ValueError
调用 result() 时抛出: 处理 7 失败
超时: TimeoutError
f.exception() 返回异常对象(无异常返回 None),f.result(timeout=...) 超时抛 TimeoutError。记住:不调用 result() / exception(),任务里的异常会被静默吞掉。
12.2.3 进程池必须加 main 守卫
这是新手最容易踩的坑。macOS 与 Windows 默认用 spawn 启动子进程——子进程会重新导入你的主模块。如果 ProcessPoolExecutor 写在模块顶层,子进程导入时又会执行一遍、再创建进程,如此递归,解释器会直接报错。故意去掉守卫,实测报错如下:
from concurrent.futures import ProcessPoolExecutor
def square(x: int) -> int:
return x * x
with ProcessPoolExecutor(max_workers=2) as ex: # 顶层代码,没有守卫
print(list(ex.map(square, [1, 2, 3, 4])))
RuntimeError:
An attempt has been made to start a new process before the
current process has finished its bootstrapping phase.
This probably means that you are not using fork to start your
child processes and you have forgotten to use the proper idiom
in the main module:
if __name__ == '__main__':
freeze_support()
...
主进程侧则会看到 BrokenProcessPool。只要用进程池,就把代码放进 if __name__ == "__main__": 里,一行解决:
import os
if __name__ == "__main__":
with ProcessPoolExecutor(max_workers=4) as ex:
print(list(ex.map(square, [1, 2, 3, 4])))
print("子进程 PID != 主进程 PID:",
ex.submit(os.getpid).result(), "!=", os.getpid())
[1, 4, 9, 16]
子进程 PID != 主进程 PID: 54239 != 54211
12.2.4 启动开销与小任务反转
进程池的代价是建池 + 序列化 + 进程间通信。实测「首次任务(含建池)」:
ThreadPoolExecutor 首次任务(含建池): 0.0002s
ProcessPoolExecutor 首次任务(含建池): 0.2047s
差了一千倍。再看两类负载(池预热后重复取最小值):
200 个小任务 : 线程池 0.0009s | 进程池 0.0166s
4 个 CPU 任务: 线程池 2.166s | 进程池 0.555s
结论完全反转:任务很小时(square 这种),线程池快得多,进程池的开销全是浪费;任务真占 CPU 时,进程池拿到接近 4 倍的加速(2.166s → 0.555s),线程池则被 GIL 卡住。「用进程」不是万能药,要先看任务粒度。
12.2.5 chunksize:减少进程间往返
进程池 map 的 chunksize 把任务分批发给子进程,一次发一批而不是一个。4000 个极小任务实测:
chunksize= 1: 0.5516s
chunksize= 100: 0.0034s
chunksize=1000: 0.0034s
chunksize=1 意味着每个任务都要一次 IPC 往返,慢了两个数量级;成批发送后开销骤降。经验值:任务越小、数量越多,chunksize 就该越大;任务本身耗时长时 chunksize=1 反而更好(负载更均衡)。默认值会按任务数自动计算,通常够用。
12.2.6 进程间不共享内存
进程有独立内存空间,普通对象改不动对方。想共享,得用 multiprocessing 的专用容器:
import multiprocessing as mp
def producer(q, n):
for i in range(n):
q.put(i * i)
def accumulate(shared, n):
for i in range(n):
shared.append(i)
if __name__ == "__main__":
q = mp.Queue()
p = mp.Process(target=producer, args=(q, 5))
p.start()
p.join()
print("Queue 收到:", [q.get() for _ in range(5)])
with mp.Manager() as manager:
shared = manager.list()
procs = [mp.Process(target=accumulate, args=(shared, 3)) for _ in range(2)]
for pr in procs:
pr.start()
for pr in procs:
pr.join()
print("Manager.list 长度:", len(shared), "内容:", list(shared))
plain = []
pr = mp.Process(target=accumulate, args=(plain, 3))
pr.start()
pr.join()
print("普通 list 长度:", len(plain))
Queue 收到: [0, 1, 4, 9, 16]
Manager.list 长度: 6 内容: [0, 1, 2, 0, 1, 2]
普通 list 长度: 0
mp.Queue:单向传值(pickle 序列化),适合生产者-消费者。mp.Manager().list()/dict():代理对象,跨进程读写,两个子进程的修改都保住了(长度 6)。- 普通
list:子进程改了主进程完全看不到(长度 0)。
代价是 Manager 会起一个服务进程、每次访问都要 IPC,比普通对象慢,能少用就少用。
12.2.7 multiprocessing.Pool 与 concurrent.futures 的关系
你可能在旧代码里见过 multiprocessing.Pool:它的 map / apply_async / imap 是更早的 API。两者都能用,但:
| 维度 | multiprocessing.Pool | concurrent.futures |
|---|---|---|
| 线程支持 | 无 | 有(ThreadPoolExecutor) |
| 统一接口 | 仅进程 | 线程 + 进程同一套 |
| 取消/异常 | 较弱 | Future 模型更清晰 |
| 推荐度 | 维护旧代码 | 新代码首选 |
新项目一律用 concurrent.futures,需要线程时无缝切换,代码风格一致。
12.2.8 os.process_cpu_count 与 os.cpu_count
决定开几个 worker 时要用 CPU 核数。3.13 新增了 os.process_cpu_count():
import os
print("os.cpu_count():", os.cpu_count())
print("os.process_cpu_count():", os.process_cpu_count())
os.cpu_count(): 10
os.process_cpu_count(): 10
区别在于:os.cpu_count() 返回系统逻辑 CPU 数;os.process_cpu_count() 返回当前进程可用的 CPU 数(在 Linux 上会遵守 sched_getaffinity 的亲和性限制,容器里尤其有用)。本机 macOS 不支持 CPU 亲和性,所以两者都返回 10;在有亲和性限制的 Linux 容器里,process_cpu_count() 会更小、更贴近实际可用核数。两者都受 PYTHON_CPU_COUNT 环境变量与 -X cpu_count 影响。
12.2.9 延伸阅读
线程、进程、异步的性能取舍与更多基准,见专题 Python 并发与性能 ;进程间通信的更多细节见 Python 的 GIL 对并发编程有哪些影响 。
小结
ThreadPoolExecutor与ProcessPoolExecutor共享submit/map/shutdown一套接口,换类名即可在线程与进程间切换。submit返回Future;map保序,as_completed抢先返回;异常存进 Future,必须调用result()/exception()才会暴露。- 进程池必须放进
if __name__ == "__main__":,否则 macOS/Windows 的 spawn 会报RuntimeError。 - 进程池启动成本高:小任务下比线程池慢得多,CPU 密集时才能拿到接近核数的加速;
chunksize用来减少 IPC 往返。 - 进程间不共享内存,用
mp.Queue/mp.Manager传递数据;新代码优先concurrent.futures而非旧的multiprocessing.Pool。 os.process_cpu_count()(3.13+)返回进程可用核数,比os.cpu_count()更贴近容器实际。
到这里,我们有了「怎么把活分给多个线程/进程」的工具。下一节转向 I/O 的另一端——网络:用 socket 和 HTTP 客户端把数据从远端取回来。
阅读导航:上一节:GIL 与线程/进程模型 · 下一节:socket、HTTP 客户端与 requests/httpx 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。