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

Spark扁平化嵌套JSON后转换列时出现列不存在矛盾报错

问题:Spark扁平化后含点号的列无法被识别

场景重现

输入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 04:34:51