如何用Python Multiprocessing并行处理代码并解决HDF5写入报错
问题描述
我有一个generate(file_path)函数,返回整数索引idx和numpy数组,简化实现如下:
def generate(file_path): temp = np.load(file_path) #从file_path字符串获取索引 idx = int(file_path.split["_"][0]) #对temp执行数学运算 result = operate(temp) return idx, result
需要遍历目录收集该函数的结果并写入HDF5文件,串行代码如下:
for path in glob.glob(directory): idx, result = generate(path) hdf5_file["results"][idx,:] = result hdf5_file.close()
尝试修改为多进程代码加速逻辑,但运行报错KeyError: "Unable to open object (object 'results' doesn't exist)",修改后的代码如下:
def generate(file_path): temp = np.load(file_path) #从file_path字符串获取索引 idx = int(file_path.split["_"][0]) #对temp执行数学运算 result = operate(temp) hdf5_path = "./result.hdf5" hdf5_file = h5py.File(hdf5_path, 'w') hdf5_file["results"][idx,:] = result hdf5_file.close() if __name__ == '__main__': ##创建HDF5文件 hdf5_path = "./output.hdf5" hdf5_file = h5py.File(hdf5_path, 'w') hdf5_file.create_dataset("results", [2000,15000], np.uint8) hdf5_file.close() path_ = "./compute/*" p = Pool(mp.cpu_count()) p.map(generate, glob.glob(path_)) hdf5_file.close() print("finished")
错误原因
- 文件路径不匹配:主进程创建的是
./output.hdf5,但子进程里打开的是./result.hdf5,两个完全不同的文件,子进程的文件里根本没有创建results数据集,自然报错。 - 文件打开模式错误:就算路径一致,子进程用
'w'模式打开会清空原有文件,导致之前创建的results数据集被删除,同样会触发找不到对象的错误。 - 语法错误:原代码中
file_path.split["_"]是错误写法,split是函数,调用需用圆括号file_path.split("_")。
解决方案
这里提供两种可靠的多进程实现方式:
方式一:子进程计算,主进程统一写入
这种方式避免多进程直接操作HDF5文件,逻辑简单且安全,适合大多数场景:
import multiprocessing as mp import numpy as np import glob import h5py def operate(temp): # 替换为你的实际运算逻辑 return temp * 2 def generate(file_path): temp = np.load(file_path) idx = int(file_path.split("_")[0]) result = operate(temp) return idx, result if __name__ == '__main__': # 1. 提前创建HDF5文件和数据集 hdf5_path = "./output.hdf5" with h5py.File(hdf5_path, 'w') as hdf5_file: hdf5_file.create_dataset("results", [2000, 15000], np.uint8) # 2. 启动多进程执行计算 path_ = "./compute/*" pool = mp.Pool(mp.cpu_count()) results = pool.map(generate, glob.glob(path_)) pool.close() pool.join() # 3. 主进程统一写入结果到HDF5 with h5py.File(hdf5_path, 'r+') as hdf5_file: for idx, result in results: hdf5_file["results"][idx, :] = result print("finished")
方式二:多进程直接写入HDF5(h5py多进程安全模式)
h5py支持多进程写入,只要所有进程以'r+'模式打开同一个文件,且不同进程写入不同的索引区域(避免并发写同一位置):
import multiprocessing as mp import numpy as np import glob import h5py def operate(temp): # 替换为你的实际运算逻辑 return temp * 2 def generate(file_path): temp = np.load(file_path) idx = int(file_path.split("_")[0]) result = operate(temp) # 以r+模式打开已存在的HDF5文件,避免清空 hdf5_path = "./output.hdf5" with h5py.File(hdf5_path, 'r+') as hdf5_file: hdf5_file["results"][idx, :] = result if __name__ == '__main__': # 1. 提前创建HDF5文件和数据集 hdf5_path = "./output.hdf5" with h5py.File(hdf5_path, 'w') as hdf5_file: hdf5_file.create_dataset("results", [2000, 15000], np.uint8) # 2. 启动多进程执行计算并写入 path_ = "./compute/*" pool = mp.Pool(mp.cpu_count()) pool.map(generate, glob.glob(path_)) pool.close() pool.join() print("finished")
额外注意点
- 若数据集体积极大,方式二的逐进程写入可避免主进程占用过多内存存储所有结果。
- 多进程写入时需确保不同进程操作的
idx无重叠,避免并发写入同一区域导致数据损坏。
内容的提问来源于stack exchange,提问作者Bill
相关产品推荐
相关产品推荐

