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

PySpark SQL中UNNEST与SPLIT函数报错问题排查

问题描述
  • 同一段查询在Athena(Presto语法)中可正常运行,但在AWS Glue的Python Spark SQL执行时,抛出AnalysisException: Column 'metric_value' does not exist错误,该列实际存在于原表中。
  • 仅当查询包含CROSS JOIN UNNEST语句时触发报错,移除该语句后任务可正常执行。
  • metric_value列示例值:{14311, 242342134, 13132}
  • 原查询代码:
df1 = spark.sql("""
select * FROM database.table
WHERE date(load_date) = current_date - interval '1' day
    """)

df1.createOrReplaceTempView("table1")
    
df2 = spark.sql("""  
SELECT load_date, load_hr, timestamp,
ne_name, object, metric_value, 
MIN(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS MIN_UTIL, 
MAX(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS MAX_UTIL, 
AVG(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS AVG_UTIL
FROM table1
cross join UNNEST(CAST(SPLIT(REGEXP_REPLACE(metric_value, '[{}]',''), ',') AS ARRAY<VARCHAR(132)>)) AS t(split_part)
WHERE metric_name = 'util'
group by load_date, load_hr, timestamp, ne_name, object_id, metric_value
""")
报错原因

Spark SQL与Presto的语法解析逻辑存在差异,在Spark中,若在UNNEST的嵌套函数中直接引用原表列(如metric_value),会因解析优先级问题导致Spark无法识别该列,尤其是在多层函数嵌套的场景下。

解决方法

将metric_value的转换逻辑提前到子查询中,避免在UNNEST语句中直接嵌套复杂函数,同时修正group by与select字段不一致的问题(原查询select取object,但group by用object_id)。修改后的代码如下:

df1 = spark.sql("""
select * FROM database.table
WHERE date(load_date) = current_date - interval '1' day
    """)

df1.createOrReplaceTempView("table1")
    
df2 = spark.sql("""  
SELECT load_date, load_hr, timestamp,
ne_name, object, metric_value, 
MIN(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS MIN_UTIL, 
MAX(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS MAX_UTIL, 
AVG(CAST((CASE WHEN trim(split_part) = '' THEN NULL ELSE trim(split_part) END) AS DOUBLE))/100 AS AVG_UTIL
FROM (
    SELECT 
        load_date, load_hr, timestamp, ne_name, object, metric_value, metric_name,
        CAST(SPLIT(REGEXP_REPLACE(metric_value, '[{}]',''), ',') AS ARRAY<VARCHAR(132)>) AS metric_array
    FROM table1
    WHERE metric_name = 'util'
) t1
cross join UNNEST(t1.metric_array) AS t(split_part)
group by load_date, load_hr, timestamp, ne_name, object, metric_value
""")

修改说明

  1. 新增子查询t1,提前将metric_value转换为数组列metric_array,让Spark能正确识别原表列
  2. 主查询中直接对t1.metric_array执行UNNEST操作
  3. 修正group by字段,将object_id改为object,与select字段保持一致,避免额外报错
预期结果
Load_dateLoad_hrtimestampne_nameobjectmetric_valueMin_UTILMAX_UTILAVG_UTIL
2023-08-10152023-08-10T09:45AP1AP1.12.5{14311, 242342134, 13132}3.429.115.8

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 02:35:54