You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何利用多进程共享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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 08:45:14