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

Scala中Spark含特殊字符与斜杠的JSON DataFrame扁平化报错解决

问题解决:Spark处理含特殊字符(斜杠/等)的JSON扁平化

问题场景

读取包含特殊字符(%、-、+、&、/)的JSON数据并转为Spark DataFrame,使用自定义扁平化函数时抛出如下错误:

AnalysisException: Column 'ExportData_Country_Country' does not exist.
Did you mean one of the following? [ExportData_Country_Country/Code,
ExportSchema_CountryCode, ExportData_Country_MaterialExternal,
ExportData_Country_MaterialExternal, ExportSchema_Application,
ExportSchema_EntityType, ExportSchema_ExportType,
ExportSchema_FileName, ExportSchema_ExportUtcDateTime,
ExportSchema_LastExportUtcDateTime]; line 1 pos 0; 'Project
[ExportSchema_Application#52091, ExportSchema_CountryCode#52092,
ExportSchema_EntityType#52093, ExportSchema_ExportType#52094,
ExportSchema_ExportUtcDateTime#52095, ExportSchema_FileName#52096,
ExportSchema_LastExportUtcDateTime#52097,
unresolvedalias(('ExportData_Country_Country / 'Code),
Some(org.apache.spark.sql.Column$$Lambda$6738/2007565455@225acd0)),
ExportData_Country_MaterialExternal#52112]
+- Generate explode(ExportData_Country_MaterialExternal#52099), true, [ExportData_Country_MaterialExternal#52112]

核心原因是字段名中的斜杠/未被正确转义,Spark将其解析为除法运算符而非字段名的一部分。

修改后的完整代码

val jsonString_new = """
{"ExportSchema":[{"Application":"fuzzitDwh","EntityType":"4","CountryCode":"AT","ExportType":"DIFF","ExportUtcDateTime":"2021-12-15T01:45:01","LastExportUtcDateTime":"2021-12-14T10:19:37","FileName":"UL_DataLake_fuzzitDwh_MaterialExternal_AT_DIFF_20211215_014501.json"}],"ExportData":[{"Country/":[{"Country/Code":"AT","MaterialExternal":[{"WholesalerCode":" A__METR","M%at-erial()+ ExternalCode":"26780","Attributes":{"Name":"KALBSFOND HELL PASTOES 1KG KNR"}},{"WholesalerCode":"A__ESVP","M%at-erial()+ ExternalCode":"38636","Attributes":{"MaterialInternalCodeFromWholesaler":"38636","Name":"B&J 465ml Netflix&Chill'd CL1x8x180EB"}},{"WholesalerCode":"A__ESVP","M%at-erial()+ ExternalCode":"39767","Attributes":{"MaterialInternalCodeFromWholesaler":" 39767","Name":"B&J 170G ChocChipCookieDghCL3b X8x147EB"}}]}]}]}
""".stripMargin

val json_df = spark.read.json(Seq(jsonString_new).toDS) 

def flatten_json(df: DataFrame): DataFrame = {
  val fields = df.schema.fields
  val fieldNames = fields.map(x => x.name)

  for (i <- fields.indices) {
    val field = fields(i)
    val fieldType = field.dataType
    val fieldName = field.name
    fieldType match {
      case _: ArrayType =>
        val fieldNamesExcludingArray = fieldNames.filter(_ != fieldName)
        // 用反引号包裹数组字段名,避免特殊字符被误解析
        val fieldNamesAndExplode = fieldNamesExcludingArray ++ Array(
          s"explode_outer(`$fieldName`) as `$fieldName`"
        )
        val explodedDf = df.selectExpr(fieldNamesAndExplode: _*)
        return flatten_json(explodedDf)
      case structType: StructType =>
        // 构造子字段引用时,父字段和子字段都用反引号包裹
        val childFieldNames = 
          structType.fieldNames.map(childname => s"`$fieldName`.`$childname`")
        val newFieldNames = fieldNames.filter(_ != fieldName) ++ childFieldNames
        import org.apache.spark.sql.functions.col

        val renamedCols = 
          newFieldNames.map { x =>
            // 将字段名中的特殊字符替换为下划线,生成合法列名
            val newName = x.replace("`", "").replace(".", "_").replace("/", "_")
            col(x).as(newName)
          }

        val explodedDf = df.select(renamedCols: _*)
        return flatten_json(explodedDf)
      case _ =>
    }
  }

  df
}
val flattened_Df = flatten_json(json_df)
flattened_Df.show(false)

关键修改点

  1. 反引号包裹字段名:所有含特殊字符的列名必须用反引号()包裹,明确告诉Spark将其作为完整字段名解析,而非运算符。比如将explode_outer($fieldName)改为explode_outer($fieldName),构造子字段引用时使用s"$fieldName.$childname"`格式。
  2. 清理特殊字符生成合法列名:重命名时将字段名中的.、/替换为下划线_,避免后续操作中出现解析问题。例如Country/Code会被重命名为Country_Code。

运行修改后的代码即可正确扁平化含特殊字符的JSON数据,不会再抛出列不存在的错误。

内容的提问来源于stack exchange,提问作者Dayita Chakrabarty

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 03:49:58