如何在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
相关产品推荐
相关产品推荐

