PySAGA多进程场景下如何序列化(Pickle)SWIG对象?
如何实现PySAGA中SWIG对象的Pickle序列化以完成多进程任务?
你尝试使用PySAGA进行Python多进程处理,为避免每个工作进程重复加载输入栅格,计划将已加载的栅格(SWIG包装的对象)通过Pickle传递给多进程任务,但任务执行失败。相关代码如下:
import os,sys import time, math from joblib import Parallel, delayed import numpy as np os.environ['SAGA_PATH'] = 'C:/Program Files/SAGA/' os.add_dll_directory(os.environ['SAGA_PATH']) sys.path.append('C:/Program Files/SAGA/') def Run_Gravitational_Process_Path_Model(dem,source,mu,md): # Get the tool: Tool = saga.SG_Get_Tool_Library_Manager().Create_Tool('sim_geomorphology', '0') Verbose = False if not Verbose: saga.SG_UI_ProgressAndMsg_Lock(True) if not Tool: print('Failed to request tool: Gravitational Process Path Model') return False # Set the parameter interface: Tool.Reset() Tool.Set_Parameter('DEM', dem) Tool.Set_Parameter('RELEASE_AREAS',source) Tool.Set_Parameter('MATERIAL', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('FRICTION_ANGLE_GRID', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('SLOPE_IMPACT_GRID', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('FRICTION_MU_GRID', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('FRICTION_MASS_TO_DRAG_GRID', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('OBJECTS', saga.SG_Get_Data_Manager().Add('Grid input file, optional')) Tool.Set_Parameter('DEPOSITION', saga.SG_Get_Create_Pointer()) Tool.Set_Parameter('MAX_VELOCITY', saga.SG_Get_Create_Pointer()) Tool.Set_Parameter('STOP_POSITIONS', saga.SG_Get_Create_Pointer()) Tool.Set_Parameter('ENDANGERED', saga.SG_Get_Create_Pointer()) Tool.Set_Parameter('PROCESS_PATH_MODEL', 1) # 'Random Walk' Tool.Set_Parameter('RW_SLOPE_THRES', 40) Tool.Set_Parameter('RW_EXPONENT', 2) Tool.Set_Parameter('RW_PERSISTENCE', 1.5) Tool.Set_Parameter('GPP_ITERATIONS', 1000) Tool.Set_Parameter('GPP_PROCESSING_ORDER', 2) # 'RAs in Parallel per Iteration' Tool.Set_Parameter('GPP_SEED', 1) Tool.Set_Parameter('FRICTION_MODEL', 5) # 'None' Tool.Set_Parameter('FRICTION_THRES_FREE_FALL', 60) Tool.Set_Parameter('FRICTION_METHOD_IMPACT', 0) # 'Energy Reduction (Scheidegger 1975)' Tool.Set_Parameter('FRICTION_IMPACT_REDUCTION', 75) Tool.Set_Parameter('FRICTION_ANGLE', 30) Tool.Set_Parameter('FRICTION_MU', mu) Tool.Set_Parameter('FRICTION_MODE_OF_MOTION', 0) # 'Sliding' Tool.Set_Parameter('FRICTION_MASS_TO_DRAG',md) Tool.Set_Parameter('FRICTION_INIT_VELOCITY', 1) Tool.Set_Parameter('DEPOSITION_MODEL', 0) # 'None' Tool.Set_Parameter('DEPOSITION_INITIAL', 20) Tool.Set_Parameter('DEPOSITION_SLOPE_THRES', 20) Tool.Set_Parameter('DEPOSITION_VELOCITY_THRES', 15) Tool.Set_Parameter('DEPOSITION_MAX', 20) Tool.Set_Parameter('DEPOSITION_MIN_PATH', 100) Tool.Set_Parameter('SINK_MIN_SLOPE', 2.5) # Execute the tool: if not Tool.Execute(): print('failed to execute tool: ' + Tool.Get_Name().c_str()) return False # Request the results: return Tool.Get_Parameter('PROCESS_AREA').asGrid(),\ Tool.Get_Parameter('DEPOSITION').asGrid(), \ Tool.Get_Parameter('MAX_VELOCITY').asGrid(), \ Tool.Get_Parameter('STOP_POSITIONS').asGrid(), \ Tool.Get_Parameter('ENDANGERED').asGrid(), def load_data(path,dem_name,source_name): saga.SG_Set_History_Depth(0) # History will not be created dem = saga.SG_Get_Data_Manager().Add_Grid('{}\\{}.tif'.format(raster_input_path,dem_name)) source = saga.SG_Get_Data_Manager().Add_Grid('{}\\{}.tif'.format(raster_input_path,source_name)) return dem,source def run_GPP(job_number): mu = float(parameters[job_number][0]) md = float(parameters[job_number][1]) print('mu=', round(mu, 2), 'md=', md) result=Run_Gravitational_Process_Path_Model(dem,source,mu,md) result[0] .Save('{}_{}_{}.tif'.format(r"Downloads\Test\process_area",round(mu,2),round(md,2))) mu_range=np.arange(0.2, 1, 0.1) md_range=np.arange(200, 230, 2.5) parameters=[] for mu in mu_range: for md in md_range: parameters.append((mu,md)) number_of_jobs =len(parameters) dem,source=load_data(raster_input_path,'dtm','release') # Run in parallel with Joblib start = time.time() Parallel(n_jobs=-1,verbose=1)(delayed(run_GPP)(i) for i in range(number_of_jobs)) # when job is done, free memory resources: saga.SG_Get_Data_Manager().Delete_All() end = time.time() print('{:.4f} s'.format(end-start))
核心问题
PySAGA中的栅格对象是SWIG包装的C++对象,这类对象默认不支持Pickle序列化/反序列化。因为SWIG对象绑定的是父进程的内存资源,跨进程传递后,子进程无法访问父进程的内存空间,导致对象失效,任务执行失败。
可行解决方案
方案1:子进程内独立加载栅格(推荐)
放弃直接传递SWIG对象,改为在每个子进程中通过文件路径重新加载栅格。虽然看似重复加载,但SAGA内部会对文件进行缓存,实际开销远低于预期,且实现简单可靠。
调整run_GPP函数
def run_GPP(job_number, raster_input_path, dem_name, source_name): mu = float(parameters[job_number][0]) md = float(parameters[job_number][1]) print('mu=', round(mu, 2), 'md=', md) # 子进程内部加载栅格数据 dem, source = load_data(raster_input_path, dem_name, source_name) result = Run_Gravitational_Process_Path_Model(dem, source, mu, md) result[0].Save('{}_{}_{}.tif'.format(r"Downloads\Test\process_area", round(mu,2), round(md,2))) # 子进程内清理内存,避免资源泄漏 saga.SG_Get_Data_Manager().Delete_All()
修改Parallel调用逻辑
# Run in parallel with Joblib start = time.time() Parallel(n_jobs=-1, verbose=1)( delayed(run_GPP)(i, raster_input_path, 'dtm', 'release') for i in range(number_of_jobs) ) # 父进程清理资源 saga.SG_Get_Data_Manager().Delete_All() end = time.time() print('{:.4f} s'.format(end-start))
方案2:自定义Pickle序列化(复杂且不推荐)
如果必须传递已加载的对象,需要为SWIG栅格对象自定义Pickle方法。但该方法需要深入了解PySAGA的内部实现,且容易因版本更新失效。大致思路是:
- 序列化时:提取栅格的元数据(如尺寸、投影、nodata值)和数据数组(转为numpy数组)
- 反序列化时:根据元数据重建SAGA栅格对象,再将numpy数组写入对象
示例代码框架(仅作参考,需根据实际情况调整):
import pickle from saga import CSG_Grid # 为CSG_Grid注册pickle方法 def csg_grid_pickle(grid): # 提取元数据 meta = { 'nx': grid.Get_NX(), 'ny': grid.Get_NY(), 'xmin': grid.Get_XMin(), 'xmax': grid.Get_XMax(), 'ymin': grid.Get_YMin(), 'ymax': grid.Get_YMax(), 'nodata': grid.Get_NoData_Value() } # 提取数据为numpy数组 data = grid.asarray() return (meta, data) def csg_grid_unpickle(state): meta, data = state # 重建栅格对象 grid = CSG_Grid() grid.Create(meta['nx'], meta['ny'], meta['xmin'], meta['xmax'], meta['ymin'], meta['ymax']) grid.Set_NoData_Value(meta['nodata']) # 将numpy数组写入栅格 grid.fromarray(data) return grid # 注册pickle方法 pickle.pickle(CSG_Grid, csg_grid_pickle, csg_grid_unpickle)
注意:该方法可能因PySAGA版本差异无法正常工作,且SAGA栅格对象可能包含更多内部状态(如投影信息)需要处理,稳定性较差。
内容的提问来源于stack exchange,提问作者Hans Baumann
相关产品推荐
相关产品推荐

