使用多进程调用geopandas.overlay无报错但始终无法完成的问题排查
问题
我尝试通过多进程调用geopandas.overlay()提升处理速度:
- 此前用自定义函数结合
functools.partial,传入唯一ID子集化DataFrame的方式能成功运行,最终合并结果。 - 改用
numpy.array_split()拆分DataFrame后,多进程启动后处理器全部关闭,进程挂起,无任何工作或退出迹象;尝试spawn上下文启动进程也没改善,但直接运行geopandas.overlay(tracks_gdf, grid_gdf)能得到预期结果。
疑问:
- 为何两种partial函数调用方式结果不同?
numpy.array_split()的返回结果是否为可迭代对象?- 如何分块将单个DataFrame传入
geopandas.overlay()以利用多进程,并得到可合并的结果?
成功运行的代码示例
def taska(id, points, crs): return make_break_points((vms_points[points.ID == id]).reset_index(drop=True), crs) points_gdf = geodataframe of points with an id field grid_gdf = geodataframe polygon grid partialA = functools.partial(taska, points=points_gdf, crs=grid_gdf.crs) partialA_results =[] with Pool(cpu_count()-4) as pool: for results in pool.map(partialA, list(points_gdf.ID.unique())): partialA_results.append(results) bpts_gdf = pd.concat(partialA_results)
出现问题的代码示例
def taskc(tracks, grid): return gpd.overlay(tracks, grid, how='union').explode().reset_index(drop=True) tracks_gdf = geodataframe of points with an id field dfs = np.array_split(tracks_gdf, (cpu_count()-4)) grid_gdf = geodataframe polygon grid partialC_results = [] partialC = functools.partial(taskc, grid=grid_gdf) with Pool(cpu_count() - 4) as pool: for results in pool.map(partialC, dfs): partialC_results.append(results) results_df = pd.concat(partialC_results)
解答
1. 两种partial调用方式结果不同的原因
- 数据传递开销与序列化差异:成功的代码中,
pool.map传递的是轻量的ID值,每个进程仅需根据ID从传入的points_gdf中筛选子集,数据传输成本极低;而出问题的代码中,np.array_split拆分后的GeoDataFrame块包含几何数据,体积大,跨进程序列化/反序列化开销极高,且GeoPandas依赖的GEOS库对象用默认pickle序列化容易出现隐式错误,直接导致进程无响应。 - 任务复杂度差异:
taska仅做简单的子集筛选和断点生成,taskc调用的overlay是计算密集型操作,大体积GeoDataFrame块会让子进程负载骤增,加上序列化问题,直接触发进程挂起。
2. numpy.array_split()的返回结果是否为可迭代对象?
是,np.array_split()返回的是包含拆分后DataFrame的列表,完全支持迭代。进程挂起和这个无关,核心原因是GeoDataFrame的跨进程序列化问题。
3. 正确分块多进程处理的方案
方案一:沿用ID分组模式(推荐,稳定性高)
复用你成功的ID子集化逻辑,避免直接传递大体积GeoDataFrame:
def taskc(id, tracks, grid): track_subset = tracks[tracks.ID == id].reset_index(drop=True) # 跳过空子集,避免无意义计算 if not track_subset.empty: return gpd.overlay(track_subset, grid, how='union').explode().reset_index(drop=True) return None tracks_gdf = geodataframe of points with an id field grid_gdf = geodataframe polygon grid partialC = functools.partial(taskc, tracks=tracks_gdf, grid=grid_gdf) partialC_results = [] with Pool(cpu_count() - 4) as pool: for result in pool.map(partialC, tracks_gdf.ID.unique()): if result is not None: partialC_results.append(result) results_df = pd.concat(partialC_results)
方案二:改进序列化(适合必须按DataFrame块拆分的场景)
用dill替代默认pickle处理GeoDataFrame序列化,配合spawn上下文:
import dill from multiprocessing import get_context def taskc(tracks, grid): return gpd.overlay(tracks, grid, how='union').explode().reset_index(drop=True) tracks_gdf = geodataframe of points with an id field grid_gdf = geodataframe polygon grid dfs = np.array_split(tracks_gdf, cpu_count()-4) # 配置spawn上下文与dill序列化 ctx = get_context('spawn') ctx.set_serializer('pickle', dill.dumps, dill.loads) partialC_results = [] partialC = functools.partial(taskc, grid=grid_gdf) with ctx.Pool(cpu_count()-4) as pool: for result in pool.map(partialC, dfs): partialC_results.append(result) results_df = pd.concat(partialC_results)
注意:该方案仍存在较大数据传输开销,数据量较大时优先选ID分组模式。
内容的提问来源于stack exchange,提问作者MrKingsley
相关产品推荐
相关产品推荐

