如何在Spark SQL中使用qualify?PySpark SQL等价写法求助
Spark SQL中QUALIFY的使用及PySpark改写方案
一、如何在Spark SQL中使用QUALIFY?
Spark SQL从3.3版本开始正式支持QUALIFY子句,它专门用于对窗口函数计算后的结果进行过滤,弥补了WHERE(仅过滤原始行)和HAVING(仅过滤聚合结果)无法直接过滤窗口函数结果的不足。
使用规则:
QUALIFY需跟在窗口函数表达式之后,语法上位于WHERE/HAVING之后(如果有这两个子句的话)- 可直接在子句中使用窗口函数,或引用窗口函数的别名
基础示例:
-- 筛选每个城市下交易时间最新的记录 SELECT id, city, transaction_ts FROM table QUALIFY transaction_ts = MAX(transaction_ts) OVER (PARTITION BY city)
二、目标SQL的PySpark改写方案
你的原始SQL语句是:
select * from table where id = 1 qualify transaction_ts = max(transaction_ts) over (partition by city)
方案1:Spark 3.3+ 直接使用QUALIFY
如果你的Spark版本≥3.3,PySpark中可以直接通过原生SQL或DataFrame API实现:
方式1:Spark SQL原生执行
spark.sql(""" SELECT * FROM table WHERE id = 1 QUALIFY transaction_ts = MAX(transaction_ts) OVER (PARTITION BY city) """).show()
方式2:DataFrame API写法
from pyspark.sql import Window from pyspark.sql.functions import max window_spec = Window.partitionBy("city") df = spark.table("table") result_df = df.filter(df.id == 1) \ .withColumn("latest_ts", max("transaction_ts").over(window_spec)) \ .filter(df.transaction_ts == df.latest_ts) \ .drop("latest_ts") result_df.show()
方案2:Spark 3.3以下版本 用子查询替代
如果Spark版本低于3.3不支持QUALIFY,可以用子查询/CTE先计算窗口函数结果,再过滤:
方式1:Spark SQL子查询写法
spark.sql(""" SELECT t.* FROM ( SELECT *, MAX(transaction_ts) OVER (PARTITION BY city) AS latest_ts FROM table WHERE id = 1 ) t WHERE t.transaction_ts = t.latest_ts """).show()
方式2:DataFrame API写法
from pyspark.sql import Window from pyspark.sql.functions import max window_spec = Window.partitionBy("city") df = spark.table("table") temp_df = df.filter(df.id == 1).withColumn("latest_ts", max("transaction_ts").over(window_spec)) result_df = temp_df.filter(temp_df.transaction_ts == temp_df.latest_ts).drop("latest_ts") result_df.show()
内容的提问来源于stack exchange,提问作者Sankar Azad
相关产品推荐
相关产品推荐

