Spark SQL中collect_set()分组语法错误排查求助
问题分析与解决
错误根源
- 语法位置错误:你错误地将
GROUP BY c.ContractNumber写在了SELECT子句的collect_set函数后面,GROUP BY是独立的SQL子句,必须放在WHERE子句之后,SELECT子句只能定义要查询的列、聚合函数和别名。 - collect_set参数错误:你用双引号包裹了
refcat.Name,这会让Spark把它当成字符串常量收集,而不是取refcat表的Name列值。
修正方案
按照SQL标准规范调整语句:
- 将GROUP BY子句移到整个SQL的末尾(WHERE之后)。
- 所有SELECT中未被聚合函数包裹的列,必须全部加入GROUP BY(Spark SQL严格遵循ANSI SQL规范,非聚合列必须是分组键的一部分)。
- 去掉collect_set参数的双引号,确保引用的是列值。
修正后的完整代码
dfc = spark.sql(""" SELECT c.Id, c.ContractNumber, c.CycleNumber, c.CycleEffectiveDate, c.CycleExpirationDate, INITCAP(c.Description) AS Description, COALESCE(GlobalAmendmentNumber, 0) AS GlobalAmendmentNumber, collect_set(refcat.Name) AS CategoryName, reftype.Name AS ContractType, refaward.name AS AwardStatus, refof.Name AS OrderFrom, refdn.NAME AS DepartmentName, CASE WHEN current_date() BETWEEN cast(c.CycleEffectiveDate AS timestamp) AND cast(c.CycleExpirationDate AS timestamp) THEN '1' WHEN current_date() < cast(c.CycleEffectiveDate AS timestamp) THEN '2' ELSE '3' END AS ContractWeight FROM contractmaster.contract c JOIN contractmaster.amendment am ON am.Id = c.Id JOIN contractmaster.taxon t ON t.HeaderId = c.Id JOIN contractmaster.referencedatum refcat ON t.CategoryId = refcat.Id JOIN contractmaster.referencedatum reftype ON c.ContractTypeId = reftype.Id JOIN contractmaster.header h ON c.Id = h.Id JOIN contractmaster.referencedatum refaward ON h.AwardStatusId = refaward.Id JOIN contractmaster.referencedatum refof ON h.OrderFromId = refof.Id JOIN contractmaster.referencedatum refdn ON h.DepartmentId = refdn.Id WHERE am.IsExecuted = True AND reftype.Type = "ContractType" AND refaward.Type = "AwardStatus" AND refof.Type = "OrderFrom" AND refdn.Type = "Department" GROUP BY c.Id, c.ContractNumber, c.CycleNumber, c.CycleEffectiveDate, c.CycleExpirationDate, c.Description, GlobalAmendmentNumber, reftype.Name, refaward.name, refof.Name, refdn.NAME, CASE WHEN current_date() BETWEEN cast(c.CycleEffectiveDate AS timestamp) AND cast(c.CycleExpirationDate AS timestamp) THEN '1' WHEN current_date() < cast(c.CycleEffectiveDate AS timestamp) THEN '2' ELSE '3' END """)
额外说明
如果觉得GROUP BY列过多,可以通过设置spark.sql.ansi.enabled=false关闭ANSI模式,允许非聚合列不加入GROUP BY,但这种做法可能导致结果不确定(因为Spark会随机选择分组内的某一行值),新手不建议使用,严格遵循规范更能保证结果正确性。
内容的提问来源于stack exchange,提问作者Mackenzie Jones
相关产品推荐
相关产品推荐

