blob: 15dd1c4415e23161be7ded2b831e0a091b0d538c [file] [view]
<!--
Licensed to the Apache Software Foundation (ASF) under one
or more contributor license agreements. See the NOTICE file
distributed with this work for additional information
regarding copyright ownership. The ASF licenses this file
to you under the Apache License, Version 2.0 (the
"License"); you may not use this file except in compliance
with the License. You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing,
software distributed under the License is distributed on an
"AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
KIND, either express or implied. See the License for the
specific language governing permissions and limitations
under the License.
-->
# Apache Sedona 上的 GeoPandas API
Apache Sedona 上的 GeoPandas API 提供了与 GeoPandas 一致的接口,可以让您的地理空间分析突破单机的局限。该 API 把熟悉的 GeoPandas DataFrame 语法与 Apache Sedona 在 Apache Spark 上的分布式处理能力结合起来,让您能够使用同样的代码模式处理行星级(planetary-scale)的数据集。
## 概览
### 什么是 Apache Sedona 的 GeoPandas API?
Apache Sedona 的 GeoPandas API 是一层兼容层,让您能够在分布式地理空间数据上使用 GeoPandas 风格的操作。您原有的 GeoPandas 代码不再受限于单机处理能力,可以充分利用 Apache Spark 集群进行大规模地理空间分析。
### 主要优势
- **熟悉的 API**:使用您已经熟悉的 GeoPandas 语法与方法
- **分布式处理**:突破单机限制,处理大规模数据集
- **惰性求值**:受益于 Apache Sedona 的查询优化与延迟执行
- **高性能**:利用分布式计算高效完成复杂的地理空间运算
- **平滑迁移**:以最少的代码改动迁移现有 GeoPandas 工作流
## 准备工作
Apache Sedona 的 GeoPandas API 通过 PySpark 的 pandas-on-Spark 集成自动管理 SparkSession。准备方式有两种:
### 方式 1:自动创建 SparkSession(推荐)
GeoPandas API 会自动使用 PySpark 的默认 SparkSession:
```python
from sedona.spark.geopandas import GeoDataFrame, read_parquet
# 无需显式配置 SparkSession,使用默认会话
# API 会自动初始化 Sedona context
```
### 方式 2:手动配置 SparkSession
如果需要自定义 SparkSession,或需要在某些环境下显式控制:
```python
from sedona.spark.geopandas import GeoDataFrame, read_parquet
from sedona.spark import SedonaContext
# 创建并配置 SparkSession
config = SedonaContext.builder().getOrCreate()
sedona = SedonaContext.create(config)
# GeoPandas API 会使用上面配置好的会话
```
### 方式 3:复用现有 SparkSession
如果已有 SparkSession(如 Databricks、EMR 等托管环境):
```python
from sedona.spark.geopandas import GeoDataFrame, read_parquet
from sedona.spark import SedonaContext
# 复用现有 SparkSession(如 Databricks 中的 `spark`)
sedona = SedonaContext.create(spark) # `spark` 即已有会话
```
### SparkSession 是如何被管理的
GeoPandas API 借助 PySpark 的 pandas-on-Spark 自动管理 SparkSession 生命周期:
1. **默认会话**:导入 `sedona.spark.geopandas` 时,会自动通过 `pyspark.pandas.utils.default_session()` 获取 PySpark 的默认会话。
2. **自动注册 Sedona 函数**:必要时 API 会自动把 Sedona 的空间函数与优化注册到 SparkSession。
3. **透明集成**:所有 GeoPandas 操作在底层都被翻译为 Spark SQL 操作,使用所配置的 SparkSession 执行。
4. **无需手动管理 context**:与传统的 Sedona 用法不同,您通常不需要显式调用 `SedonaContext.create()`,除非有自定义配置需求。
这种设计让 API 更加易用,把 SparkSession 管理的复杂度隐藏起来,同时仍能充分发挥分布式处理的能力。
### S3 配置
在使用 S3 数据时,GeoPandas API 使用 Spark 内置的 S3 支持,而不是 s3fs 等外部库。可通过 Spark 配置开启对公开 S3 桶的匿名访问:
```python
from sedona.spark import SedonaContext
# 公开 S3 桶的匿名访问
config = (
SedonaContext.builder()
.config(
"spark.hadoop.fs.s3a.bucket.bucket-name.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider",
)
.getOrCreate()
)
sedona = SedonaContext.create(config)
```
需要鉴权访问 S3 时,请使用合适的 AWS 凭据提供者:
```python
# 使用 IAM 角色(在 EC2/EMR 上推荐)
config = (
SedonaContext.builder()
.config(
"spark.hadoop.fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.InstanceProfileCredentialsProvider",
)
.getOrCreate()
)
# 使用 access key(不推荐用于生产环境)
config = (
SedonaContext.builder()
.config("spark.hadoop.fs.s3a.access.key", "your-access-key")
.config("spark.hadoop.fs.s3a.secret.key", "your-secret-key")
.getOrCreate()
)
```
## 基本用法
### 引入 API
不再直接导入 GeoPandas,而是从 Sedona 的 GeoPandas 模块导入:
```python
# 传统的 GeoPandas 导入
# import geopandas as gpd
# Sedona GeoPandas API 的导入
import sedona.spark.geopandas as gpd
# 或
from sedona.spark.geopandas import GeoDataFrame, read_parquet
```
### 读取数据
API 支持读取多种地理空间格式,包括来自云存储的 Parquet 文件。要以匿名凭据访问 S3,请配置 Spark 使用匿名 AWS 凭据:
```python
from sedona.spark import SedonaContext
# 配置 Spark 进行匿名 S3 访问
config = (
SedonaContext.builder()
.config(
"spark.hadoop.fs.s3a.bucket.wherobots-examples.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider",
)
.getOrCreate()
)
sedona = SedonaContext.create(config)
# 直接从 S3 加载 GeoParquet 文件
s3_path = "s3://wherobots-examples/data/onboarding_1/nyc_buildings.parquet"
nyc_buildings = gpd.read_parquet(s3_path)
# 显示基本信息
print(f"Dataset shape: {nyc_buildings.shape}")
print(f"Columns: {nyc_buildings.columns.tolist()}")
nyc_buildings.head()
```
### 空间过滤
使用 `cx` 进行分布式坐标范围筛选;如需将相交的几何裁剪到掩膜边界,
可继续使用 `clip`:
```python
from shapely.geometry import box
# 中央公园的边界框
central_park_bbox = box(
-73.973,
40.764, # 左下角(经度,纬度)
-73.951,
40.789, # 右上角(经度,纬度)
)
# 筛选与坐标范围相交的要素
central_park_buildings = nyc_buildings.cx[
-73.973:-73.951,
40.764:40.789,
]
# 将筛选后的几何裁剪到精确的多边形边界
central_park_buildings = central_park_buildings.clip(central_park_bbox)
# 显示结果
print(
central_park_buildings[["BUILD_ID", "PROP_ADDR", "height_val", "geometry"]].head()
)
```
**针对大规模数据的替代方案——使用空间连接:**
```python
# 用边界框创建一个 GeoDataFrame
bbox_gdf = gpd.GeoDataFrame({"id": [1]}, geometry=[central_park_bbox], crs="EPSG:4326")
# 使用空间连接筛选边界框内的建筑
central_park_buildings = nyc_buildings.sjoin(bbox_gdf, predicate="intersects")
```
## 高级操作
### 空间连接
可使用与 GeoPandas 相同的语法执行空间连接:
```python
# 加载两份数据集
left_df = gpd.read_parquet("s3://bucket/left_data.parquet")
right_df = gpd.read_parquet("s3://bucket/right_data.parquet")
# 使用 distance 谓词的空间连接
result = left_df.sjoin(right_df, predicate="dwithin", distance=50)
# 其他空间谓词
intersects_result = left_df.sjoin(right_df, predicate="intersects")
contains_result = left_df.sjoin(right_df, predicate="contains")
```
### 坐标参考系操作
在不同坐标参考系(CRS)之间转换几何对象:
```python
# GeoParquet 通常会保留 CRS 元数据,仅在缺失时进行设置
buildings = gpd.read_parquet("buildings.parquet")
if buildings.crs is None:
buildings = buildings.set_crs("EPSG:4326")
# 转换为投影坐标系以便计算面积
buildings_projected = buildings.to_crs("EPSG:3857")
# 计算面积
buildings_projected["area"] = buildings_projected.geometry.area
```
### 几何运算
应用几何变换与分析:
```python
# 缓冲区操作
buffered = buildings.geometry.buffer(100) # 100 米缓冲
# 几何属性
buildings["is_valid"] = buildings.geometry.is_valid
buildings["is_simple"] = buildings.geometry.is_simple
buildings["bounds"] = buildings.geometry.bounds
# 距离计算
from shapely.geometry import Point
reference_point = Point(-73.9857, 40.7484) # 时代广场
buildings["distance_to_times_square"] = buildings.geometry.distance(reference_point)
# 面积与周长(需要使用投影 CRS)
buildings_projected = buildings.to_crs("EPSG:3857") # Web Mercator
buildings_projected["area"] = buildings_projected.geometry.area
buildings_projected["perimeter"] = buildings_projected.geometry.length
```
## 性能考虑
### 何时仍使用传统 GeoPandas:
- 小数据集(< 1GB)
- 对本地数据的简单操作
- 需要完整的功能覆盖
- 单机处理已经足够
### 何时使用 Apache Sedona 的 GeoPandas API:
- 大规模数据集(> 1GB)
- 复杂的地理空间分析
- 需要分布式处理
- 数据保存在云存储(S3、HDFS 等)
## 已支持的操作
Apache Sedona 的 GeoPandas API 已实现最常用的 GeoSeries 与 GeoDataFrame 操作:
### 数据 I/O
- `read_parquet()` —— 读取 GeoParquet 文件
- `read_file()` —— 读取多种地理空间格式
- `to_parquet()` —— 写出为 Parquet 格式
### 空间操作
- `sjoin()` —— 多种谓词的空间连接
- `cx` —— 基于坐标范围的空间筛选
- `clip()` —— 使用标量、矩形或分布式掩膜裁剪几何
- `overlay()` —— 支持全部五种 GeoPandas 模式的分布式数据框叠加
- `buffer()` —— 几何缓冲
- `distance()` —— 距离计算
- `intersects()`、`contains()`、`within()` —— 空间谓词
- `geom_equals_identical()` —— 对所有已存储坐标维度执行精确的结构相等性比较
- `sindex` —— 空间索引(功能有限)
### CRS 操作
- `set_crs()` —— 设置坐标参考系
- `to_crs()` —— 在 CRS 之间转换
- `crs` —— 访问 CRS 信息
### 几何属性
- `area`、`length`、`bounds` —— 几何度量
- `is_valid`、`is_simple`、`is_empty` —— 几何校验
- `centroid`、`envelope`、`boundary` —— 几何属性
- `x`、`y`、`z`、`has_z` —— 坐标访问
- `total_bounds`、`estimate_utm_crs` —— 包围盒与 CRS 工具
- `hilbert_distance()` —— 基于几何包围盒中点生成分布式空间排序键
### 空间运算
- `buffer()` —— 几何缓冲
- `distance()` —— 距离计算
- `intersects()`、`contains()`、`within()` —— 空间谓词
- `intersection()` —— 几何相交
- `make_valid()` —— 几何校验与修复
- `is_valid_coverage()` —— 校验多边形覆盖的内部是否互不重叠、共享边界是否完全匹配
- `invalid_coverage_edges()` —— 按输入行返回导致覆盖无效的边,并可选择检测狭窄间隙
- `sample_points()` —— 使用原生分布式表达式按面积对多边形采样、按长度对线采样
- `GeoSeries.explode()` 与 `GeoDataFrame.explode()` —— 将多部件几何展开为多行,
其中 GeoDataFrame 方法会保留对应的属性列
- `cx` —— 基于坐标范围的空间筛选
- `clip()` —— 分布式几何裁剪
- `overlay()` —— GeoDataFrame 之间的分布式相交、差集、恒等、对称差和并集
- `sindex` —— 空间索引(功能有限)
```python
from shapely.geometry import box
coverage = gpd.GeoSeries([box(0, 0, 1, 1), box(1, 0, 2, 1)])
assert coverage.is_valid_coverage()
with_gap = gpd.GeoSeries([box(0, 0, 1, 1), box(1.1, 0, 2.1, 1)])
assert not with_gap.is_valid_coverage(gap_width=0.2)
invalid_edges = with_gap.invalid_coverage_edges(gap_width=0.2)
```
`invalid_coverage_edges()` 会让几何行保持分布式执行。Sedona 对按
`gap_width` 扩展的包围盒执行空间自连接来寻找候选邻居,聚合候选并在 executor
上执行校验,再将诊断结果映射回输入行。该方法
保留所有输入行、各级索引和 CRS;没有无效边或多边形分量的行返回空线,缺失
几何保持缺失。`is_valid_coverage()` 通过分布式归约返回一个 Python `bool`。
在高密度或大量重叠的覆盖中,或 `gap_width` 较大时,每个目标的候选数组和
executor 内存开销都会增加。`gap_width` 必须是有限的非负数。
`hilbert_distance()` 使用原生 Spark 表达式,并让逐行排序键保持分布式执行。
未提供 `total_bounds` 时,该方法会通过一次分布式聚合计算所有包围盒中点的
范围,只有这个有界摘要会返回 driver。由于 Spark 没有无符号整数类型,
与 GeoPandas 无符号 32 位值相同的结果会存储在 `int64` Series 中。
`overlay()` 会让两个输入始终保持分布式执行。候选对由 Sedona 空间连接
生成,差集分支则在 executor 上按源行聚合相匹配的掩膜几何。该实现不会收集
几何行、创建 Python UDF 或缓存中间候选对。构造结果时会先对两个完整输入
lineage 执行一次分布式的校验与元数据聚合;只有有界的摘要会返回 driver,
叠加几何行仍保持惰性求值。组合模式会有意执行多次空间连接,而不是隐式持久化
可能很大的候选对关系,从而以重复计算换取避免隐藏的缓存内存与生命周期成本。
输出行顺序不受保证。
与 GeoPandas 一样,每个输入必须只包含一个基本几何族,且默认修复无效多边形。
JTS 与 GEOS 对拓扑等价的结果可能采用不同的部件或坐标顺序。
Sedona 的 JTS 结构修复可能与 GeoPandas 的 GEOS 线网修复采用不同方式拆分
无效多边形;有效输入的拓扑应保持一致。`keep_geom_type=None` 的过滤行为与
`True` 相同,但 Sedona 不会仅为了在删除
低维几何时发出 GeoPandas 的条件警告而提前执行完整叠加。MultiIndex 列、
重复的单级列标签,以及会产生重复输出标签的属性后缀都会被明确拒绝。包括
MultiIndex 在内的输入行索引会被新的分布式索引替代。除 `difference` 外的
模式始终把活动输出列命名为 `geometry`,空结果与空间上不相交的结果也不例外,
从而避免 GeoPandas 的包围盒快速路径产生特殊命名。空的左输入会返回保留类型
的空结果,而不会像部分 GeoPandas 版本那样在访问第一个几何时抛出异常。
`identity` 遵循 GeoPandas 1.1+ 的 dtype 语义,保留左侧属性的 dtype,并提升
可空右侧属性的 dtype。`union` 和 `symmetric_difference` 即使某个逻辑差集
分支没有输出行,也会使用稳定的可空属性 dtype,从而避免仅为特化 dtype 而
触发提前执行。
`GeoSeries.explode()` 与 `GeoDataFrame.explode()` 使用 Sedona 原生的
`ST_Dump` 表达式与 Spark `posexplode`。几何部件及 GeoDataFrame 属性始终
保持分布式处理;这两个操作都不会收集源数据行,也不会使用 Python UDF。
为了保留与 GeoPandas 一致的行顺序和部件顺序并重建分布式索引,需要执行一次
全局排序与 shuffle;对于 `GeoDataFrame`,保留的属性列也会参与该 shuffle。
### 分布式几何聚合
- `GeoDataFrame.dissolve()` —— 对行进行分组,并使用 Sedona 原生分布式聚合
对每组几何执行并集
- `sedona.spark.geopandas.tools.collect()` —— 将分布式 GeoSeries 聚合为一个
几何,并在需要时使用同构多部件几何
```python
from sedona.spark.geopandas.tools import collect
by_region = buildings.dissolve(
by="region",
aggfunc={"population": "sum", "name": "first"},
)
all_building_parts = collect(buildings.geometry)
```
`dissolve` 支持 `unary` 并集方法,并会明确拒绝定点精度 `grid_size`、
`coverage` 与 `disjoint_subset` 并集。多个属性聚合会保留为 pandas-on-Spark
的两级 `MultiIndex`;GeoPandas 则把对应元组标签放在单级 object 索引中。
原生 `first`、`last` 与 `count` 聚合支持所有属性类型;`nunique` 支持数值、
布尔与字符串列;`min`、`max`、`sum`、`mean`、`median`、`std` 与 `var`
支持数值和布尔列。字符串 `sum` 与所有 `prod` 聚合都会被明确拒绝。可调用
聚合仅限 Python 内置的 `min`、`max`、`sum`,以及 NumPy 的 `amin`、
`amax`、`min`、`max`、`sum`、`mean`、`median`、`std` 与 `var`。
两个操作都在 Spark executor 上聚合输入行。`collect` 只会把聚合元数据和
该 API 所需的单个几何结果传输到 driver。
### 数据转换
- `to_geopandas()` —— 转换为传统 GeoPandas
- `GeoDataFrame.to_wkb()`、`GeoDataFrame.to_wkt()` —— 在分布式
pandas-on-Spark DataFrame 中把所有几何列序列化为 WKB/WKT
- `points_from_xy()`、`GeoSeries.from_xy()` —— 通过坐标列创建分布式
GeoSeries,且不收集分布式输入
- `geom_type` —— 获取几何类型
与 GeoPandas 不同,`GeoDataFrame.to_wkb()` 和 `GeoDataFrame.to_wkt()` 返回
惰性的 pandas-on-Spark DataFrame。仅在需要本地 pandas DataFrame 时,才对结果
调用 `.to_pandas()`。
`points_from_xy()` 会让 pandas-on-Spark 坐标 Series 保持分布式执行。本地
列表、NumPy 数组和 pandas 对象会先在 driver 上物化,因此大型坐标列应使用
分布式 Series。类似地,大型逐行 `sample_points()` 数量向量也应使用分布式
`size` Series。采样输出与逐行中间计算量会随请求的数量增长;当前线采样器会先
物化该行的采样点,再把它们收集为 MultiPoint。
## 完整工作流示例
```python
import sedona.spark.geopandas as gpd
from sedona.spark import SedonaContext
# 配置 Spark 进行匿名 S3 访问
config = (
SedonaContext.builder()
.config(
"spark.hadoop.fs.s3a.bucket.wherobots-examples.aws.credentials.provider",
"org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider",
)
.getOrCreate()
)
sedona = SedonaContext.create(config)
# 加载数据
DATA_DIR = "s3://wherobots-examples/data/geopandas_blog/"
overture_size = "1M"
postal_codes_path = DATA_DIR + "postal-code/"
overture_path = DATA_DIR + overture_size + "/" + "overture-buildings/"
postal_codes = gpd.read_parquet(postal_codes_path)
buildings = gpd.read_parquet(overture_path)
# 空间分析(GeoParquet 通常会保留 CRS 元数据)
if buildings.crs is None:
buildings = buildings.set_crs("EPSG:4326")
buildings_projected = buildings.to_crs("EPSG:3857")
# 计算面积并过滤
buildings_projected["area"] = buildings_projected.geometry.area
large_buildings = buildings_projected[buildings_projected["area"] > 1000]
result = large_buildings.sjoin(postal_codes, predicate="intersects")
# 按邮政编码聚合
summary = (
result.groupby("postal_code")
.agg({"area": "sum", "BUILD_ID": "count"})
.rename(columns={"BUILD_ID": "building_count"})
)
print(summary.head())
```
## 资源与贡献
完整且最新的 API 文档(包含方法签名、参数与示例)请参阅:
**📚 [GeoPandas API 文档](https://sedona.apache.org/latest/api/pydocs/sedona.spark.geopandas.html)**
Apache Sedona 的 GeoPandas API 是一个开源项目,欢迎参与贡献:可在 [GitHub issue tracker](https://github.com/apache/sedona/issues/2230) 报告 bug、提出新需求或贡献代码。更多贡献指南请参阅 [贡献者指南](../community/develop.md)。