如何在PySpark DataFrame中利用经纬度创建Shapely Point列
解决方案
方法一:结合Shapely与PySpark UDF生成Point对象
如果后续需要基于Shapely进行空间运算,可通过UDF实现,但需注意分布式环境下的序列化问题,且集群所有节点需提前安装shapely库。
步骤实现:
- 初始化SparkSession并加载数据
- 定义生成Shapely Point的函数(注意Shapely的Point参数顺序为
(经度, 纬度),即(Lon, Lat)) - 注册UDF并添加新列
from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from shapely.geometry import Point from pyspark.sql.types import StringType # 初始化SparkSession spark = SparkSession.builder.appName("ShapelyPointGenerator").getOrCreate() # 构造示例数据 sample_data = [(1, 20.2, 78.3), (2, 20.3, 78.4), (3, 20.4, 78.4), (4, 20.5, 78.4), (5, 20.6, 78.4)] df = spark.createDataFrame(sample_data, schema=["TS", "Lat", "Lon"]) # 定义生成Point的函数,转为WKT字符串便于Spark存储 def build_shapely_point(lat, lon): return Point(lon, lat).wkt # 转为WKT字符串,避免序列化问题 # 注册UDF point_udf = udf(build_shapely_point, StringType()) # 添加Point列 result_df = df.withColumn("Point", point_udf(col("Lat"), col("Lon"))) # 查看结果 result_df.show(truncate=False)
进阶:自定义UDT存储原生Shapely对象
若需保留原生Shapely Point对象而非字符串,可自定义UserDefinedType(UDT)实现序列化:
from pyspark.sql.types import UserDefinedType from shapely.wkb import dumps, loads class PointUDT(UserDefinedType): def sqlType(self): return StringType() def serialize(self, obj): return dumps(obj) def deserialize(self, datum): return loads(datum) # 重新注册UDF使用自定义UDT point_udf = udf(build_shapely_point, PointUDT()) result_df = df.withColumn("Point", point_udf(col("Lat"), col("Lon")))
方法二:使用PySpark内置函数生成Point格式字符串(无额外依赖)
若仅需生成Point(xxx yyy)格式的字符串,无需依赖Shapely,直接用内置字符串拼接函数即可,性能更优:
from pyspark.sql.functions import concat, lit, col result_df = df.withColumn( "Point", concat(lit("Point("), col("Lat"), lit(" "), col("Lon"), lit(")")) ) result_df.show(truncate=False)
方案选择建议
- 若需后续空间运算:选方法一,确保集群安装Shapely,优先用WKT字符串存储(序列化更稳定)
- 仅需格式化输出:选方法二,无依赖、性能高,适合大数据场景
内容的提问来源于stack exchange,提问作者Sameer J
相关产品推荐
相关产品推荐

