Python多进程共享操作NumPy数组并行计算结果异常问题
问题背景
需要使用multiprocessing加速基于NumPy数组的影像分段计算任务,整体工作流如下:
- 现有三个输入数组:
id_array:存储同属一个分段的像素IDclass_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个核心错误:
- 异步任务未等待执行完成:
map_async是非阻塞调用,进程池上下文退出时只会调用close()不会等待异步任务跑完,直接读结果时任务还没执行,所以结果全是初始值。要么改用阻塞的pool.map,要么在map_async后调用返回对象的.get()方法等待执行完成。 - 初始化参数不匹配:
init_worker定义只接收5个参数,但initargs传了6个值(多传了shape参数),子进程初始化直接报错,且传入的是RawArray原始对象而非包装好的NumPy视图,子进程无法正确操作数组。需要在init_worker内将RawArray按对应dtype和shape重新包装为NumPy数组。 - 结果读取方式错误:直接对
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
相关产品推荐
相关产品推荐

