Python中利用多进程/process_map加速Numpy密集计算的正确方式
针对Numpy密集计算的并行加速方案
问题核心分析
你的测试结果里多进程反而变慢、多线程加速有限,本质是两个原因:
- 多进程通信开销过大:
process_map默认会把每个numpy数组通过pickle序列化后传递给子进程,500×500的数组虽然单个体积不大,但频繁的序列化/反序列化开销会抵消并行计算的收益,尤其当单任务计算时间不足以覆盖通信成本时。 - 多线程受GIL限制:Python的全局解释器锁(GIL)会让CPU密集型任务的多线程无法真正并行,仅能通过numpy底层C函数释放GIL的间隙获得少量加速,无法利用多核。
正确的并行加速方案
1. 优先用Joblib实现多进程(最省心)
Joblib专门针对数值计算场景优化了并行逻辑,默认用loky后端处理numpy数组,会自动使用共享内存减少数据传递开销,无需手动处理拆分或内存共享:
import time import numpy as np from tqdm.auto import tqdm from joblib import Parallel, delayed def mydataset(size, length): for ii in range(length): yield np.random.rand(*size) def calc(mat): # 模拟密集计算,替换为你的真实逻辑 avg = np.mean(mat) std = np.std(mat) for _ in range(999): avg = np.mean(mat) std = np.std(mat) return avg, std def main(): ds = list(mydataset((500,500), 100)) # 单循环基准 t0 = time.time() res1 = [calc(mat) for mat in tqdm(ds)] print(f'for loop: {time.time() - t0:.2f}s') # Joblib多进程,n_jobs=-1自动用满所有核心 t0 = time.time() res2 = Parallel(n_jobs=-1, verbose=0)(delayed(calc)(mat) for mat in tqdm(ds)) print(f'joblib multi-process: {time.time() - t0:.2f}s') if __name__ == '__main__': main()
2. 手动用共享内存优化多进程(极致性能)
如果需要更精细的控制,可通过multiprocessing.Array创建共享内存,让子进程直接访问父进程的数组,完全避免序列化开销:
import time import numpy as np from tqdm.auto import tqdm from concurrent.futures import ProcessPoolExecutor import multiprocessing as mp def mydataset(size, length): for ii in range(length): yield np.random.rand(*size) # 将数组存入共享内存 def create_shared_arr(arr): shared_buf = mp.Array('d', arr.size, lock=False) np_arr = np.frombuffer(shared_buf, dtype=np.float64).reshape(arr.shape) np_arr[:] = arr[:] return shared_buf, arr.shape # 从共享内存读取并计算 def calc_shared(arr_info): shared_buf, shape = arr_info mat = np.frombuffer(shared_buf, dtype=np.float64).reshape(shape) avg = np.mean(mat) std = np.std(mat) for _ in range(999): avg = np.mean(mat) std = np.std(mat) return avg, std def main(): ds = list(mydataset((500,500), 100)) t0 = time.time() res1 = [calc(mat) for mat in tqdm(ds)] print(f'for loop: {time.time() - t0:.2f}s') # 准备共享内存数据 shared_data = [create_shared_arr(mat) for mat in ds] t0 = time.time() with ProcessPoolExecutor() as executor: res2 = list(tqdm(executor.map(calc_shared, shared_data), total=len(ds))) print(f'multi-process with shared memory: {time.time() - t0:.2f}s') if __name__ == '__main__': main()
3. 批量处理任务(降低通信频率)
如果单任务计算时间较短,可将多个数组打包成一批处理,减少进程间通信的次数,让计算时间占比远高于通信开销:
import time import numpy as np from tqdm.auto import tqdm from concurrent.futures import ProcessPoolExecutor def mydataset(size, length): for ii in range(length): yield np.random.rand(*size) def calc_batch(mats): results = [] for mat in mats: avg = np.mean(mat) std = np.std(mat) for _ in range(999): avg = np.mean(mat) std = np.std(mat) results.append((avg, std)) return results # 拆分数据集为批量 def split_batches(lst, batch_size): for i in range(0, len(lst), batch_size): yield lst[i:i+batch_size] def main(): ds = list(mydataset((500,500), 100)) batch_size = 10 # 根据核心数调整,28核可设为4-5 t0 = time.time() res1 = [calc(mat) for mat in tqdm(ds)] print(f'for loop: {time.time() - t0:.2f}s') t0 = time.time() batches = list(split_batches(ds, batch_size)) with ProcessPoolExecutor() as executor: batch_results = list(tqdm(executor.map(calc_batch, batches), total=len(batches))) # 合并批量结果 res2 = [item for sublist in batch_results for item in sublist] print(f'multi-process with batch: {time.time() - t0:.2f}s') if __name__ == '__main__': main()
关键注意事项
- CPU密集型任务必须用多进程,多线程仅适合IO密集场景;
- 优先优化计算逻辑:比如减少数组重复遍历(例如用
np.mean和np.var替代两次遍历计算std),避免不必要的循环; - 若用
process_map,可通过chunksize参数设置批量大小,减少通信次数,例如process_map(calc, ds, chunksize=10)。
内容的提问来源于stack exchange,提问作者LiTuX
相关产品推荐
相关产品推荐

