PySpark中Spark SQL触发ParseError异常求助(AWS Glue环境)
解决AWS Glue PySpark SQL ParseError问题
核心问题分析
你的SQL在客户端能运行但Glue中报错,核心原因是Spark SQL与你使用的客户端SQL(如PostgreSQL)存在语法兼容性差异,加上Spark SQL对语法的严格性要求,导致解析失败。以下是具体问题和修复步骤:
具体修复点
1. 替换类型转换语法
Spark SQL不支持PostgreSQL风格的::INT类型转换,需改用CAST()函数:
原代码片段:
date_part('second', aud.created_time)::INT
修改为:
CAST(date_part('second', aud.created_time) AS INT)
2. 调整INTERVAL语法
Spark SQL的INTERVAL格式为INTERVAL <数值> <单位>,不支持带引号的'10 sec'写法:
原代码片段:
INTERVAL '10 sec'
修改为:
INTERVAL 10 SECONDS
3. 替换数组构造方式
Spark SQL中创建数组需使用array()函数,而非直接用方括号[value]:
原代码片段(多处出现):
max(CASE WHEN lower(field) = 'annotation' THEN [value] ELSE NULL END)
修改为:
max(CASE WHEN lower(field) = 'annotation' THEN array(value) ELSE NULL END)
4. GROUP BY中禁用SELECT别名
Spark SQL默认不允许在GROUP BY子句中直接引用SELECT里定义的别名,需替换为原始列名:
原GROUP BY中的role_name是r.name AS role_name的别名,需改为r.name:
原代码片段:
GROUP BY aud.case_id, cas.account_id, role_name, r.team_name, ...
修改为:
GROUP BY aud.case_id, cas.account_id, r.name, r.team_name, ...
5. 限定未明确的列名
WHERE子句中的field未指定表别名,可能引发列名歧义(多表可能存在同名列),需统一加上aud.前缀:
原代码片段(多处出现):
AND NOT (field = 'created_time' AND username = 'System')
修改为:
AND NOT (aud.field = 'created_time' AND aud.username = 'System')
修复后的完整代码示例
def getSCMCaseHistLanding(self): landingDF = self.spark.sql("""SELECT aud.case_id, r.name AS role_name, cas.account_id, cas.created_time AS case_created_time, cas.last_updated_time, cas.screening_decision, aud.created_time AS record_created_time, aud.username, row_number() OVER (PARTITION BY aud.case_id ORDER BY date_trunc('minute', aud.created_time) + CAST(date_part('second', aud.created_time) AS INT) / 10 * INTERVAL 10 SECONDS) rank, max(CASE WHEN aud.field = 'status_id' THEN aud.value ELSE NULL END) AS d_status, max(CASE WHEN aud.field = 'status_id' THEN SPLIT_PART(SPLIT_PART(aud.description, 'from ', 2), ' to', 1) ELSE NULL END) AS src_status, max(CASE WHEN lower(aud.field) = 'annotation' THEN array(aud.value) ELSE NULL END) AS annotation, max(CASE WHEN lower(aud.field) = 'reason_id' THEN array(aud.value) ELSE NULL END) AS reason, max(CASE WHEN lower(aud.field) = 'approver_id' THEN array(aud.value) ELSE NULL END) AS approver, r.team_name AS approver_source, max(CASE WHEN lower(aud.field) = 'decision_id' THEN array(aud.value) ELSE NULL END) AS decision, max(CASE WHEN aud.field = 'status_id' THEN aud.created_time ELSE NULL END) AS lsd, last_value(lsd IGNORE NULLS) OVER (PARTITION BY aud.case_id ORDER BY record_created_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS last_status_changed_date, max(CASE WHEN aud.field = 'assigned_to' THEN REVERSE(SPLIT_PART(REVERSE(aud.value), ' ', 1)) ELSE NULL END) AS assigned_to_t, last_value(assigned_to_t IGNORE NULLS) OVER (PARTITION BY aud.case_id ORDER BY record_created_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS assigned_to, max(CASE WHEN aud.field = 'assigned_to' THEN aud.username ELSE NULL END) AS assigned_by, max(CASE WHEN (aud.field = 'assigned_to' AND aud.description LIKE 'Workbasket%') OR (aud.field = 'assigned_to' AND aud.description LIKE '%by%') THEN 'Auto Assign' WHEN aud.field = 'assigned_to' AND aud.description LIKE '%themself%' THEN 'Get Next' END) AS assignment_method, max( CASE WHEN aud.field = 'assigned_to' THEN aud.description ELSE NULL END) AS assignment_detail, max(CASE WHEN lower(aud.field) = 'state_id' THEN array(aud.value) ELSE NULL END) AS ll_state, last_value(d_status IGNORE NULLS) OVER (PARTITION BY aud.case_id ORDER BY record_created_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS dest_status, last_value(ll_state IGNORE NULLS) OVER (PARTITION BY aud.case_id ORDER BY record_created_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS l_state, CASE WHEN l_state IS NULL THEN 'Open' ELSE l_state END AS state, max(CASE WHEN lower(aud.field) = 'responsive_action_id' THEN array(aud.value) ELSE NULL END) AS responsive_action, max(CASE WHEN skill_type.type_name = 'Language' THEN skill.skill_name ELSE NULL END) AS language_skill, max( CASE WHEN skill_type.type_name = 'Business' THEN skill.skill_name ELSE NULL END) AS business_skill, max(CASE WHEN lower(aud.field) = 'annotation' AND lower(aud.value) LIKE 'f+ applied by jarvis%' THEN 'jarvis' WHEN lower(aud.field) = 'annotation' AND (lower(aud.value) LIKE '%bulk action%' OR lower(aud.value) LIKE '%moving cases to invalid data%') THEN 'bulk action' WHEN lower(aud.field) = 'annotation' AND lower(aud.value) LIKE '%t+ decision recommended by scr model%' THEN 'SCR T+ Model' WHEN lower(aud.field) = 'annotation' AND lower(aud.value) LIKE '%f+ applied by scr model%' THEN 'SCR F+ Model' ELSE NULL END) AS source_of_action, max(CASE WHEN lower(aud.field) = 'annotation' AND lower(aud.value) LIKE '%duplicate case%' THEN 'Y' ELSE 'N' END) AS duplicate_case, max(CASE WHEN aud.field = 'state_id' AND aud.description = 'Case reopened' THEN date_trunc('minute', aud.created_time) + CAST(date_part('second', aud.created_time) AS INT) / 10 * INTERVAL 10 SECONDS ELSE NULL END) AS case_reopened_time, max(0) AS active_record_flag FROM cm_spectre_case_audit aud JOIN cm_spectre_case cas ON cas.case_id = aud.case_id LEFT JOIN v_cm_user_snapshot u ON CASE WHEN aud.field = 'assigned_to' THEN REVERSE(SPLIT_PART(REVERSE(aud.value), ' ', 1)) ELSE aud.username END = u.alias AND aud.created_time BETWEEN u.start_time AND u.end_time LEFT JOIN cm_role r ON u.current_role_id = r.role_id LEFT JOIN cm_lookup_skills lookup_skill ON lookup_skill.case_id = cas.case_id LEFT JOIN cm_skill skill ON skill.skill_id = lookup_skill.skill_id LEFT JOIN cm_skill_type skill_type ON skill_type.type_id = skill.type_id WHERE 1 = 1 AND NOT (aud.field = 'created_time' AND aud.username = 'System') AND NOT (aud.field = 'annotation' AND aud.username = 'System') AND NOT (aud.field = 'skill_name_') AND NOT (aud.field = 'assigned_to' AND lower(aud.description) LIKE '%unassigned%') AND NOT (aud.field = 'attachment') AND NOT (aud.field = 'accept_list_id') AND NOT (aud.field LIKE 'screening_match_id%') GROUP BY aud.case_id, cas.account_id, r.name, r.team_name, cas.created_time, cas.last_updated_time, cas.screening_decision, aud.created_time, aud.username """) return landingDF
额外排查建议
- 若仍报错,可将SQL拆分为多个小片段逐步测试,定位具体出错的语句块;
- 检查Glue使用的Spark版本,部分语法可能因版本差异存在支持性问题;
- 确保临时表的列名、数据类型与SQL中引用的完全匹配,避免隐式类型转换引发的解析问题。
内容的提问来源于stack exchange,提问作者Jatin
相关产品推荐
相关产品推荐

