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

多进程无锁修改共享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]))
问题原因
  1. 伪共享(False Sharing):数组是连续内存存储,不同行的元素可能落在同一个CPU缓存行中。当多个进程修改不同行时,会频繁触发缓存行失效与同步,这是性能瓶颈的核心原因,与锁无关。
  2. 内存带宽饱和:32个进程同时写入大数组,内存带宽被占满,导致进程等待内存资源,反而比单进程效率更低。
  3. 冗余内存操作:测试代码中for j in range(100)重复赋值同一行,放大了不必要的内存写入开销。
  4. 类型不匹配: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 22:12:11