Spark SQL中filter_col字段无法访问的报错原因及解决咨询
Spark SQL过滤条件执行失败的原因及解决方法
错误原因
报错cannot resolve 'filter_col'是因为SQL的执行顺序限制:Spark SQL会先执行FROM子句(包括嵌套子查询),然后执行WHERE过滤,最后才执行SELECT投影生成新列。你原来的SQL中,filter_col是在SELECT里定义的别名,WHERE阶段执行时这个列还未被创建,自然无法识别。
解决方法
方法一:将计算逻辑移到子查询中
把生成filter_col的CASE逻辑放到内层子查询里,让外层WHERE能引用到这个已生成的列:
dfsql = spark.sql(""" select tiers, poid, etat, pays, zyada, date, max_date, cnt_tiers, filter_col from ( select *, MAX(date) OVER (PARTITION BY tiers) AS max_date, COUNT(*) OVER (PARTITION BY tiers) as cnt_tiers, CASE WHEN cnt_tiers = 1 THEN True WHEN cnt_tiers > 1 and date = max_date THEN True ELSE False END AS filter_col from tierstbl where poid is not null and etat is not null and pays is not null and zyada is not null ) a where filter_col = True """)
方法二:使用CTE(公共表表达式)
用CTE封装带计算列的数据集,逻辑更清晰易读:
dfsql = spark.sql(""" WITH tier_with_calc AS ( select *, MAX(date) OVER (PARTITION BY tiers) AS max_date, COUNT(*) OVER (PARTITION BY tiers) as cnt_tiers, CASE WHEN cnt_tiers = 1 THEN True WHEN cnt_tiers > 1 and date = max_date THEN True ELSE False END AS filter_col from tierstbl where poid is not null and etat is not null and pays is not null and zyada is not null ) select tiers, poid, etat, pays, zyada, date, max_date, cnt_tiers, filter_col from tier_with_calc where filter_col = True """)
方法三:直接在WHERE中复用CASE逻辑(不推荐)
如果不想嵌套子查询,可以直接把CASE的判断逻辑写到WHERE里,但会导致代码冗余,不利于维护:
dfsql = spark.sql(""" select a.tiers, a.poid, a.etat, a.pays, a.zyada, a.date, a.max_date, a.cnt_tiers, CASE WHEN a.cnt_tiers = 1 THEN True WHEN a.cnt_tiers > 1 and a.date = a.max_date THEN True ELSE False END AS filter_col from ( select *, MAX(date) OVER (PARTITION BY tiers) AS max_date, COUNT(*) OVER (PARTITION BY tiers) as cnt_tiers from tierstbl where poid is not null and etat is not null and pays is not null and zyada is not null ) a where (a.cnt_tiers = 1) OR (a.cnt_tiers > 1 AND a.date = a.max_date) """)
内容的提问来源于stack exchange,提问作者amine jisung
相关产品推荐
相关产品推荐

