如何利用多进程共享Numpy数组加速网格聚合计算?
多进程加速numpy数组处理(共享内存方案)
问题背景
我有两个维度相同、静态且无需修改的数组arrayone和arraytwo,还有一个预构建的masterarray用来存储最终的整数计算结果。目前遍历ndarray行列的处理速度还过得去,但想通过多进程来加速这个过程,同时要避免内存消耗过高,让多进程在共享内存中遍历任意行(i)的列(j),并把结果写入masterarray。之前参考过Stack Overflow的相关方案,但sharedmem存在稳定性问题,所以来求助。
原参考代码:
def gridagg(): masterarray = np.empty([1228,2606,208]) for index, val in np.ndenumerate(arrayone): selection = arraytwo[index[0]][index[1]] piece = stacked[selection[:,0], selection[:,1]].tolist() piece = [j for i in piece for j in i] comparray = np.array(piece) if index[1] == 0: compiled = comparray else: stage1 = comparray stage2 = compiled if index[1] == 1: compiled = np.stack([stage2, stage1]) else: compiled = np.vstack([stage2, stage1[None, :]]) if index[1] == 2605: masterarray[index[0], :] = compiled
解决方案:用Python官方shared_memory实现稳定共享内存多进程
我们可以利用Python 3.8+内置的multiprocessing.shared_memory模块(官方稳定,无第三方依赖),结合按行拆分任务的方式实现加速——因为原代码中每行的处理是独立的(整行的compiled由该行所有列逐步构建,与其他行无关),非常适合并行化。
步骤1:重构单行处理逻辑
把原代码中针对单行的处理抽成独立函数,优化冗余逻辑:
import numpy as np from multiprocessing import Pool, shared_memory def process_row(row_idx): # 连接到主进程创建的共享内存 shm = shared_memory.SharedMemory(name=globals()['shm_name']) master = np.ndarray(globals()['master_shape'], dtype=globals()['dtype'], buffer=shm.buf) compiled = None col_count = globals()['master_shape'][1] for col_idx in range(col_count): selection = globals()['arraytwo'][row_idx][col_idx] # 用numpy flatten替代列表推导,提升效率 piece = globals()['stacked'][selection[:,0], selection[:,1]].flatten() comparray = np.array(piece, dtype=globals()['dtype']) if col_idx == 0: compiled = comparray else: if col_idx == 1: compiled = np.stack([compiled, comparray]) else: compiled = np.vstack([compiled, comparray[None, :]]) # 将结果写入共享内存的masterarray对应行 master[row_idx, :] = compiled shm.close()
步骤2:主进程初始化共享内存并启动多进程
if __name__ == "__main__": # ---------------------- # 假设以下数组已在主进程加载完成 # arrayone = ... # arraytwo = ... # stacked = ... # ---------------------- # 配置masterarray的参数(根据实际数据类型调整dtype) dtype = np.int32 master_shape = (1228, 2606, 208) # 创建共享内存,大小匹配masterarray的字节数 shm = shared_memory.SharedMemory(create=True, size=np.prod(master_shape)*dtype.itemsize) # 创建共享内存的numpy视图 masterarray = np.ndarray(master_shape, dtype=dtype, buffer=shm.buf) # 全局变量传递给子进程(Unix下fork继承,Windows下需显式传递) globals()['shm_name'] = shm.name globals()['master_shape'] = master_shape globals()['dtype'] = dtype globals()['arrayone'] = arrayone globals()['arraytwo'] = arraytwo globals()['stacked'] = stacked # 启动进程池,进程数建议等于CPU核心数 with Pool(processes=4) as pool: # 提交所有行的处理任务 pool.map(process_row, range(master_shape[0])) # 处理完成后释放共享内存 shm.close() shm.unlink()
关键注意事项
- Windows系统适配:Windows没有fork机制,静态数组(
arrayone/arraytwo/stacked)无法通过继承传递,需额外将它们也放入共享内存,或者用multiprocessing.Manager(效率略低)。 - 内存优化:如果
stacked是超大数组,必须也用shared_memory共享,否则每个子进程会复制一份,导致内存暴涨。 - 进程数设置:进程数不要超过CPU物理核心数,避免上下文切换抵消加速效果。
- 结果验证:先测试前10行的处理结果,和原单进程代码对比,确保逻辑一致后再跑全量数据。
- 效率提升:原代码中的列表推导
[j for i in piece for j in i]替换为np.flatten(),能大幅减少CPU开销。
内容的提问来源于stack exchange,提问作者TornadoEric
相关产品推荐
相关产品推荐

