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

Spark Scala:如何修改嵌套Struct结构,将经纬度转为数组格式

在Scala中高效转换Spark嵌套Struct为指定JSON格式

问题背景

现有嵌套Struct结构的Spark DataFrame,需要将每个城市的lat和long从键值对格式{"lat": x, "long": y}转换为数组格式[x, y],最终输出指定结构的JSON。原to_json输出无法满足需求,且国家数量、每个国家的城市数量不固定,需要动态处理。

解决方案

利用Spark内置函数递归遍历嵌套Struct字段,动态转换目标结构,避免硬编码字段名,同时保证执行效率。

步骤1:导入依赖包

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

步骤2:定义递归转换函数

该函数会自动识别嵌套Struct:

  • 若遇到包含lat和long的城市Struct,将其转换为数组
  • 若遇到包含多个城市的国家Struct,递归处理每个城市字段
  • 其他类型字段直接保留原格式
def transformCityStruct(structField: StructField): Column = {
  structField.dataType match {
    case structType: StructType if structType.fieldNames.contains("lat") && structType.fieldNames.contains("long") =>
      array(col(s"${structField.name}.lat"), col(s"${structField.name}.long"))
    case structType: StructType =>
      val cityCols = structType.fields.map(f => transformCityStruct(f).alias(f.name))
      struct(cityCols: _*)
    case _ => col(structField.name)
  }
}

步骤3:转换DataFrame并生成JSON

假设原DataFrame名为df,执行以下代码完成转换:

// 获取Country字段的结构定义
val countryStruct = df.schema("Country").dataType.asInstanceOf[StructType]
// 构造转换后的Country列
val transformedCountry = struct(
  countryStruct.fields.map(f => transformCityStruct(f).alias(f.name)): _*
)
// 生成新DataFrame并输出JSON
val resultDF = df.withColumn("Country", transformedCountry)
val targetJson = resultDF.select(to_json(col("Country"))).head().getString(0)

优势说明

  • 动态适配:无需硬编码国家、城市字段名,自动适配任意数量的国家和城市
  • 高效执行:基于Spark内置函数实现,避免UDF的性能损耗,Spark可对执行计划进行优化

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 16:05:41