如何让multiprocessing Pool空闲时充分利用空余核心处理大计算量任务
你遇到的性能问题核心来自两点:
- 循环计算本身是强依赖上一步结果的迭代逻辑,无法拆分为多进程并行,大任务的性能上限完全取决于numpy矩阵运算的多线程加速效率
- 多进程执行小任务时,为了避免进程间numpy多线程资源竞争,通常会默认/手动将numpy单进程线程数限制为1,等到只剩大任务单进程运行时,这个线程限制没有被放开,导致只能用1个核心,对应你6核设备用5核时20%的使用率(1/5)
解决方案
1. 核心优化思路
将原有的一次性提交所有任务的逻辑拆为两个阶段:
- 第一阶段:并行执行所有小维度的A/B计算任务,此时每个进程内numpy线程数设为1,避免多进程资源抢占
- 第二阶段:所有小任务执行完成后,放开numpy线程数限制为所有可用核心,单独执行大维度的C/D计算任务,充分利用空闲核心加速矩阵运算
2. 调整后的代码示例
import os # 先设置线程数相关环境变量,必须在import numpy前配置 os.environ['OMP_NUM_THREADS'] = '1' os.environ['OPENBLAS_NUM_THREADS'] = '1' os.environ['MKL_NUM_THREADS'] = '1' import multiprocessing as mp import psutil import copy import numpy as np from time import time def my_func(args): low_index = args[0][0] up_index = args[0][1] params = args[1][0] A = args[1][1] B = args[1][2] print("PID:", mp.current_process()) for k in range(low_index, up_index): a = params[k] A = a*A + (np.dot(A, B))*(np.dot(B, B)) B = a*B + (np.dot(B, A))*(np.dot(A, A)) return A,B if __name__ == '__main__': ts = time() params = np.linspace(1, 10, 1000) n_dim = 1000 A = np.random.rand(n_dim, n_dim) B = np.random.rand(n_dim, n_dim) C = np.random.rand(5*n_dim, 5*n_dim) D = np.random.rand(5*n_dim, 5*n_dim) ncpus = psutil.cpu_count(logical=False) number_processes = ncpus - 1 total_items = params.shape[0] n_chunck = int(total_items/number_processes) intervals = [ [k*n_chunck, (k+1)*n_chunck] for k in range(number_processes) ] intervals[-1][-1] = total_items # 拆分小任务和大任务 small_task_args = [] for k in range(number_processes - 1): small_task_args.append( [intervals[k], (params, copy.deepcopy(A), copy.deepcopy(B))] ) big_task_args = [intervals[-1], (params, copy.deepcopy(C), copy.deepcopy(D))] # 第一阶段:并行跑小任务 pool = mp.Pool(processes = number_processes - 1) small_results = pool.map(my_func, small_task_args) pool.close() pool.join() # 第二阶段:放开numpy线程限制跑大任务 os.environ['OMP_NUM_THREADS'] = str(number_processes) os.environ['OPENBLAS_NUM_THREADS'] = str(number_processes) os.environ['MKL_NUM_THREADS'] = str(number_processes) # 重新加载numpy让配置生效,也可以用对应后端API直接调整线程数 import importlib importlib.reload(np) big_result = my_func(big_task_args) # 合并结果 results = small_results + [big_result] print(f"总耗时:{time() - ts}")
3. 可选优化
- 如果动态修改环境变量不生效,可以把大任务单独提交给一个新的进程池,进程池启动前先设置好线程数环境变量即可
- 如果你使用MKL/OpenBLAS后端,可以直接调用对应API调整线程数,不用重新加载numpy,比如
mkl.set_num_threads(number_processes)
内容的提问来源于stack exchange,提问作者Zarathustra
相关产品推荐
相关产品推荐

