未知结构JSON数组扁平化:转换为多行存入数据库方案
处理未知结构的JSON扁平化与数组展开
这个问题确实很头疼——当JSON结构不固定时,硬编码字段路径完全行不通,得用动态递归遍历的思路来解决。下面我给你分享一个基于Spark Scala的通用方案,不管JSON嵌套多深、数组出现在哪一层,都能自动完成扁平化,把数组的每个元素转换成单独一行,同时重复所有非数组字段的内容。
核心思路
- 递归遍历结构:自动识别JSON中的嵌套对象、数组字段和普通字段;
- 扁平化嵌套对象:把嵌套层级的字段提取到顶级(比如
level.productReference.prodID直接变成prodID,或者保留路径前缀,看你需求); - 展开数组字段:遇到数组时,用
explode函数把数组拆分成多行,每一行对应数组的一个元素,同时保留当前所有非数组字段的值; - 循环处理:递归执行上述步骤,直到所有字段都是普通数据类型(没有嵌套对象或数组)。
代码实现
先导入必要的Spark依赖包:
import org.apache.spark.sql.{DataFrame, SparkSession} import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._
然后写一个通用的扁平化函数:
def flattenDataFrame(df: DataFrame): DataFrame = { val fields = df.schema.fields val flatFields = fields.flatMap { field => field.dataType match { // 处理嵌套对象:提取子字段到顶级 case structType: StructType => structType.fields.map { nestedField => // 如果需要保留字段路径,把别名改成 s"${field.name}_${nestedField.name}" col(s"${field.name}.${nestedField.name}").alias(nestedField.name) } // 暂时保留数组字段,后面单独处理 case _: ArrayType => Seq(col(field.name)) // 普通字段直接保留 case _ => Seq(col(field.name)) } } val flattened = df.select(flatFields: _*) // 检查是否还有嵌套结构,递归处理 val hasNested = flattened.schema.fields.exists(_.dataType.isInstanceOf[StructType]) // 检查是否有数组字段,需要展开 val hasArray = flattened.schema.fields.exists(_.dataType.isInstanceOf[ArrayType]) hasNested match { case true => flattenDataFrame(flattened) case false if hasArray => // 提取所有数组字段,依次展开(多个数组会产生笛卡尔积,符合需求) val arrayCols = flattened.schema.fields.filter(_.dataType.isInstanceOf[ArrayType]).map(_.name) val exploded = arrayCols.foldLeft(flattened) { (tempDf, colName) => tempDf.withColumn(colName, explode(col(colName))) } // 展开数组后可能还有嵌套对象,继续递归扁平化 flattenDataFrame(exploded) case _ => flattened } }
使用示例
用你提供的JSON数据来测试:
// 初始化SparkSession val spark = SparkSession.builder() .appName("DynamicJsonFlattener") .master("local[*]") // 生产环境去掉这个,用集群配置 .getOrCreate() // 加载示例JSON数据 val jsonSample = """ { "level":{"productReference":{ "prodID":"1234", "unitOfMeasure":"EA" }, "states":[ { "state":"SELL", "effectiveDateTime":"2015-10-09T00:55:23.6345Z", "stockQuantity":{ "quantity":1400.0, "stockKeepingLevel":"A" } }, { "state":"HELD", "effectiveDateTime":"2015-10-09T00:55:23.6345Z", "stockQuantity":{ "quantity":800.0, "stockKeepingLevel":"B" } } ] }} """ val rawDf = spark.read.json(spark.sparkContext.parallelize(Seq(jsonSample))) // 执行扁平化和数组展开 val resultDf = flattenDataFrame(rawDf) // 查看结果 resultDf.show()
执行后你会得到和示例完全一致的输出:
+------+-------------+-----+-------------------------+--------+-----------------+ |prodID|unitOfMeasure|state|effectiveDateTime |quantity|stockKeepingLevel| +------+-------------+-----+-------------------------+--------+-----------------+ |1234 |EA |SELL |2015-10-09T00:55:23.6345Z|1400.0 |A | |1234 |EA |HELD |2015-10-09T00:55:23.6345Z|800.0 |B | +------+-------------+-----+-------------------------+--------+-----------------+
额外说明
- 字段路径保留:如果需要区分不同层级的同名字段(比如两个嵌套对象都有
id),可以修改扁平化时的别名规则,比如把alias(nestedField.name)改成alias(s"${field.name}_${nestedField.name}"); - 多数组处理:如果JSON中有多个数组字段,这个函数会依次展开,产生笛卡尔积(这是符合需求的,因为每个数组元素都要和其他字段组合);
- 性能优化:这个方案基于Spark DataFrame API,Spark的 Catalyst 优化器会自动处理底层的执行计划,比直接用RDD更高效。
内容的提问来源于stack exchange,提问作者John
相关产品推荐
相关产品推荐

