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

未知结构JSON数组扁平化:转换为多行存入数据库方案

处理未知结构的JSON扁平化与数组展开

这个问题确实很头疼——当JSON结构不固定时,硬编码字段路径完全行不通,得用动态递归遍历的思路来解决。下面我给你分享一个基于Spark Scala的通用方案,不管JSON嵌套多深、数组出现在哪一层,都能自动完成扁平化,把数组的每个元素转换成单独一行,同时重复所有非数组字段的内容。

核心思路

  1. 递归遍历结构:自动识别JSON中的嵌套对象、数组字段和普通字段;
  2. 扁平化嵌套对象:把嵌套层级的字段提取到顶级(比如level.productReference.prodID直接变成prodID,或者保留路径前缀,看你需求);
  3. 展开数组字段:遇到数组时,用explode函数把数组拆分成多行,每一行对应数组的一个元素,同时保留当前所有非数组字段的值;
  4. 循环处理:递归执行上述步骤,直到所有字段都是普通数据类型(没有嵌套对象或数组)。

代码实现

先导入必要的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:12:57