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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 01:05:22