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

Python多进程共享操作NumPy数组并行计算结果异常问题

问题背景

需要使用multiprocessing加速基于NumPy数组的影像分段计算任务,整体工作流如下:

  • 现有三个输入数组:
    • id_array:存储同属一个分段的像素ID
    • class_array:分类结果数组,存储影像分类得到的整数类别值
    • prob_array:概率数组,存储对应类别的置信概率
  • 针对每个分段需要完成两项计算:
    • 统计分段内像素的众数类别
    • 仅筛选分段内属于众数类别的像素,计算其概率中位数,将结果赋值给分段内所有像素
可正常运行的单进程实现
import numpy as np

id_array = np.array([[1, 1, 2, 2, 2],
                     [1, 1, 2, 2, 4],
                     [3, 3, 4, 4, 4],
                     [3, 3, 4, 4, 4]])
class_array = np.array([[7, 7, 6, 8, 8],
                        [5, 7, 7, 8, 8],
                        [8, 8, 5, 5, 8],
                        [9, 9, 8, 7, 7]])
prob_array = np.array([[0.7, 0.3, 0.9, 0.5, 0.1],
                       [0.4, 0.6, 0.3, 0.5, 0.9],
                       [0.8, 0.6, 0.2, 0.2, 0.3],
                       [0.4, 0.4, 0.6, 0.3, 0.7]])

all_ids = np.unique(id_array)
dst_classes = np.zeros_like(class_array)
dst_probs = np.zeros_like(prob_array)

for my_id in all_ids:
    segment = np.where(id_array == my_id)
    class_data = class_array[segment]
    # 统计分段内的众数类别
    majority = np.bincount(class_data.flatten()).argmax()
    # 提取分段内概率值
    prob_data = prob_array[segment]
    # 提取分段内众数类别对应的概率值
    majority_probs = prob_data[np.where(class_data == majority)]
    # 计算概率中位数
    median_prob = np.nanmedian(majority_probs)
    # 写入结果
    dst_classes[segment] = majority
    dst_probs[segment] = median_prob

print(dst_classes)
print(dst_probs)

单进程代码在小数据量下运行正常,但真实业务数据包含约400万个分段,单进程计算需耗时约一周,因此编写了基于共享内存的多进程版本:

import numpy as np
import multiprocessing as mp

WORKER_DICT = dict()
NODATA = 0

def shared_array_from_np_array(data_array, init_value=None):
    raw_array = mp.RawArray(np.ctypeslib.as_ctypes_type(data_array.dtype), data_array.size)
    shared_array = np.frombuffer(raw_array, dtype=data_array.dtype).reshape(data_array.shape)
    if init_value:
        np.copyto(shared_array, np.full_like(data_array, init_value))
        return raw_array, shared_array
    else:
        np.copyto(shared_array, data_array)
        return raw_array, shared_array

def init_worker(id_array, class_array, prob_array, class_results, prob_results):
    WORKER_DICT['id_array'] = id_array
    WORKER_DICT['class_array'] = class_array
    WORKER_DICT['prob_array'] = prob_array
    WORKER_DICT['class_results'] = class_results
    WORKER_DICT['prob_results'] = prob_results
    WORKER_DICT['shape'] = id_array.shape
    mp.freeze_support()

def worker(id):
    id_array = WORKER_DICT['id_array']
    class_array = WORKER_DICT['class_array']
    prob_array = WORKER_DICT['prob_array']
    class_result = WORKER_DICT['class_results']
    prob_result = WORKER_DICT['prob_results']
    # 提取对应ID的分段索引
    segment = np.where(id_array == id)
    # 提取对应位置数据,掩膜无值区域
    class_data = np.ma.masked_equal(class_array[segment], NODATA)
    # 计算众数类别
    majority_class = np.bincount(class_data.flatten()).argmax()
    # 提取概率值
    probs = prob_array[segment]
    majority_probs = probs[np.where(class_array[segment] == majority_class)]
    med_majority_probs = np.nanmedian(majority_probs)
    class_result[segment] = majority_class
    prob_result[segment] = med_majority_probs
    return

if __name__ == '__main__':
    # 分段ID数组
    id_ra, id_array = shared_array_from_np_array(np.array(
        [[1, 1, 2, 2, 2],
         [1, 1, 2, 2, 4],
         [3, 3, 4, 4, 4],
         [3, 3, 4, 4, 4]]))
    # 分类结果数组
    cl_ra, class_array = shared_array_from_np_array(np.array(
        [[7, 7, 6, 8, 8],
         [5, 7, 7, 8, 8],
         [8, 8, 5, 5, 8],
         [9, 9, 8, 7, 7]]))
    # 概率数组
    pr_ra, prob_array = shared_array_from_np_array(np.array(
        [[0.7, 0.3, 0.9, 0.5, 0.1],
         [0.4, 0.6, 0.3, 0.5, 0.9],
         [0.8, 0.6, 0.2, 0.2, 0.3],
         [0.4, 0.4, 0.6, 0.3, 0.7]]))
    cl_res, class_results = shared_array_from_np_array(class_array, 0)
    pr_res, prob_results = shared_array_from_np_array(prob_array, 0.)
    unique_ids = np.unique(id_array)
    init_args = (id_ra, cl_ra, pr_ra, cl_res, pr_res, id_array.shape)
    with mp.Pool(processes=2, initializer=init_worker, initargs=init_args) as pool:
        pool.map_async(worker, unique_ids)
    print('Majorities:', cl_res)
    print('Probabilities:', pr_res)
遇到的问题

多进程版本运行后无法正确读取计算结果,尝试使用如下代码读取结果:

np.frombuffer(cl_res)
np.frombuffer(pr_res)

读取结果存在异常:

  • cl_res仅返回10个值(预期应为20个)且数值完全随机
  • pr_res的值与原始输入prob_array完全一致

问题原因与修复方案

代码存在3个核心错误:

  1. 异步任务未等待执行完成:map_async是非阻塞调用,进程池上下文退出时只会调用close()不会等待异步任务跑完,直接读结果时任务还没执行,所以结果全是初始值。要么改用阻塞的pool.map,要么在map_async后调用返回对象的.get()方法等待执行完成。
  2. 初始化参数不匹配:init_worker定义只接收5个参数,但initargs传了6个值(多传了shape参数),子进程初始化直接报错,且传入的是RawArray原始对象而非包装好的NumPy视图,子进程无法正确操作数组。需要在init_worker内将RawArray按对应dtype和shape重新包装为NumPy数组。
  3. 结果读取方式错误:直接对RawArray调用np.frombuffer未指定dtype,也没有reshape回原数组形状,默认按float类型解析会得到乱码值和错误长度。

修复后的可运行代码

import numpy as np
import multiprocessing as mp

WORKER_DICT = dict()
NODATA = 0

def shared_array_from_np_array(data_array, init_value=None):
    raw_array = mp.RawArray(np.ctypeslib.as_ctypes_type(data_array.dtype), data_array.size)
    shared_array = np.frombuffer(raw_array, dtype=data_array.dtype).reshape(data_array.shape)
    if init_value is not None:
        np.copyto(shared_array, np.full_like(data_array, init_value))
    else:
        np.copyto(shared_array, data_array)
    return raw_array, data_array.dtype, data_array.shape

def init_worker(id_raw, id_dtype, id_shape, 
                class_raw, class_dtype, class_shape,
                prob_raw, prob_dtype, prob_shape,
                res_cl_raw, res_cl_dtype, res_cl_shape,
                res_pr_raw, res_pr_dtype, res_pr_shape):
    WORKER_DICT['id_array'] = np.frombuffer(id_raw, dtype=id_dtype).reshape(id_shape)
    WORKER_DICT['class_array'] = np.frombuffer(class_raw, dtype=class_dtype).reshape(class_shape)
    WORKER_DICT['prob_array'] = np.frombuffer(prob_raw, dtype=prob_dtype).reshape(prob_shape)
    WORKER_DICT['class_results'] = np.frombuffer(res_cl_raw, dtype=res_cl_dtype).reshape(res_cl_shape)
    WORKER_DICT['prob_results'] = np.frombuffer(res_pr_raw, dtype=res_pr_dtype).reshape(res_pr_shape)

def worker(seg_id):
    id_array = WORKER_DICT['id_array']
    class_array = WORKER_DICT['class_array']
    prob_array = WORKER_DICT['prob_array']
    class_result = WORKER_DICT['class_results']
    prob_result = WORKER_DICT['prob_results']
    
    segment = np.where(id_array == seg_id)
    class_data = class_array[segment]
    if len(class_data) == 0:
        return
    majority_class = np.bincount(class_data.flatten()).argmax()
    probs = prob_array[segment]
    majority_probs = probs[class_data == majority_class]
    med_prob = np.nanmedian(majority_probs)
    class_result[segment] = majority_class
    prob_result[segment] = med_prob

if __name__ == '__main__':
    mp.freeze_support()
    test_id = np.array([[1, 1, 2, 2, 2],
                        [1, 1, 2, 2, 4],
                        [3, 3, 4, 4, 4],
                        [3, 3, 4, 4, 4]])
    test_class = np.array([[7, 7, 6, 8, 8],
                           [5, 7, 7, 8, 8],
                           [8, 8, 5, 5, 8],
                           [9, 9, 8, 7, 7]])
    test_prob = np.array([[0.7, 0.3, 0.9, 0.5, 0.1],
                          [0.4, 0.6, 0.3, 0.5, 0.9],
                          [0.8, 0.6, 0.2, 0.2, 0.3],
                          [0.4, 0.4, 0.6, 0.3, 0.7]])
    
    id_ra, id_dtype, id_shape = shared_array_from_np_array(test_id)
    cl_ra, cl_dtype, cl_shape = shared_array_from_np_array(test_class)
    pr_ra, pr_dtype, pr_shape = shared_array_from_np_array(test_prob)
    res_cl_ra, res_cl_dtype, res_cl_shape = shared_array_from_np_array(test_class, init_value=0)
    res_pr_ra, res_pr_dtype, res_pr_shape = shared_array_from_np_array(test_prob, init_value=0.)

    unique_ids = np.unique(test_id)
    init_args = (id_ra, id_dtype, id_shape,
                 cl_ra, cl_dtype, cl_shape,
                 pr_ra, pr_dtype, pr_shape,
                 res_cl_ra, res_cl_dtype, res_cl_shape,
                 res_pr_ra, res_pr_dtype, res_pr_shape)
    
    with mp.Pool(processes=mp.cpu_count(), initializer=init_worker, initargs=init_args) as pool:
        pool.map(worker, unique_ids, chunksize=1000)
    
    dst_classes = np.frombuffer(res_cl_ra, dtype=res_cl_dtype).reshape(res_cl_shape)
    dst_probs = np.frombuffer(res_pr_ra, dtype=res_pr_dtype).reshape(res_pr_shape)
    
    print("众数分类结果:")
    print(dst_classes)
    print("\n概率中位数结果:")
    print(dst_probs)

400万分段场景性能优化建议

  • 不要循环对每个分段调用np.where(id_array == id),这种全数组扫描时间复杂度为O(N*M),是单进程跑一周的核心原因。先把三个数组展平为一维,按id排序后用np.split一次性切分所有分段,时间复杂度可降到O(NlogN),速度能提升数个量级。
  • chunksize不要用默认值1,400万任务下会产生巨量IPC开销,建议设置为1000~10000,根据CPU核心数调整。
  • 内存足够的前提下,不要让子进程直接写共享内存,让worker返回每个分段的计算结果(id、众数、中位数),主进程统一写入数组,可避免多进程写共享内存的锁开销。

内容的提问来源于stack exchange,提问作者s6hebern

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 14:42:21