Spark扁平化嵌套JSON后转换列时出现列不存在矛盾报错
场景重现
输入JSON
{ "order": { "name": "John Doe", "address": { "state": "NY" }, "orders": [ { "order_id": 1001, "quarter": "2023-06-30" }, { "order_id": 1002, "quarter": "2023-06-29" } ] } }
扁平化后结果
+-------------------+----------+---------------------+------------------------+ |order.address.state|order.name|order.orders.order_id|order.orders.quarter | +-------------------+----------+---------------------+------------------------+ |NY |John Doe |[1001, 1002] |[2023-06-30, 2023-06-29]| |NY |John Doe |[1001, 1002] |[2023-06-30, 2023-06-29]| +-------------------+----------+---------------------+------------------------+
错误信息
Exception in thread "main" org.apache.spark.sql.AnalysisException: Column 'order.address.state' does not exist. Did you mean one of the following? [order.address.state, order.orders.quarter, order.name, order.orders.order_id];
问题根源
Spark的列名解析逻辑会默认把包含.的字符串解析为嵌套结构体路径,比如functions.col("order.address.state")会被解析为:查找名为order的结构体列,再找其下的address结构体,最后找state字段。但扁平化后,order.address.state是一个独立的顶层列名,并非嵌套结构,因此Spark找不到对应的嵌套路径,抛出矛盾错误。
解决方案
方法1:用反引号包裹列名
在引用含点号的列时,用反引号()包裹完整列名,明确告诉Spark这是单个列名而非嵌套路径: 修改applyTransformation`方法中的列引用代码:
input = input.withColumn(columnName, functions.concat(functions.col("`" + columnName + "`"), functions.lit("_suffix")));
方法2:扁平化时替换列名中的点号
在递归扁平化过程中,将列名里的.替换为无歧义字符(比如_),避免触发Spark的嵌套解析逻辑:
修改flattenNestedJson方法中生成列名的代码:
// 原代码 String fullColumnName = prefix.isEmpty() ? columnName : prefix + "." + columnName; // 修改后 String fullColumnName = prefix.isEmpty() ? columnName : prefix + "_" + columnName;
后续转换、反扁平化操作统一使用替换后的列名即可。
方法3:使用Column.quoted方法(Spark 3.x+)
Spark 3.x及以上版本支持Column.quoted()方法,自动处理含特殊字符的列名:
input = input.withColumn(columnName, functions.concat(functions.col(columnName).quoted(), functions.lit("_suffix")));
额外注意
在unflattenJson方法中,同样需要用反引号包裹含点号的列名:
Dataset<Row> withStructColumn = input.withColumn("value", functions.struct( input.col("`order.address.state`").alias("state"), input.col("`order.name`").alias("name"), input.col("`order.orders.order_id`").alias("order_id"), input.col("`order.orders.quarter`").alias("quarter") ));
内容的提问来源于stack exchange,提问作者Ajith Kannan

