Python下NetCDF数据提取并行化求助:ThreadPoolExecutor性能劣化
ERA5 NetCDF数据并行处理性能劣化问题
背景与需求
我使用xarray读取ECMWF ERA5的NetCDF数据,数据包含x、y、z、t四个维度,以及r、h、g三个变量。需要按其中三个维度创建目录,将对应变量的最后一维(t维度)所有数据写入纯文本文件。
串行工作代码
串行版本运行正常,简化代码如下:
import xarray as xr def extract_data(data,z,x,y,t): out_str = '' for t_val in t: for z_val in z: out_str += data.sel(x=x, y=y, t=t_val, z=z_val).data return(out_str) def data_get_write(dir_path, r, h, g, x_val, y_val, z, t): r_str = extract_data(r, z, x_val, y_val, t) with open(dir_path+'r_file', 'w') as f_r: f_r.write(r_str) h_str = extract_data(h, z, x_val, y_val, t) with open(dir_path+'h_file', 'w') as f_h: f_h.write(h_str) g_str = extract_data(g, z, x_val, y_val, t) with open(dir_path+'g_file', 'w') as f_g: f_g.write(g_str) return True if __name__ == '__main__': ds=xr.open_dataset('data.nc') r=ds['r'] h=ds['h'] g=ds['g'] z=ds['z'] t=ds['t'] dir_list=[] x_list=[] y_list=[] for val_x in x: for val_y in y: dir_path = f"./output/{val_x}_{val_y}/" # 示例目录生成逻辑 dir_list.append(dir_path) x_list.append(val_x) y_list.append(val_y) list_len=len(dir_list) r_list=[r]*list_len h_list=[h]*list_len g_list=[g]*list_len z_list=[z]*list_len t_list=[t]*list_len results = map(data_get_write, dir_list, r_list, h_list, g_list, x_list, y_list, z_list, t_list)
并行尝试与问题
我尝试用concurrent.futures.ThreadPoolExecutor替换map调用实现并行,代码片段如下:
import concurrent.futures # ... 其余代码与串行版本一致 ... max_workers = min(32, int(arguments["-n"]) + 4) with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: results = executor.map(data_get_write, dir_list, r_list, h_list, g_list, x_list, y_list, z_list, t_list)
但并行版本运行速度远慢于串行:测试集为600MB的小文件,在小型服务器上串行耗时约2分40秒,并行耗时约6分30秒,调整核心数(1到20)无明显性能变化。
选择ThreadPoolExecutor是因为实际数据可达数GB,担心使用进程池时传递数据会引发OOM错误。测试集的x、y维度约20个值,t维度6个值;实际场景中t可达上千,x、y约百级,数据量会比测试集大4-5个数量级。
需求
怀疑当前并行实现有误,或者ThreadPoolExecutor不适用该场景。希望找到标准库支持的并行方案,避免大量改写为asyncio,同时规避OOM风险。
内容的提问来源于stack exchange,提问作者Oliver Henriot
相关产品推荐
相关产品推荐

