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 """)
修改说明
- 新增子查询
t1,提前将metric_value转换为数组列metric_array,让Spark能正确识别原表列 - 主查询中直接对
t1.metric_array执行UNNEST操作 - 修正
group by字段,将object_id改为object,与select字段保持一致,避免额外报错
预期结果
| Load_date | Load_hr | timestamp | ne_name | object | metric_value | Min_UTIL | MAX_UTIL | AVG_UTIL |
|---|---|---|---|---|---|---|---|---|
| 2023-08-10 | 15 | 2023-08-10T09:45 | AP1 | AP1.12.5 | {14311, 242342134, 13132} | 3.4 | 29.1 | 15.8 |
内容的提问来源于stack exchange,提问作者cpljp
相关产品推荐
相关产品推荐

