如何使用Dask在GeoDataFrame中递归查找3阶皇后邻接多边形?
Dask 皇后邻接3阶扩展问题修复方案
原代码核心问题分析
- 强制加载全量数据:直接对Dask GeoDataFrame调用
iterrows()会触发全量计算,把所有数据加载到内存,直接导致内存溢出崩溃。 - 共享状态并发冲突:全局
main_chunk_ids集合在多个delayed任务中被修改,会引发数据不一致,且集合修改无法被Dask追踪,导致内存泄漏。 - 递归任务爆炸:递归调用
delayed会生成指数级任务,超出Dask调度器处理能力,最终崩溃。 - 逐行查询效率极低:每个单元格单独查询邻接,没有利用Dask的并行批量处理优势。
修正后的代码实现
import dask_geopandas as dgpd import geopandas as gpd import pandas as pd import dask from tqdm import tqdm import gc queen_out = {} def expand_queen_contiguity(current_ids, combined_dask, max_order=3): current_order = 0 # 用Dask Bag管理分布式ID集合,避免共享状态冲突 result_ids = dask.bag.from_sequence(list(current_ids)) while current_order < max_order: # 筛选当前阶所有单元格的几何 current_cells = combined_dask[combined_dask['uID'].isin(result_ids)] # 批量空间连接实现皇后邻接查询(intersects等价于皇后邻接) adjacency = dgpd.sjoin(combined_dask, current_cells, how='inner', predicate='intersects') # 提取邻接ID并去重 new_ids = adjacency['uID_left'].unique() # 合并结果并去重 result_ids = dask.bag.concat([result_ids, new_ids.to_bag()]).distinct() current_order += 1 gc.collect() # 计算最终结果 return set(result_ids.compute()) for n1 in tqdm(range(1), total=1): # 读取主块数据 main_chunk = gpd.read_parquet(f"./out/singapore/tess_chunk_{int(n1)}.pq") main_chunk_ids = set(main_chunk['uID']) # 合并主块与相邻块(用pd.concat替代append提升效率) combined_chunks = [main_chunk] for n2 in w.neighbors[n1]: neigh_chunk = gpd.read_parquet(f"./out/singapore/tess_chunk_{int(n2)}.pq") combined_chunks.append(neigh_chunk) combined_chunks = pd.concat(combined_chunks, ignore_index=True) # 转为Dask GeoDataFrame并构建空间索引(加速查询) combined_chunks_dask = dgpd.from_geopandas(combined_chunks, npartitions=8) combined_chunks_dask = combined_chunks_dask.set_geometry('geometry') combined_chunks_dask.sindex # 提前构建空间索引 # 执行3阶皇后邻接扩展 queen_area = expand_queen_contiguity(main_chunk_ids, combined_chunks_dask) queen_out[n1] = queen_area
关键改进点
- 迭代替代递归:用循环逐阶扩展,避免生成嵌套任务,降低调度压力。
- 批量空间查询:用
dgpd.sjoin实现并行邻接查询,比逐行检查效率提升数倍。 - 分布式集合管理:用Dask Bag处理ID集合的去重与合并,彻底解决共享状态的并发问题。
- 空间索引优化:提前构建Dask GeoDataFrame的空间索引,大幅减少邻接查询的计算量。
- 懒加载保留:所有操作保持Dask的懒加载特性,仅在最后
compute()时计算结果,降低内存占用。
额外优化建议
- 空间分区:如果数据量极大,可按空间范围或文件大小分区(如
dgpd.from_geopandas(..., partition_size="100MB")),让Dask仅查询相关分区。 - 并行读取块:用Dask delayed封装块读取操作,实现多块并行加载:
@dask.delayed def read_chunk(chunk_id): return gpd.read_parquet(f"./out/singapore/tess_chunk_{int(chunk_id)}.pq") chunks = [read_chunk(n1)] + [read_chunk(n2) for n2 in w.neighbors[n1]] combined_chunks = pd.concat(dask.compute(*chunks), ignore_index=True)
内容的提问来源于stack exchange,提问作者Reuben
相关产品推荐
相关产品推荐

