多进程无锁修改共享numpy数组效率低,如何高效提速?
问题描述
在多进程场景中使用pool.map共享大型numpy数组,每个进程仅修改数组的不同行,理论上无需锁,因此选用multiprocessing.sharedctypes.RawArray而非带锁的multiprocessing.Array以避免性能损耗。但测试发现:
- 单进程写入耗时0.05s
- 使用32核多进程写入耗时1.8s,且
RawArray与Array耗时完全相同,进程写入时似乎仍有全局锁定的性能损耗。
测试代码如下:
import time import numpy as np import multiprocessing from multiprocessing.sharedctypes import RawArray def modify_array(i): arr = np.frombuffer(np_x, dtype=np.float32).reshape(np_x_shape) start_time_local = time.perf_counter() # Do some processing, each process write to a different row for j in range(100): arr[i, ...] = i print(f"Writting time {time.perf_counter() - start_time_local}") def pool_initializer(X, X_shape): global np_x np_x = X global np_x_shape np_x_shape = X_shape if __name__ == "__main__": n_processes = multiprocessing.cpu_count() # 1 # Original numpy array array_shape = (80, 1920 * 1080 * 3) data = np.ones(array_shape, dtype=np.float32) # Allocate the shared Array X = RawArray('i', np.array(array_shape).prod().item()) X_np = np.frombuffer(X, dtype=np.float32).reshape(array_shape) # Copy data to the shared array np.copyto(X_np, data) # Create the processes with multiprocessing.Pool(processes=n_processes, initializer=pool_initializer, initargs=(X, array_shape)) as pool: pool.map(modify_array, range(array_shape[0]))
问题原因
- 伪共享(False Sharing):数组是连续内存存储,不同行的元素可能落在同一个CPU缓存行中。当多个进程修改不同行时,会频繁触发缓存行失效与同步,这是性能瓶颈的核心原因,与锁无关。
- 内存带宽饱和:32个进程同时写入大数组,内存带宽被占满,导致进程等待内存资源,反而比单进程效率更低。
- 冗余内存操作:测试代码中
for j in range(100)重复赋值同一行,放大了不必要的内存写入开销。 - 类型不匹配:
RawArray使用了'i'(int类型),但实际存储的是float32,会导致数据异常,同时可能引发额外的内存处理开销。
优化方案
1. 修复RawArray类型错误
将RawArray('i', ...)改为RawArray('f', ...),匹配float32类型(ctype中'f'对应单精度浮点数)。
2. 消除伪共享
给数组的每行添加填充字节,让每行独占完整的CPU缓存行(通常64字节),避免不同行的元素共享缓存行。
3. 优化任务分配与内存访问
- 去掉冗余循环:直接一次性赋值,减少内存写入次数。
- 让每个进程处理连续的多行,而非单行,减少进程调度开销,提升缓存命中率。
4. 限制进程数量
内存带宽有限,过多进程会导致带宽饱和。建议使用cpu_count()//2或更小的进程数,测试找到最优值。
修改后的示例代码
import time import numpy as np import multiprocessing from multiprocessing.sharedctypes import RawArray def modify_array(start_end): start_idx, end_idx = start_end arr = np.frombuffer(np_x, dtype=np.float32).reshape(np_x_shape) start_time_local = time.perf_counter() # 一次性赋值连续多行,去掉冗余循环 arr[start_idx:end_idx, :np_x_original_cols] = np.arange(start_idx, end_idx)[:, None] print(f"处理行 {start_idx}-{end_idx}: 写入耗时 {time.perf_counter() - start_time_local:.6f}s") def pool_initializer(X, X_shape, original_cols): global np_x np_x = X global np_x_shape np_x_shape = X_shape global np_x_original_cols np_x_original_cols = original_cols if __name__ == "__main__": # 限制进程数量,避免内存带宽饱和 n_processes = multiprocessing.cpu_count() // 2 print(f"使用 {n_processes} 个进程") # 原始数组参数 original_rows = 80 original_cols = 1920 * 1080 * 3 original_shape = (original_rows, original_cols) # 计算缓存行对齐的填充列数(缓存行默认64字节,float32占4字节) cache_line_size = 64 elem_size = np.dtype(np.float32).itemsize elem_per_cache_line = cache_line_size // elem_size pad_cols = (elem_per_cache_line - (original_cols % elem_per_cache_line)) % elem_per_cache_line padded_cols = original_cols + pad_cols padded_shape = (original_rows, padded_cols) # 分配共享内存,使用正确的float32类型 total_elements = np.prod(padded_shape) X = RawArray('f', total_elements) X_np = np.frombuffer(X, dtype=np.float32).reshape(padded_shape) # 初始化共享数组的有效区域 data = np.ones(original_shape, dtype=np.float32) X_np[:, :original_cols] = data # 分割任务:每个进程处理连续的多行 row_chunks = [] chunk_size = original_rows // n_processes for i in range(n_processes): start = i * chunk_size end = start + chunk_size if i != n_processes - 1 else original_rows row_chunks.append((start, end)) # 启动进程池 with multiprocessing.Pool(processes=n_processes, initializer=pool_initializer, initargs=(X, padded_shape, original_cols)) as pool: pool.map(modify_array, row_chunks) # 提取处理后的有效数据(去掉填充部分) final_data = X_np[:, :original_cols].copy()
说明
优化后,缓存同步开销大幅降低,内存访问效率提升,多进程性能会显著优于单进程。需要注意的是,即使使用RawArray无锁共享,CPU的缓存一致性机制仍会带来一定开销,但这是硬件层面的必要代价,远低于锁的开销。
内容的提问来源于stack exchange,提问作者Neabfi
相关产品推荐
相关产品推荐

