Spark Scala解析复杂嵌套JSON转为DataFrame的问题求助
嵌套JSON完全展平为Spark DataFrame的实现方案
现有代码的问题
- 你当前实现的
flattenStructSchema函数仅支持递归展开Struct类型嵌套,未处理Array类型字段,因此data(二维字符串数组)、meta.view.approvals(结构体数组)这类数组类型字段无法被展开,会被保留为原复杂类型。 - 拆分
meta和data字段分别处理会丢失两者的行关联关系,若业务需要保留元数据和业务数据的对应关系,不建议拆分处理。
通用全嵌套展平实现方案
你可以使用支持同时处理Struct、Array嵌套的递归展平函数,实现全层级自动展开:
import org.apache.spark.sql.Column import org.apache.spark.sql.functions.{col, explode_outer} import org.apache.spark.sql.types.{ArrayType, StructType} def flattenSchema(schema: StructType, prefix: String = null, explodeArray: Boolean = true): Array[Column] = { schema.fields.flatMap(f => { val fullColName = if (prefix == null) f.name else s"$prefix.${f.name}" f.dataType match { // 递归处理Struct嵌套结构 case st: StructType => flattenSchema(st, fullColName, explodeArray) // 处理Array类型:先展开再递归处理数组内元素 case at: ArrayType if explodeArray => val explodedCol = explode_outer(col(fullColName)).as(fullColName) at.elementType match { case st: StructType => flattenSchema(st, fullColName, explodeArray) case innerAt: ArrayType => flattenSchema(new StructType().add(f.name, innerAt), prefix, explodeArray) case _ => Array(explodedCol.as(fullColName.replace(".", "_"))) } // 基础类型直接返回,替换列名中的点为下划线避免语法冲突 case _ => Array(col(fullColName).as(fullColName.replace(".", "_"))) } }) }
业务场景使用示例
// 1. 完整展平整个JSON的所有嵌套结构,保留行关联 val fullyFlattenedDf = df.select(flattenSchema(df.schema):_*) // 2. 仅展开Struct嵌套、保留Array字段不拆分的调用方式 val structOnlyFlattenedDf = df.select(flattenSchema(df.schema, explodeArray = false):_*) // 3. 单独处理二维数组data字段、转为多行单值的特殊场景实现 val dataExplodedDf = df.select(explode_outer(col("data")).as("data_row")) .select(explode_outer(col("data_row")).as("data_single_value"))
注意事项
- 如果不需要保留空数组对应的行,可以将
explode_outer替换为explode提升性能 - 如果你使用的是Socrata平台导出的标准JSON格式,
data二维数组的取值顺序和meta.view.columns中定义的列名一一对应,你可以先提取列名列表,再通过withColumn批量将data数组转为多列,避免多次explode破坏关联关系。
内容的提问来源于stack exchange,提问作者Roxane Felton
相关产品推荐
相关产品推荐

