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

如何在Apache Spark Scala中将含数组的嵌套JSON扁平化为单行DataFrame

实现思路

Spark默认的explode方法会将数组元素拆分为多行,要实现单行扁平化,核心是不拆分数组,而是将数组的每个索引位作为单独列生成,把层级路径+索引作为最终列名,通过递归遍历DataFrame的Schema结构自动生成所有扁平化列的表达式即可。

完整实现代码
import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.functions.{col, size}
import org.apache.spark.sql.types._

// 递归生成扁平化的列表达式与对应别名
def flattenSchema(df: DataFrame, schema: StructType, prefix: String = ""): Seq[(String, String)] = {
  schema.fields.flatMap { field =>
    val colFullPath = if (prefix.isEmpty) field.name else s"${prefix}_${field.name}"
    field.dataType match {
      // 嵌套结构体类型,递归处理内部字段
      case st: StructType => flattenSchema(df, st, colFullPath)
      // 数组类型,按索引逐个提取元素
      case arr: ArrayType =>
        val maxArrLen = getArrayMaxLength(colFullPath, df)
        arr.elementType match {
          // 数组元素为结构体,递归处理结构体字段,拼上当前索引
          case st: StructType =>
            (0 until maxArrLen).flatMap(idx => flattenSchema(df, st, s"${colFullPath}_$idx"))
          // 数组元素为基础类型,直接按索引取元素生成列
          case _ =>
            (0 until maxArrLen).map(idx => (s"$colFullPath[$idx]", s"${colFullPath}_$idx"))
        }
      // 基础类型直接返回对应列
      case _ => Seq((colFullPath, colFullPath))
    }
  }
}

// 辅助方法:计算指定数组列的全局最大长度,适配多数据场景
def getArrayMaxLength(colPath: String, df: DataFrame): Int = {
  df.select(size(col(colPath))).as[Int].reduce(math.max)
}
使用示例
// 1. 读取JSON数据,可替换为spark.read.json("文件路径")读本地/分布式JSON文件
val jsonStr = """{
 "name":"John",
 "age":30,
 "bike":{
    "name":"Bajaj", "models":["Dominor", "Pulsar"]
    },
 "cars": [
   { "name":"Ford", "models":[ "Fiesta", "Focus", "Mustang" ] },
   { "name":"BMW", "models":[ "320", "X3", "X5" ] },
   { "name":"Fiat", "models":[ "500", "Panda" ] }
 ]
}"""
val rawDf = spark.read.json(Seq(jsonStr).toDS())

// 2. 生成所有扁平化列的选择表达式
val flatSelectCols = flattenSchema(rawDf, rawDf.schema).map {
  case (expr, alias) => col(expr).alias(alias)
}

// 3. 执行选择得到最终单行扁平化DataFrame
val flatDf = rawDf.select(flatSelectCols: _*)

// 验证输出
flatDf.show()
输出说明

最终生成的DataFrame列名完全符合要求,示例输出列包括:

  • name、age
  • bike_name、bike_models_0、bike_models_1
  • cars_0_name、cars_0_models_0、cars_0_models_1、cars_0_models_2
  • cars_1_name、cars_1_models_0、cars_1_models_1、cars_1_models_2
  • cars_2_name、cars_2_models_0、cars_2_models_1

所有数据保留在同一行,不会拆分出多条记录。如果是批量处理多组JSON数据,数组长度不足的索引位会自动填充null值,无需额外适配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 10:27:02