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

如何在GeoPandas中通过分割数据库实现多进程交集运算

解决方案:分割土地利用数据实现多核并行处理

你的代码只占用1核的核心原因是任务队列仅包含1个任务(snp_distributions里只有一组键值对),所以即使开启8进程池,也只有1个进程在工作。下面是修改后的代码,实现将土地利用数据分割为6份,并行处理空间交集后合并最终结果:

修改核心思路

  1. 新增GeoDataFrame分割函数,将大体积土地利用数据拆分为6份
  2. 主进程提前读取完整土地利用数据并分割,避免子进程重复读取大文件
  3. 构建多进程任务队列,每个任务对应一份分割后的土地利用数据片段
  4. 子进程独立处理单份数据的空间交集运算,保存临时结果
  5. 所有子进程完成后,合并所有临时结果为最终文件

修改后的完整代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 03:05:57