NumPy 多核并行编程实战:用多进程、多线程与第三方库榨干 CPU 性能
2026/9/21 0:45:23 网站建设 项目流程
  • 科学计算
  • 数据分析

【免费下载链接】numpy

The fundamental package for scientific computing with Python.

项目地址:https://gitcode.com/gh_mirrors/nu/numpy
点击查看免费下载

导读

NumPy 通过向量化操作在 Python 中实现了高性能数值计算,但向量化本身并不能自动利用多核处理器的全部算力。本文以 NumPy 官方用户指南《Writing Performant NumPy Code with Multi-Core CPUs》为主体,系统讲解利用 Python 标准库(concurrent.futures的多进程/多线程执行器)与第三方库(Dask、joblib、threadpoolctl)实现多核并行的方法,并以 Mandelbrot 集合生成为贯穿全文的实战案例。读完本文,你将掌握多进程/多线程的选型依据、规避 GIL 与 CPU 过载(oversubscription)陷阱的实操技巧,以及一套可直接复制运行的并行计算模板。

引言:为什么向量化还不够

NumPy 的设计目标是通过向量化操作在 Python 中实现高性能数值计算——将循环下推到 C 层执行,避免 Python 解释器的逐元素开销。然而,向量化并不总是能充分利用多核处理器的能力:单次 NumPy 调用通常运行在单核上,若要发挥多核优势,还需要额外的并行策略。

本文围绕三条主线展开:

  • 在 Python 中使用多核处理器的一般概念(多进程与多线程的取舍);
  • 使用 Python 标准库借助 NumPy 利用多核处理器(附完整示例代码);
  • 面向多核处理的第三方库(Dask、joblib、threadpoolctl)。

在 Python 中使用多核处理器的一般概念

多进程(Multiprocessing)

多进程技术允许同时执行多个进程,每个进程拥有独立的 Python 解释器和独立的内存空间。Python 提供的高层 API 是concurrent.futures.ProcessPoolExecutor

优点:

  • 绕过全局解释器锁(GIL),实现真正的并行执行;
  • 各进程内存空间相互隔离,避免意外共享数据。

缺点:

  • 每个进程都需要独立的内存空间,内存占用更高;
  • 进程间共享数据困难,对象需要经过序列化(pickling)才能传递。
通用建议一:降低进程创建开销

与线程相比,进程创建需要初始化新的 Python 解释器和内存空间,开销更高。缓解策略:

  • 使用进程池复用已有进程,而不是为每个任务新建进程。concurrent.futures.ProcessPoolExecutor正是为此设计的;
  • 谨慎选择启动方式(start method)。除非确认应用安全,否则避免显式选择fork——从多线程进程中 fork 会引发死锁或崩溃。Python 3.14 已将 POSIX 平台的默认启动方式从fork改为forkserver,正是为了规避多线程进程的常见不兼容问题(详见 Python 官方multiprocessing文档的 "Contexts and start methods" 一节)。
通用建议二:降低通信开销

进程间通信(IPC)因数据序列化与传输会产生显著开销;在 Python 中,只有可 pickle 的对象才能跨进程传递。因此,需要频繁序列化数据的程序并不适合多进程。

  • 尽量减少进程间传输的数据量;
  • 对于需要被多个进程访问的大数据,使用共享内存结构,如multiprocessing.shared_memorymultiprocessing.Arraymultiprocessing.Value
  • 通过合理的负载均衡(见下文)确保所有进程都被高效利用,避免空闲。
通用建议三:Pickling 考量

使用多进程时,worker 函数及其参数必须可 pickle。这在处理复杂数据结构或动态定义的函数时可能成为限制。如遇 pickling 问题:

  • 重构代码以使用更简单的数据结构或函数——例如将 worker 函数定义在模块顶层,避免使用 lambda 或嵌套函数;
  • 考虑第三方库joblib:其默认后端loky依赖cloudpickle做序列化,能处理比标准pickle模块更广的 Python 对象(如 lambda 函数)。

多线程(Multithreading)

多线程允许多个线程在同一进程内运行,共享同一内存空间

需要特别说明的是:自由线程(free-threaded)Python于 Python 3.13 作为实验特性引入,在 Python 3.14 成为受支持(非实验性)特性。当与显式设计为线程安全的库结合时,线程也可以实现真正的并行执行。NumPy 官方文档在 线程安全说明 中指出:NumPy 支持通过标准库threading模块在多线程环境中使用,许多 NumPy 底层操作会释放 GIL,因此与许多纯 Python 代码不同,NumPy 能够利用多线程并行性。

Python 提供的高层 API 是concurrent.futures.ThreadPoolExecutor。此外还有concurrent.futures.InterpreterPoolExecutor(在独立线程中运行多个解释器、避免解释器间共享 Python 对象),但该执行器目前尚不可用于 NumPy(见 NumPy 议题 gh-24755)。

优点:

  • 线程共享同一内存空间,内存占用更低;
  • 线程间通信更容易。

缺点:

  • 多个线程同时修改共享数据并伴随其他线程读取时,可能产生竞态条件(race condition)
  • 若所用 Python 库非线程安全或对自由线程构建支持有限,性能提升受限。
通用建议一:避免竞态条件

竞态条件指多个线程同时更新共享数据导致结果不可预测。规避策略:

  • 通过线程局部存储(thread-local storage)或显式向线程传递数据,最小化线程间共享的数据量
  • 尽可能使用不可变 NumPy 数组或只读访问模式,减少显式同步需求;
  • 使用线程安全数据结构或同步原语(锁、信号量、条件变量)管理共享数据访问。注意:这些同步机制使用不当会造成死锁,需谨慎使用。

NumPy 的 线程安全文档 也强调:跨线程共享 NumPy 数组需要格外小心,例如一个线程在读取数组时另一线程调整数组大小会造成未定义行为;最稳妥的做法是每个线程持有自己的数组对象、共享数组只读访问,或在需要变更时自行加锁。

通用建议二:避免 CPU 过载(oversubscription)

部分 NumPy 操作——如矩阵乘法与线性代数函数(见 线性代数参考文档)——会使用底层 BLAS 库(如 OpenBLAS、MKL)提供的多线程。若这些操作运行在已经占满所有 CPU 核的外层线程池中,就会发生 CPU 过载:外层线程池与 BLAS 线程竞争同一批 CPU 资源,反而降低性能。

NumPy 官方 全局配置选项文档 印证了这一背景:NumPy 自身的函数调用通常刻意限制为单线程,但高性能线性代数依赖 OpenBLAS、MKL 等 BLAS 后端,这些后端可能使用多线程,线程数可通过OMP_NUM_THREADS等环境变量或threadpoolctl包控制。numpy.linalg 包文档 也明确写道:BLAS/LAPACK 库是多线程且依赖处理器的,可能需要环境变量或threadpoolctl等外部包来控制线程数。

规避策略:

  • 例如使用threadpoolctl将 BLAS 线程数限制为 1。

多进程与多线程的通用建议

负载均衡(Balance processing load)

若处理负载在 worker 间分配不均,部分 worker 提前完成而空闲、其余仍在工作,会导致资源利用率低下、整体执行时间变长。策略:

  • 采用动态任务分配:任务在 worker 空闲时就绪时再分配,而非预先静态分配;
  • 检查chunksize参数:任务块既不能太小(造成过多开销),也不能太大(导致负载失衡)。
确定正确的 CPU 数量

Python 提供os.cpu_count(系统 CPU 数)和os.process_cpu_count(当前进程可用 CPU 数)两个函数。但在某些环境(如 Docker 容器、HPC 集群)中,二者可能无法反映进程实际可用的 CPU 数。

更准确的做法是使用joblib.cpu_count():它会考虑 CPU 亲和性设置、Linux CFS 调度器配额等约束,返回值更贴近真实可用核数(详见下文 joblib 一节)。

用 Python 标准库利用多核:Mandelbrot 集合实战

本节演示如何用 Python 标准库结合 NumPy 利用多核处理器。示例选用Mandelbrot 集合生成

Mandelbrot 集合定义为满足如下条件的复数集合:由迭代函数

z_{n+1} = z_n^2 + c, z_0 = 0

产生的序列不发散到无穷大。若经过固定迭代次数后z_n的绝对值仍有界(通常以不超过阈值2判断),则认为c属于 Mandelbrot 集合。

按照该定义,复平面上每个点可独立计算,天然适合并行计算。下图的热色代表复平面上每个点序列发散所用的迭代次数。

多进程示例

以下代码演示如何使用concurrent.futures.ProcessPoolExecutor跨多个进程并行生成 Mandelbrot 集合。

该示例优先考虑清晰性而非效率。实践中,在进程间传输大型 NumPy 数组代价高昂;定义共享内存数组或在每个进程内创建数组可能是更高效的实现。

from concurrent.futures import ProcessPoolExecutor import numpy as np from numpy.typing import NDArray def mandelbrot_block( c_block: NDArray[np.complex128], max_iter: int ) -> NDArray[np.int64]: z = np.zeros(c_block.shape, dtype=np.complex128) steps = np.zeros(c_block.shape, dtype=np.int64) for _ in range(max_iter): mask = np.abs(z) <= 2 z[mask] = z[mask] * z[mask] + c_block[mask] steps[mask] += 1 return steps def mandelbrot_set( arr: NDArray[np.complex128], max_iter: int, n_workers: int, ) -> NDArray[np.int64]: n_workers = min(n_workers, arr.size) arrs = np.array_split(arr, n_workers) with ProcessPoolExecutor(max_workers=n_workers) as pool: futures = [ pool.submit(mandelbrot_block, _arr, max_iter) for _arr in arrs ] results = [future.result() for future in futures] return np.concatenate(results) if __name__ == '__main__': xmin, xmax, ymin, ymax = -2.0, 1.0, -1.5, 1.5 nx, ny = 800, 800 max_iter = 10000 n_workers = 10 real = np.linspace(xmin, xmax, nx, dtype=np.float64) imag = np.linspace(ymin, ymax, ny, dtype=np.float64) arr = (real[:, np.newaxis] + 1j * imag[np.newaxis, :]).ravel() mandelbrot_image = mandelbrot_set(arr, max_iter, n_workers) mandelbrot_image = mandelbrot_image.reshape((nx, ny))

实现要点解读:

  • 任务切分np.array_split(arr, n_workers)将整个复数平面按 worker 数均匀切块,每块交给一个进程独立计算——这正是"每个点可独立计算"这一并行友好特性的落地;
  • 进程池管理ProcessPoolExecutor作为上下文管理器使用,退出时自动回收进程,避免反复创建进程的开销(呼应前文"降低创建开销");
  • 结果合并:各进程返回的np.int64步数数组通过np.concatenate拼回,再reshape((nx, ny))还原为图像形状;
  • 入口保护if __name__ == '__main__':是多进程编程的必备结构——Windows 与spawn/forkserver启动方式下,子进程会重新导入主模块,若无此保护会引发递归创建进程的严重问题;
  • 数据传递成本:每个块数组都要经历 pickle 序列化→跨进程传输→反序列化的完整链路,这正是文档强调"清晰性优先于效率"的原因。

多线程示例

与多进程示例类似,下面用concurrent.futures.ThreadPoolExecutor跨多个线程并行生成 Mandelbrot 集合。

环境准备:安装自由线程版 Python

运行多线程示例前,需要确保已安装Python 3.13 或更高版本的自由线程(free-threaded)构建。根据 Python 官方文档,可通过以下方式验证当前 Python 构建是否为自由线程版:

  • 在终端运行python -VV,检查输出中是否显示free-threading build
  • 在 Python shell 中检查sys._is_gil_enabled()的值,应为False
代码示例
import sys from concurrent.futures import ThreadPoolExecutor import numpy as np def mandelbrot_block(start: int, stop: int, max_iter: int) -> None: z_target = np.zeros(stop - start, dtype=np.complex128) indexes = slice(start, stop) arr_target = SHARED_readonly_arr[indexes] steps_target = SHARED_updating_steps[indexes] threshold = 2.0 for _ in range(max_iter): mask = np.abs(z_target) <= threshold z_target[mask] = z_target[mask] * z_target[mask] + arr_target[mask] steps_target[mask] += 1 SHARED_updating_steps[indexes] = steps_target return None def mandelbrot_set( total_size: int, max_iter: int, n_workers: int, ) -> None: chunksize = total_size // n_workers with ThreadPoolExecutor(max_workers=n_workers) as pool: futures = [ pool.submit( mandelbrot_block, start, min(start + chunksize, total_size), max_iter ) for start in range(0, total_size, chunksize) ] _ = [future.result() for future in futures] if __name__ == '__main__': print("Python version is free-threaded:", not sys._is_gil_enabled()) assert not sys._is_gil_enabled() xmin, xmax, ymin, ymax = -2.0, 1.0, -1.5, 1.5 nx, ny = 800, 800 max_iter = 10000 n_workers = 10 real = np.linspace(xmin, xmax, nx, dtype=np.float64) imag = np.linspace(ymin, ymax, ny, dtype=np.float64) SHARED_readonly_arr = (real[:, np.newaxis] + 1j * imag[np.newaxis, :]).ravel() SHARED_readonly_arr.flags.writeable = False SHARED_updating_steps = np.zeros(SHARED_readonly_arr.shape, dtype=np.int64) mandelbrot_set(SHARED_readonly_arr.size, max_iter, n_workers) mandelbrot_image = SHARED_updating_steps.reshape((nx, ny))

实现要点解读:

  • GIL 检查:入口处assert not sys._is_gil_enabled()强制要求在自由线程构建下运行——若在标准 GIL 构建下执行,线程间的数值循环无法真正并行,示例失去意义;
  • 共享数组设计:本实现刻意在多个线程间共享数组。SHARED_readonly_arr是保存待求值复数集合的只读数组(通过flags.writeable = False显式置为只读),SHARED_updating_steps是保存各点迭代次数的更新数组——只读共享 + 分块写入(每线程只写属于自己的切片indexes)的设计,正是前文"避免竞态条件"建议(不可变数组/只读访问 + 最小化共享写入)的代码级体现;
  • 任务切分方式:与多进程版用np.array_split切分数据不同,多线程版按索引区间[start, stop)切分任务,每个线程只读写自己负责的切片,互不重叠,天然避免写入冲突;
  • 无返回值设计:结果直接写入共享数组SHARED_updating_steps,因此 worker 返回None,线程间零序列化开销——这也是多线程相较多进程最大的成本优势(共享内存、无 pickling)。

面向多核处理的第三方库

在许多实际场景中,第三方库比 Python 标准库提供更便捷、更高效的并行方案。

Dask

Dask 是开源并行计算库,不仅支持单机,还支持机器集群并行计算。它提供与 NumPyndarrayAPI 高度相似的DaskArray:如果你熟悉 NumPy,可以轻松上手DaskArray,将原有的数组式写法迁移到更大规模、分布式的并行计算中。

joblib

joblib提供了一系列便于并行化任务的辅助函数。其两个核心能力:

  • 默认后端loky依赖cloudpickle进行序列化,能处理比标准pickle模块更广的 Python 对象(例如lambda 函数),有效缓解前文提到的 pickling 限制;
  • joblib.cpu_count()返回当前进程可用的 CPU 数,会考虑 CPU 亲和性设置与 Linux CFS 调度器配额等约束。在 Docker 容器等资源受限环境中,该值通常比os.cpu_count/os.process_cpu_count更准确,可用于确定进程池/线程池的max_workers规模。

threadpoolctl

threadpoolctl提供控制 Python 中线程池行为的工具,包括 BLAS、OpenMP 等库使用的底层线程池。它允许你在同时使用多个会动用线程的库时,避免前文所述的CPU 过载(oversubscription)问题——例如在外部线程池中运行numpy.linalg系列函数(numpy.linalg模块)时,先用threadpoolctl将 BLAS 线程数限制为 1,即可防止外层线程池与 BLAS 内部线程争抢 CPU 核。NumPy 官方的 全局配置选项文档 也推荐用threadpoolctl控制线性代数后端线程数,并提及OMP_NUM_THREADS等环境变量在 OpenBLAS/MKL 场景下的等效作用。

小结与选型参考

维度多进程(ProcessPoolExecutor)多线程(ThreadPoolExecutor)
并行本质多解释器、多内存空间,真并行单进程多线程,需自由线程构建(Python 3.13+,3.14 正式支持)才可绕过 GIL
内存开销高(每进程独立内存)低(共享内存)
数据通信需 pickle 序列化,代价高共享数组直接读写,零序列化
主要风险pickling 限制、进程创建开销竞态条件、CPU 过载(BLAS 线程竞争)
适用场景数据可切分、传输量小、函数可 pickle大量共享只读数据、需要低延迟通信

写作高性能多核 NumPy 代码时,可以遵循如下决策路径:

  1. 先向量化:确保单核上的 NumPy 操作已用足向量化能力;
  2. 判断并行瓶颈:任务是 CPU 密集(选多进程或自由线程多线程)还是 I/O 密集(线程即可);
  3. 确定 CPU 规模:容器/HPC 环境下用joblib.cpu_count()而非os.cpu_count
  4. 防止线程叠加:并行任务中调用 BLAS 密集运算时,用threadpoolctl将 BLAS 线程限制为 1,避免 CPU 过载;
  5. 控制任务粒度:调好chunksize,既不过细(调度开销)也不过粗(负载失衡),并优先采用动态任务分配。

上述标准库与第三方库两条技术路线在 用户指南目录 中有系统编排,配合 NumPy 线程安全说明 与 全局配置选项文档 阅读,可以形成从概念到落地的完整多核优化知识闭环。

  • 科学计算
  • 数据分析

【免费下载链接】numpy

The fundamental package for scientific computing with Python.

项目地址:https://gitcode.com/gh_mirrors/nu/numpy
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询