如何在Scala/Spark中实现EPSG坐标转换:东向/北向转经纬度
Scala/Spark 实现EPSG:27700转EPSG:4326坐标转换
Spark本身没有内置坐标转换能力,以下两种方案可实现你需要的坐标转换需求:
方案一:使用GeoSpark(推荐,适配分布式处理)
GeoSpark是Spark生态的地理空间处理库,原生支持坐标参考系转换,适合大规模数据场景。
1. 添加依赖
SBT 依赖
libraryDependencies += "org.datasyslab" % "geospark-sql_2.12" % "1.4.1" libraryDependencies += "org.datasyslab" % "geospark" % "1.4.1"
Maven 依赖
<dependency> <groupId>org.datasyslab</groupId> <artifactId>geospark-sql_2.12</artifactId> <version>1.4.1</version> </dependency> <dependency> <groupId>org.datasyslab</groupId> <artifactId>geospark</artifactId> <version>1.4.1</version> </dependency>
2. 实现代码
import org.apache.spark.sql.SparkSession import org.datasyslab.geosparksql.utils.GeoSparkSQLRegistrator object CoordinateTransform { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("EPSGTransform") .master("local[*]") // 生产环境请移除该配置 .getOrCreate() // 注册GeoSpark SQL扩展函数 GeoSparkSQLRegistrator.registerAll(spark) // 构造测试数据 val rawData = Seq( ("94489", 276164, 84185), ("94555", 428790, 92790), ("94806", 357501, 173246), ("99118", 439545, 336877), ("76202", 357353, 170708) ).toDF("node_id", "easting", "northing") // 将东/北向坐标转为EPSG:27700的Point类型 rawData.createOrReplaceTempView("raw_coords") val pointDF = spark.sql( """ |SELECT node_id, easting, northing, | ST_Point(easting, northing) AS point_27700 |FROM raw_coords |""".stripMargin ) // 转换坐标到EPSG:4326,提取经纬度 pointDF.createOrReplaceTempView("point_coords") val resultDF = spark.sql( """ |SELECT node_id, easting, northing, | ST_X(ST_Transform(point_27700, 'EPSG:27700', 'EPSG:4326')) AS longitude, | ST_Y(ST_Transform(point_27700, 'EPSG:27700', 'EPSG:4326')) AS latitude |FROM point_coords |""".stripMargin ) // 输出结果 resultDF.show(false) spark.stop() } }
方案二:自定义UDF结合Proj4j库
若不想引入GeoSpark,可基于Proj4j库封装坐标转换逻辑为Spark UDF。
1. 添加依赖
SBT 依赖
libraryDependencies += "org.osgeo" % "proj4j" % "1.1.0"
Maven 依赖
<dependency> <groupId>org.osgeo</groupId> <artifactId>proj4j</artifactId> <version>1.1.0</version> </dependency>
2. 实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{udf, col} import org.osgeo.proj4j.{CRSFactory, CoordinateTransform, CoordinateTransformFactory, ProjCoordinate} object Proj4jCoordinateTransform { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("Proj4jTransform") .master("local[*]") // 生产环境请移除该配置 .getOrCreate() // 初始化Proj4j转换实例 val ctFactory = new CoordinateTransformFactory() val crsFactory = new CRSFactory() val sourceCRS = crsFactory.createFromName("EPSG:27700") val targetCRS = crsFactory.createFromName("EPSG:4326") val transform: CoordinateTransform = ctFactory.createTransform(sourceCRS, targetCRS) // 定义坐标转换UDF val transformUdf = udf((easting: Double, northing: Double) => { val srcCoord = new ProjCoordinate(easting, northing) val targetCoord = new ProjCoordinate() transform.transform(srcCoord, targetCoord) (targetCoord.x, targetCoord.y) }) // 构造测试数据 val rawData = Seq( ("94489", 276164.0, 84185.0), ("94555", 428790.0, 92790.0), ("94806", 357501.0, 173246.0), ("99118", 439545.0, 336877.0), ("76202", 357353.0, 170708.0) ).toDF("node_id", "easting", "northing") // 应用UDF并提取经纬度 val resultDF = rawData.withColumn("lat_lon", transformUdf(col("easting"), col("northing"))) .select( col("node_id"), col("easting"), col("northing"), col("lat_lon._1").alias("longitude"), col("lat_lon._2").alias("latitude") ) // 输出结果 resultDF.show(false) spark.stop() } }
最终输出结果
两种方案都会生成与你Python代码一致的结果:
| node_id | easting | northing | longitude | latitude |
|---|---|---|---|---|
| 94489 | 276164 | 84185 | -3.752811 | 50.644154 |
| 94555 | 428790 | 92790 | -1.593413 | 50.734016 |
| 94806 | 357501 | 173246 | -2.613059 | 51.456587 |
| 99118 | 439545 | 336877 | -1.413188 | 52.927852 |
| 76202 | 357353 | 170708 | -2.614883 | 51.433757 |
内容的提问来源于stack exchange,提问作者Suhani Bhatia
相关产品推荐
相关产品推荐

