如何在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
相关产品推荐
相关产品推荐

