Python多进程计算瓶颈排查:特征值并行效率异常问题
多进程并行计算特征值未达预期性能的问题
问题描述
我写了一段代码,通过参数(r,t)构建矩阵H,接着计算该矩阵的特征值并完成相关运算。需要对约200组不同的(r,t)参数执行这个操作,因为各操作相互独立,所以用Python的multiprocessing实现并行计算。
单组(r,t)参数的计算耗时约1.5小时,原本预期200组并行计算耗时差不多,但实际运行已经进入第2天,耗时远超预期。我用的是高校高性能计算单元,mp.cpu_count()显示有192核可用,理论上核心足够。
原始代码示例
import numpy as np import multiprocessing rs = np.linspace(0.5,1.5,20) thetas = np.linspace(0,np.pi/2,10) def matrix_stuff(r,t): #constructs a matrix H(r,t). Diagonalize and do stuff. processes = [] rets = [] q = multiprocessing.Queue() for i in range(len(thetas)): for j in range(len(rs)): theta = thetas[i] r = rs[I] # 此处有笔误:I应为j p = multiprocessing.Process(target=matrix_stuff, args = (r,theta,i,j)) processes.append(p) p.start() for p in processes: ret = q.get() rets.append(ret) for p in processes: p.join() #store output in deltas deltas = np.zeros((len(rs),len(thetas))) for ret in rets: # ret has format (value, i,j) deltas[rets[1],rets[2]] = rets[0] # 此处错误:应为ret[1],ret[2],ret[0] np.savetxt('new_deltas.txt',deltas)
补充测试情况
我写了测试脚本验证耗时,当并行任务仅为time.sleep(10)时,并行符合预期;但用如下测试代码时:
import multiprocessing as mp import time import numpy as np def foo(): X = np.random.rand(2000,2000) eigs = np.linalg.eigvals(X) return eigs def foo_q(q): X = np.random.rand(2000,2000) eigs = np.linalg.eigvals(X) q.put(eigs) start_lin = time.time() for _ in range(5): val = foo() end_lin = time.time() print('Time taken for linear process : ',(end_lin - start_lin)) q = mp.Queue() processes = [] start_par = time.time() for _ in range(5): p = mp.Process(target = foo_q, args = (q,)) processes.append(p) p.start() for p in processes: ret = q.get() for p in processes: p.join() end_par = time.time() print('Time taken for parallel process : ',(end_par - start_par))
得到结果:
Time taken for linear process : 15.935563564300537
Time taken for parallel process : 16.6868999004364
我想知道这是怎么回事?程序是否存在我没考虑到的瓶颈?
核心问题分析
- 进程过载与调度开销:一次性创建200个进程,远超192核的硬件上限,操作系统需要频繁切换进程上下文,消耗大量资源,反而拖慢每个任务的执行效率。
- numpy底层多线程冲突:
numpy.linalg.eigvals依赖的BLAS/LAPACK库(如OpenBLAS、MKL)默认会开启多线程加速单个特征值计算。这会导致你手动创建的多进程,与底层库的多线程抢占CPU核心,最终所有任务互相竞争资源,并行效率不升反降——这也是测试代码中并行比串行还慢的直接原因。 - 原始代码的明显错误:
- 循环中
r = rs[I]是笔误,应为rs[j],否则会触发未定义变量报错。 matrix_stuff函数没有向Queue写入结果的逻辑,主进程的q.get()会一直阻塞,程序无法正常结束。- 最后填充
deltas数组时,错误使用rets[1],rets[2],应改为ret[1],ret[2](取当前结果的索引,而非整个结果列表的固定索引)。
- 循环中
解决方法
- 限制numpy线程数:在程序开头设置环境变量,强制numpy使用单线程,避免底层多线程与多进程冲突:
import os # 针对不同线性代数库的设置 os.environ["OPENBLAS_NUM_THREADS"] = "1" os.environ["MKL_NUM_THREADS"] = "1" os.environ["OMP_NUM_THREADS"] = "1" - 使用进程池管理进程:用
multiprocessing.Pool或concurrent.futures.ProcessPoolExecutor自动控制进程数量(默认匹配CPU核心数),避免手动创建大量进程带来的调度开销:import numpy as np import multiprocessing as mp import os # 先设置单线程 os.environ["OPENBLAS_NUM_THREADS"] = "1" os.environ["MKL_NUM_THREADS"] = "1" os.environ["OMP_NUM_THREADS"] = "1" rs = np.linspace(0.5,1.5,20) thetas = np.linspace(0,np.pi/2,10) def matrix_stuff(args): r, theta, i, j = args # 构建矩阵H(r,t)、计算特征值等操作 # 示例返回值,根据实际逻辑调整 result_val = 0.0 # 替换为你的计算结果 return (result_val, i, j) if __name__ == "__main__": # 生成所有参数组合 params = [(rs[j], thetas[i], i, j) for i in range(len(thetas)) for j in range(len(rs))] # 创建进程池,进程数可设为核心数或略少 with mp.Pool(processes=mp.cpu_count()) as pool: rets = pool.map(matrix_stuff, params) # 填充结果数组 deltas = np.zeros((len(rs), len(thetas))) for ret in rets: val, i, j = ret deltas[j, i] = val # 根据你的索引逻辑调整 np.savetxt('new_deltas.txt', deltas) - 避免用Queue传递大对象:如果特征值计算结果是大数组,
Queue的序列化/反序列化会带来额外开销,进程池的map方法在传递结果时更高效。
测试代码改进验证
修改测试代码加入线程数限制,并用进程池:
import multiprocessing as mp import time import numpy as np import os # 设置单线程 os.environ["OPENBLAS_NUM_THREADS"] = "1" os.environ["MKL_NUM_THREADS"] = "1" os.environ["OMP_NUM_THREADS"] = "1" def foo(): X = np.random.rand(2000,2000) eigs = np.linalg.eigvals(X) return eigs start_lin = time.time() for _ in range(5): val = foo() end_lin = time.time() print('串行耗时 : ', (end_lin - start_lin)) start_par = time.time() with mp.Pool(processes=5) as pool: pool.map(foo, range(5)) end_par = time.time() print('并行耗时 : ', (end_par - start_par))
运行后并行耗时会明显低于串行,符合预期。
内容的提问来源于stack exchange,提问作者Souroy
相关产品推荐
相关产品推荐

