如何用Python multiprocessing库并行化深度数组遍历循环?
测井数据并行计算优化方案
我在处理带深度序列的测井数据时,经常需要遍历每个深度点执行计算,再把结果存入数组生成新的测井曲线。以下是串行计算的示例代码:
import numpy as np depth = np.linspace(5000,6000,2001) y1 = np.random.random((len(depth),3)) y2 = np.random.random(len(depth)) def fun_1(y1, y2): return y1 + y2 def fun_2(y1, y2): return sum(y1 * y2) result_1 = np.zeros_like(y1) result_2 = np.zeros_like(y2) for i in range(len(depth)): prov_result_1 = fun_1(y1[i], y2[i]) prov_result_2 = fun_2(y1[i], y2[i]) result_1[i,:] = prov_result_1 result_2[i] = prov_result_2
当深度数组规模较大、计算函数复杂时,这种串行循环的耗时会急剧增加。我尝试用Python的multiprocessing库并行化这个循环,但之前找到的方案要么无法适配我的场景,要么过于复杂。我设想按处理器数量分割数据,并行计算后再拼接结果,但需要通用的实现方式,目前的手动分割示例代码如下:
num_pool = 2 depth_cut = (depth[-1]-depth[0])/num_pool depth_parallel = [depth[depth <= depth[0] + depth_cut], depth[depth > depth[0] + depth_cut]] y1_parallel = [y1[depth <= depth[0] + depth_cut], y1[depth > depth[0] + depth_cut]] y2_parallel = [y2[depth <= depth[0] + depth_cut], y2[depth > depth[0] + depth_cut]]
通用并行化实现方案
核心思路
放弃手动按深度分割的方式,改用按索引均分数据的通用分片逻辑,结合multiprocessing.Pool实现进程级并行,最后合并各分片的计算结果。
完整实现代码
import numpy as np from multiprocessing import Pool, cpu_count # 生成测试数据(和原示例一致) depth = np.linspace(5000, 6000, 2001) y1 = np.random.random((len(depth), 3)) y2 = np.random.random(len(depth)) # 调整计算函数为支持数组分片的版本 def fun_1(y1_slice, y2_slice): return y1_slice + y2_slice def fun_2(y1_slice, y2_slice): # 用np.sum替代内置sum,支持按行计算每个深度点的结果 return np.sum(y1_slice * y2_slice, axis=1) # 单个分片的计算逻辑:输入分片数据,返回该分片的两个结果 def process_single_slice(args): y1_slice, y2_slice = args slice_res1 = fun_1(y1_slice, y2_slice) slice_res2 = fun_2(y1_slice, y2_slice) return slice_res1, slice_res2 def parallel_calculate(y1, y2): # 获取当前CPU核心数,作为进程池大小 process_num = cpu_count() # 按索引均分数据,自动适配任意长度的输入 index_slices = np.array_split(np.arange(len(y1)), process_num) # 准备每个进程要处理的分片数据 task_list = [(y1[idx], y2[idx]) for idx in index_slices] # 创建进程池并执行并行任务 with Pool(process_num) as pool: slice_results = pool.map(process_single_slice, task_list) # 合并所有分片的结果 final_res1 = np.vstack([res[0] for res in slice_results]) final_res2 = np.concatenate([res[1] for res in slice_results]) return final_res1, final_res2 # 执行并行计算 result_1_parallel, result_2_parallel = parallel_calculate(y1, y2) # 可选:验证并行结果和串行结果一致 # 生成串行计算结果 result_1_serial = np.zeros_like(y1) result_2_serial = np.zeros_like(y2) for i in range(len(depth)): result_1_serial[i,:] = fun_1(y1[i], y2[i]) result_2_serial[i] = fun_2(y1[i:i+1], y2[i:i+1])[0] print(np.allclose(result_1_parallel, result_1_serial)) # 输出True表示结果一致 print(np.allclose(result_2_parallel, result_2_serial)) # 输出True表示结果一致
关键细节说明
- 通用数据分片:用
np.array_split按索引均分数据,不管深度序列是否均匀,都能保证分片大小尽量一致,避免手动分割的局限性。 - 计算函数适配:把原
fun_2中的内置sum改为np.sum并指定axis=1,让函数能直接处理数组分片,一次计算整个分片的所有深度点结果。 - 进程池管理:用
with Pool(...)自动管理进程的创建和销毁,map方法自动把任务分配给各个进程,并收集结果。 - 结果合并:根据结果数组的形状,用
np.vstack合并二维的result_1分片,用np.concatenate合并一维的result_2分片。
注意事项
- 如果你的计算函数本身可以用numpy向量化实现(比如
result_1 = y1 + y2),优先用向量化操作,比并行循环效率更高;只有当函数是复杂的非向量化逻辑时,并行化才更有意义。 - 避免在子进程中共享大内存对象,通过传递分片数据的方式减少内存开销和进程间通信成本。
内容的提问来源于stack exchange,提问作者Lucas Oliveira
相关产品推荐
相关产品推荐

