如何用multiprocessing高效压缩现有h5py文件?
问题原因
你遇到的文件锁错误,本质是h5py的File对象不能跨进程共享。HDF5文件句柄是进程独占的,子进程继承父进程的文件对象后,多个进程同时操作同一个句柄会触发HDF5的锁机制,导致冲突。你之前的代码中,子进程直接使用全局的hf_in和hf_out,违反了HDF5的进程安全规则。
可行解决方案
方案1:多进程处理+临时文件合并(单机多核首选)
核心思路:每个子进程独立打开输入文件(HDF5允许多进程只读访问),处理一部分数据集并写入临时压缩文件,最后主进程将所有临时文件的数据集合并到最终文件。这样既利用多核并行压缩,又避免了多进程写入同一个文件的冲突。
import h5py import numpy as np from tqdm import tqdm import multiprocessing import os def process_key_chunk(keys, temp_file_idx): # 子进程独立打开输入文件(只读模式)和临时输出文件 with h5py.File('raw_data.hdf5', 'r') as hf_in, \ h5py.File(f'temp_compress_{temp_file_idx}.hdf5', 'w') as hf_temp: for key in keys: # 读取数据集并压缩写入临时文件 data = hf_in[key][()] hf_temp.create_dataset(key, data=data, compression='gzip', compression_opts=9) if __name__ == "__main__": # 主进程获取所有数据集key with h5py.File('raw_data.hdf5', 'r') as hf_in: all_keys = list(hf_in.keys()) # 划分key到多个进程 cpu_count = multiprocessing.cpu_count() key_chunks = [all_keys[i:i+cpu_count] for i in range(0, len(all_keys), cpu_count)] # 启动多进程并行压缩 with multiprocessing.Pool(processes=cpu_count) as pool: tasks = [] for idx, chunk in enumerate(key_chunks): if chunk: # 避免空chunk tasks.append(pool.apply_async(process_key_chunk, args=(chunk, idx))) # 等待所有压缩任务完成 for task in tqdm(tasks, desc="Compressing chunks"): task.get() # 合并临时文件到最终压缩文件 with h5py.File('compressed_data_parallel.hdf5', 'w') as hf_out: for temp_idx in range(len(key_chunks)): temp_file = f'temp_compress_{temp_idx}.hdf5' if os.path.exists(temp_file): with h5py.File(temp_file, 'r') as hf_temp: # 复制临时文件中的数据集到最终文件 for key in tqdm(hf_temp.keys(), desc=f"Merging temp {temp_idx}"): hf_temp.copy(key, hf_out) # 删除临时文件 os.remove(temp_file)
优点
- 无需额外依赖,仅用Python标准库和h5py
- 完全规避多进程文件写入冲突,安全可靠
- 压缩过程完全并行,最大化利用多核CPU
方案2:MPI并行写入(适合集群/单机)
如果可以安装MPI环境,使用mpi4py结合h5py的mpio驱动,就能直接并行写入同一个HDF5文件,不需要临时文件。
步骤1:安装依赖
pip install mpi4py h5py
同时需要安装MPI运行时(比如OpenMPI)。
步骤2:并行压缩代码
from mpi4py import MPI import h5py from tqdm import tqdm comm = MPI.COMM_WORLD rank = comm.Get_rank() total_ranks = comm.Get_size() if rank == 0: # 主进程获取所有数据集key with h5py.File('raw_data.hdf5', 'r') as hf_in: all_keys = list(hf_in.keys()) # 将key均匀分配到各个进程 key_chunks = [all_keys[i::total_ranks] for i in range(total_ranks)] else: key_chunks = None # 分发key chunk到每个进程 local_keys = comm.scatter(key_chunks, root=0) # 使用mpio驱动打开文件,支持多进程并行访问 with h5py.File('raw_data.hdf5', 'r', driver='mpio', comm=comm) as hf_in, \ h5py.File('compressed_data_mpi.hdf5', 'w', driver='mpio', comm=comm) as hf_out: for key in tqdm(local_keys, desc=f"Rank {rank} processing"): data = hf_in[key][()] hf_out.create_dataset(key, data=data, compression='gzip', compression_opts=9) # 等待所有进程完成 comm.Barrier()
运行方式
mpiexec -n 4 python your_script.py # -n 指定进程数,比如4核就写4
优点
- 无需临时文件,直接并行写入最终文件
- 支持集群分布式压缩,适合超大规模数据
方案3:多线程辅助(谨慎使用)
h5py的gzip压缩会释放GIL,因此多线程也能利用多核,但多线程写入同一个HDF5文件存在数据损坏风险,仅适合测试或非关键场景:
import h5py import numpy as np from tqdm import tqdm import threading def compress_single_key(key, hf_in, hf_out): data = hf_in[key][()] hf_out.create_dataset(key, data=data, compression='gzip', compression_opts=9) if __name__ == "__main__": with h5py.File('raw_data.hdf5', 'r') as hf_in, \ h5py.File('compressed_data_thread.hdf5', 'w') as hf_out: threads = [] for key in hf_in.keys(): t = threading.Thread(target=compress_single_key, args=(key, hf_in, hf_out)) threads.append(t) t.start() # 等待所有线程完成 for t in tqdm(threads, desc="Compressing with threads"): t.join()
注意
- 此方案可能导致HDF5文件损坏,生产环境不推荐
- 性能提升可能不如多进程方案
内容的提问来源于stack exchange,提问作者Carl H
相关产品推荐
相关产品推荐

