如何在GeoPandas中通过分割数据库实现多进程交集运算
解决方案:分割土地利用数据实现多核并行处理
你的代码只占用1核的核心原因是任务队列仅包含1个任务(snp_distributions里只有一组键值对),所以即使开启8进程池,也只有1个进程在工作。下面是修改后的代码,实现将土地利用数据分割为6份,并行处理空间交集后合并最终结果:
修改核心思路
- 新增GeoDataFrame分割函数,将大体积土地利用数据拆分为6份
- 主进程提前读取完整土地利用数据并分割,避免子进程重复读取大文件
- 构建多进程任务队列,每个任务对应一份分割后的土地利用数据片段
- 子进程独立处理单份数据的空间交集运算,保存临时结果
- 所有子进程完成后,合并所有临时结果为最终文件
修改后的完整代码
import os import geopandas as gpd import pandas as pd from multiprocessing import Pool import time import tempfile def split_gdf(gdf, n_splits): """将GeoDataFrame按行分割为指定份数""" split_size = len(gdf) // n_splits splits = [] for i in range(n_splits): start = i * split_size # 最后一份包含剩余所有行,避免数据遗漏 end = start + split_size if i < n_splits - 1 else len(gdf) splits.append(gdf.iloc[start:end].copy()) return splits def process_land_use_segment(args): """处理单份土地利用数据片段的空间交集与排放量计算""" selected_snp, land_use_category, gdf_land_use_segment, output_path, temp_dir = args # 读取网格排放密度文件 grid_density_file = os.path.join(output_path, f"emiss_{selected_snp}_density.shp") gdf_density = gpd.read_file(grid_density_file) # 统一CRS为EPSG:4326,确保空间运算兼容性 if gdf_density.crs != "EPSG:4326": gdf_density = gdf_density.to_crs("EPSG:4326") if gdf_land_use_segment.crs != "EPSG:4326": gdf_land_use_segment = gdf_land_use_segment.to_crs("EPSG:4326") # 执行空间交集运算 intersection = gpd.overlay(gdf_density, gdf_land_use_segment, how='intersection') # 计算各污染物排放量 for pollutant in ["CO", "NH3", "NMVOC", "NOx", "SO2", "PMc", "PM2_5"]: if f"{pollutant}_m2" in intersection.columns and "area" in intersection.columns: intersection[f"{pollutant}_emissions"] = intersection[f"{pollutant}_m2"] * intersection["area"] # 保存临时结果文件 temp_file = os.path.join(temp_dir, f"{selected_snp}_{land_use_category}_temp_{os.getpid()}.shp") intersection.to_file(temp_file) return temp_file if __name__ == '__main__': start_time = time.time() # 路径配置 output_path = "/outpath/" land_use_path = "/pathlu/" n_processes = 6 # 匹配分割份数,充分利用6个核心 # 目标SNP与土地利用类别 selected_snp = "SNP1" land_use_category = "ind" # 读取完整土地利用数据并初始化CRS land_use_file = f"lu_ll_{land_use_category}.shp" gdf_land_use_full = gpd.read_file(os.path.join(land_use_path, land_use_file)) if gdf_land_use_full.crs is None: gdf_land_use_full.crs = "EPSG:4326" # 将土地利用数据分割为6份 gdf_splits = split_gdf(gdf_land_use_full, n_processes) # 创建临时目录存储中间结果,自动清理 with tempfile.TemporaryDirectory() as temp_dir: # 构建多进程任务参数列表 tasks = [ (selected_snp, land_use_category, split, output_path, temp_dir) for split in gdf_splits ] # 启动进程池并行处理 with Pool(n_processes) as pool: temp_files = pool.map(process_land_use_segment, tasks) # 合并所有临时结果 merged_gdf = gpd.GeoDataFrame() for temp_file in temp_files: gdf = gpd.read_file(temp_file) merged_gdf = pd.concat([merged_gdf, gdf], ignore_index=True) # 保存最终合并结果 final_output = os.path.join(output_path, f"{selected_snp}_{land_use_category}_emissions.shp") merged_gdf.to_file(final_output) end_time = time.time() print(f"执行完成,耗时 {end_time - start_time:.2f} 秒,即 {(end_time - start_time)/60:.2f} 分钟")
关键细节说明
- 分割逻辑:按行均分土地利用数据,最后一份自动包含剩余行,避免数据遗漏
- 效率优化:主进程提前读取完整土地利用数据,子进程直接处理分割后的片段,减少重复IO开销
- 临时文件管理:使用
tempfile.TemporaryDirectory自动清理临时文件,避免磁盘冗余 - CRS一致性:强制统一网格数据与土地利用数据的CRS,避免空间运算报错
- 进程匹配:设置进程数与分割份数一致,确保6个核心同时工作
内容的提问来源于stack exchange,提问作者Nat
相关产品推荐
相关产品推荐

