如何同时加载多个大Pickle文件以优化加载时长?
解决大Pickle文件并行加载与按month维度拼接的方案
先修正原代码的核心bug
原代码中file_path在循环内被重新赋值为具体文件路径,导致后续循环执行os.listdir(file_path)时会触发错误,需将目录路径与文件路径分开存放:
import pickle import os import xarray as xr # 单独存放目录路径,避免循环中被覆盖 dir_path = "/pickle/file/dir/" loaded_data = [] pickle_files = sorted([f for f in os.listdir(dir_path) if f.endswith(".pkl")]) for filename in pickle_files: file_path = os.path.join(dir_path, filename) with open(file_path, 'rb') as f: data = pickle.load(f) loaded_data.append(data) final_data = xr.concat(loaded_data, dim='month')
并行加载优化方案(替代父子脚本,更高效)
父子脚本模式需手动处理进程间通信,用multiprocessing实现多进程并行加载更简洁,能有效利用多核资源破解单线程加载瓶颈:
实现代码
import pickle import os import xarray as xr from multiprocessing import Pool dir_path = "/pickle/file/dir/" def load_pickle_file(filename): """子进程执行的单个文件加载函数""" file_path = os.path.join(dir_path, filename) with open(file_path, 'rb') as f: data = pickle.load(f) # 若每个Pickle内的字典列表需转换为xarray对象,可在此完成 # 示例:return xr.Dataset(data)(需根据实际数据结构调整) return data if __name__ == "__main__": pickle_files = sorted([f for f in os.listdir(dir_path) if f.endswith(".pkl")]) # 进程数建议设为CPU核心数-1,避免资源耗尽 with Pool(processes=os.cpu_count()-1) as pool: loaded_data = pool.map(load_pickle_file, pickle_files) # 按month维度完成拼接 final_data = xr.concat(loaded_data, dim='month') # 可选:将拼接结果保存为NetCDF(比Pickle更适合xarray数据长期存储) # final_data.to_netcdf("combined_monthly_data.nc")
关键注意事项
- 内存要求:12个4GB文件加载后总内存需求约48GB,需确保机器有足够内存;若内存不足,可分批次拼接后写入磁盘释放内存
- 数据格式适配:如果每个Pickle中的字典列表需先转换为xarray的
Dataset/DataArray,请在load_pickle_file函数内完成转换,确保返回值可被xr.concat处理
父子脚本模式实现(按需选择)
若坚持使用父子脚本,可通过subprocess调用子脚本加载单个文件,将结果写入临时NetCDF文件,父脚本再读取临时文件完成拼接:
子脚本(load_single_file.py)
import pickle import xarray as xr import sys def main(): input_path = sys.argv[1] output_path = sys.argv[2] with open(input_path, 'rb') as f: data = pickle.load(f) ds = xr.Dataset(data) # 根据实际数据结构调整转换逻辑 ds.to_netcdf(output_path) if __name__ == "__main__": main()
父脚本
import os import sys import xarray as xr import subprocess from tempfile import TemporaryDirectory dir_path = "/pickle/file/dir/" pickle_files = sorted([f for f in os.listdir(dir_path) if f.endswith(".pkl")]) # 用临时目录存放子脚本生成的中间文件 with TemporaryDirectory() as temp_dir: temp_file_paths = [] for idx, filename in enumerate(pickle_files): input_file = os.path.join(dir_path, filename) temp_file = os.path.join(temp_dir, f"temp_month_{idx}.nc") # 调用子脚本处理单个文件 subprocess.run([sys.executable, "load_single_file.py", input_file, temp_file], check=True) temp_file_paths.append(temp_file) # 读取中间文件并完成拼接 loaded_datasets = [xr.open_dataset(f) for f in temp_file_paths] final_data = xr.concat(loaded_datasets, dim='month')
父子脚本优缺点
- 优势:进程完全隔离,避免多进程共享内存的潜在问题
- 劣势:依赖磁盘IO存储中间文件,速度略低于多进程直接内存交换
内容的提问来源于stack exchange,提问作者Ehsan
相关产品推荐
相关产品推荐

