Python中并行处理GeoDataFrame行生成多GeoDataFrame的解决方案
并行处理GeoDataFrame行并拼接多结果GeoDataFrame
问题概述
现有函数myfunc,接收GeoDataFrame单行数据,处理后返回多行结构的新GeoDataFrame。需并行执行该函数对输入GeoDataFrame的每一行处理,最终拼接所有结果。此前尝试apply方法无效(函数返回GeoDataFrame而非单行值),使用joblib.Parallel时触发PicklingError: Could not pickle the task to send it to the workers错误。
核心问题分析
PicklingError的根源是:默认序列化机制无法处理GeoPandas几何对象、rasterio资源对象,或函数依赖未传入的全局变量(全局变量序列化易出问题)。
解决方案
方案1:JobLib + Loky后端(推荐)
通过指定loky后端并启用高版本pickle协议,同时消除函数对全局变量的依赖,解决序列化问题。
修改后的代码
import geopandas as gpd import joblib from joblib import Parallel, delayed import rasterio from shapely.geometry import LineString import os def myfunc(row, raster_file, crs): line = row.geometry dist_list = [0, row.dist1, row.dist2, line.length] final_coords = [] with rasterio.open(raster_file) as raster_obj: for dist in dist_list: Pt = line.interpolate(dist) row_idx, col_idx = raster_obj.index(Pt.x, Pt.y) window = rasterio.windows.Window(col_idx, row_idx, 1, 1) raster_value = raster_obj.read(1, window=window)[0][0] final_coords.append((Pt.x, Pt.y, raster_value)) line_strings = [LineString(final_coords[i:i+2]) for i in range(len(final_coords) - 1)] gpd_out = gpd.GeoDataFrame(line_strings, columns=["geometry"], geometry="geometry", crs=crs) gpd_out["Dist"] = dist_list[:-1] return gpd_out # 主执行逻辑 work_folder = r"define/your/work/folder" gpd_in = gpd.read_file(os.path.join(work_folder,"Lines.shp")) raster_file = os.path.join(work_folder, 'random_raster.tif') crs = gpd_in.crs # 并行配置:使用CPU数-1避免资源耗尽 num_workers = joblib.cpu_count() - 1 results = Parallel( n_jobs=num_workers, backend="loky", verbose=10, pickle_protocol=5 )(delayed(myfunc)(row, raster_file, crs) for row in gpd_in.itertuples()) # 拼接所有结果 gpd_out = gpd.pd.concat(results, ignore_index=True).set_crs(crs)
方案2:Multiprocessing + Dill序列化
若JobLib仍有问题,可使用multiprocessing配合dill库(支持更多Python对象序列化)。
修改后的代码
import geopandas as gpd import multiprocessing as mp import dill import rasterio from shapely.geometry import LineString import os def myfunc(row, raster_file, crs): line = row.geometry dist_list = [0, row.dist1, row.dist2, line.length] final_coords = [] with rasterio.open(raster_file) as raster_obj: for dist in dist_list: Pt = line.interpolate(dist) row_idx, col_idx = raster_obj.index(Pt.x, Pt.y) window = rasterio.windows.Window(col_idx, row_idx, 1, 1) raster_value = raster_obj.read(1, window=window)[0][0] final_coords.append((Pt.x, Pt.y, raster_value)) line_strings = [LineString(final_coords[i:i+2]) for i in range(len(final_coords) - 1)] gpd_out = gpd.GeoDataFrame(line_strings, columns=["geometry"], geometry="geometry", crs=crs) gpd_out["Dist"] = dist_list[:-1] return gpd_out # 主执行逻辑 work_folder = r"define/your/work/folder" gpd_in = gpd.read_file(os.path.join(work_folder,"Lines.shp")) raster_file = os.path.join(work_folder, 'random_raster.tif') crs = gpd_in.crs # 初始化进程池,使用dill处理序列化 pool = mp.Pool(processes=mp.cpu_count()-1) tasks = [(row, raster_file, crs) for row in gpd_in.itertuples()] results = pool.starmap(myfunc, tasks) pool.close() pool.join() # 拼接结果 gpd_out = gpd.pd.concat(results, ignore_index=True).set_crs(crs)
关键注意事项
- 消除全局变量依赖:将
raster_file、crs等必要参数传入函数,避免函数引用全局变量,这是解决序列化错误的核心。 - Pickle协议版本:Python 3.8+使用
pickle_protocol=5可更好处理复杂对象与大数据。 - CRS一致性:确保每个返回的GeoDataFrame指定正确CRS,避免拼接后坐标系统丢失。
- 进程数设置:使用
cpu_count()-1预留系统资源,防止卡顿。
内容的提问来源于stack exchange,提问作者MindOverMatter
相关产品推荐
相关产品推荐

