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

如何在PySpark DataFrame中移除含无效Polygon值的行?

问题描述

在PySpark DataFrame上执行操作时触发报错,推测是数据中存在异常行。

DataFrame结构

root
|-- geo_name: string (nullable = true)
|-- geo_latitude: double (nullable = true)
|-- geo_longitude: double (nullable = true)
|-- geo_bst: integer (nullable = true)
|-- geo_bvw: integer (nullable = true)
|-- geometry_type: string (nullable = true)
|-- geometry_polygon: string (nullable = true)
|-- geometry_multipolygon: string (nullable = true)
|-- polygon: geometry (nullable = false)

转换代码

通过以下代码将CSV中的geometry_polygon列转换为geometry类型的polygon列:

station_groups_gdf.createOrReplaceTempView("station_gdf")
spatial_station_groups_gdf = spark_sedona.sql("SELECT *, ST_PolygonFromText(station_gdf.geometry_polygon, ',') AS polygon FROM station_gdf")

示例数据

-RECORD 0-------------------------------------
geo_name              | Neckarkanal         
 geo_latitude          | 49.486697           
 geo_longitude         | 8.504944            
 geo_bst               | 0                    
 geo_bvw               | 0                    
 geometry_type         | Polygon              
 geometry_polygon      | 8.4937, 49.4892, ...
 geometry_multipolygon | null                
 polygon               | POLYGON ((8.4937 ...  

报错信息

仅调用df.show()时触发错误:

java.lang.IllegalArgumentException: Points of LinearRing do not form a closed linestring

需求

希望定位并删除这些无效行,尝试过类似dataframe.where(dataframe.polygon == valid).show()的写法,但需避免全量加载DataFrame导致报错,求最优实现方式。


最优实现方案

核心思路是在生成polygon列的阶段就完成有效性校验与过滤,避免全量加载数据时触发异常。

方法1:结合TRY与ST_IsValid函数(推荐)

利用Spark SQL的TRY函数捕获ST_PolygonFromText的转换异常,再通过Sedona的ST_IsValid校验几何有效性,最后过滤出有效行:

station_groups_gdf.createOrReplaceTempView("station_gdf")

spatial_station_groups_gdf = spark_sedona.sql("""
    SELECT *
    FROM (
        SELECT 
            *, 
            TRY(ST_PolygonFromText(station_gdf.geometry_polygon, ',')) AS polygon
        FROM station_gdf
    ) t
    WHERE polygon IS NOT NULL AND ST_IsValid(polygon)
""")
  • TRY会在转换失败时返回NULL,避免直接抛出异常中断任务
  • ST_IsValid确保生成的polygon符合几何规范(比如闭合的LinearRing)
  • 子查询结构保证转换与校验都在分布式计算阶段完成,不会全量加载数据

方法2:自定义UDF处理异常(适配低版本Sedona)

如果Sedona版本不支持TRY函数,可自定义UDF捕获转换异常并标记有效行:

from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType

def is_valid_polygon(wkt_str):
    try:
        polygon = spark_sedona.ST_PolygonFromText(wkt_str, ',')
        return spark_sedona.ST_IsValid(polygon)
    except:
        return False

is_valid_udf = udf(is_valid_polygon, BooleanType())

# 先标记有效行,过滤后再生成polygon列
valid_stations_df = station_groups_gdf.withColumn("is_valid", is_valid_udf("geometry_polygon"))
spatial_station_groups_gdf = valid_stations_df.filter("is_valid == true") \
    .withColumn("polygon", spark_sedona.ST_PolygonFromText("geometry_polygon", ','))
  • 先过滤无效行再生成polygon,避免无效数据触发转换异常
  • UDF需覆盖空值、格式错误等所有可能的异常场景

方法3:定位无效行(用于排查)

如果需要先定位异常行的具体内容,可单独筛选无效记录:

station_groups_gdf.createOrReplaceTempView("station_gdf")

invalid_rows_df = spark_sedona.sql("""
    SELECT *
    FROM (
        SELECT 
            *, 
            TRY(ST_PolygonFromText(station_gdf.geometry_polygon, ',')) AS polygon
        FROM station_gdf
    ) t
    WHERE polygon IS NULL OR NOT ST_IsValid(polygon)
""")
  • 可查看这些行的geometry_polygon内容,分析异常原因(比如首尾点不重合、坐标格式错误等)

内容的提问来源于stack exchange,提问作者pi_janes

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:15:30