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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 07:37:31