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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 08:31:19