Spark Scala中基于Apache Sedona实现坐标空间相交查询
解决方案:使用Apache Sedona实现经纬度与GeoJSON区域的空间匹配
步骤1:加载并预处理源数据
首先处理CSV格式的源数据,统一字段结构、转换经纬度格式,并生成Sedona可识别的点几何对象。
源数据格式
1111:150458,025.22826N,055.30022E,348,39,JOB_ONBOARD 2222:150448,025.22746N,055.29962E,32,48,CAR_AVAILABLE 3333,20072023:150612,025.30559N,055.38272E,130,50,CAR_AVAILABLE 4444,20072023:150740,025.21794N,055.28569E,0,0,JOB_ONBOARD
加载源数据代码
import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ // 定义兼容两种字段格式的Schema val sourceSchema = StructType(Seq( StructField("id", StringType), StructField("datetime", StringType, nullable = true), StructField("lat_str", StringType), StructField("lon_str", StringType), StructField("col1", IntegerType), StructField("col2", IntegerType), StructField("status", StringType) )) // 加载并预处理源数据 val sourceDF = spark.read .option("delimiter", ",") .option("header", "false") .schema(sourceSchema) .csv("path/to/your/source_data.csv") // 统一id和datetime字段格式 .withColumn("id", when(col("datetime").isNull, col("id")).otherwise(col("id"))) .withColumn("datetime", when(col("datetime").isNull, lit(null)).otherwise(col("datetime"))) // 将经纬度字符串转换为数值(处理N/S/E/W标识) .withColumn("latitude", expr("cast(regexp_replace(lat_str, '[N|S]', '') as double) * case when lat_str like '%S' then -1 else 1 end")) .withColumn("longitude", expr("cast(regexp_replace(lon_str, '[E|W]', '') as double) * case when lon_str like '%W' then -1 else 1 end")) // 生成WGS84坐标系的点几何对象 .withColumn("point", expr("ST_Point(longitude, latitude)"))
步骤2:正确加载GeoJSON数据
完善GeoJSON解析逻辑,提取区域标识信息(如社区名称)并转换为Sedona几何对象:
val inputLocation = "C:\\communities_0.geojson" // 加载GeoJSON并解析几何与属性 val communitiesDF = spark.read .option("multiline", "true") // 适配多行格式的GeoJSON .json(inputLocation) .selectExpr("explode(features) as features") .select( "features.properties.community_name", // 替换为你的GeoJSON中实际的区域名称字段(如name/community) expr("ST_GeomFromGeoJSON(features.geometry) as polygon") ) // 若GeoJSON坐标系非WGS84,需转换:.withColumn("polygon", expr("ST_Transform(polygon, 'EPSG:4326')"))
步骤3:执行空间相交匹配
使用ST_Intersects函数完成点与区域的空间连接,得到最终匹配结果:
// 空间连接匹配所属区域 val resultDF = sourceDF.join( communitiesDF, expr("ST_Intersects(sourceDF.point, communitiesDF.polygon)"), "left_outer" // 保留所有源数据行,即使不在任何区域内 ) // 输出符合预期格式的结果 resultDF.select( when(col("datetime").isNull, col("id")).otherwise(concat(col("id"), ",", col("datetime"))).alias("id_datetime"), col("lat_str"), col("lon_str"), col("col1"), col("col2"), col("status"), col("community_name").alias("community") ).write .option("delimiter", ",") .mode("overwrite") .csv("path/to/your/result.csv")
关键注意事项
- 坐标系一致性:确保点和多边形使用同一坐标系(GeoJSON默认WGS84/EPSG:4326),不一致时用
ST_Transform转换。 - 字段适配:
community_name需替换为你的GeoJSON properties中实际的区域名称字段。 - 源数据兼容:处理前两行与后两行字段数不一致的问题,统一Schema后再进行后续操作。
内容的提问来源于stack exchange,提问作者rahul kumar
相关产品推荐
相关产品推荐

