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

如何同时加载多个大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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 14:31:10