将指定Spark SQL查询转换为等效的PySpark DataFrame API实现代码
# 导入依赖的PySpark工具类 from pyspark.sql import Window from pyspark.sql.functions import row_number, desc # 定义窗口规则:和原SQL的窗口逻辑完全对齐,按se10分区,按se_aggrtr_indctr降序排序 window_spec = Window.partitionBy("se10").orderBy(desc("se_aggrtr_indctr")) # 完整转换后的处理逻辑 result_df = ( fraud_details_data_whole # 对应内层子查询的GROUP BY逻辑:四个字段分组等价于对这四个字段去重 .select("se10", "se3", "se_aggrtr_indctr", "key_swipe_ind") .distinct() # 对应ROW_NUMBER窗口函数打行号 .withColumn("rn", row_number().over(window_spec)) # 对应外层WHERE过滤条件 .filter("rn < 2") # 选中最终返回的字段 .select("se10", "se3", "se_aggrtr_indctr", "key_swipe_ind") )
说明:上述代码和原SQL语义100%对齐,没有逻辑差异。原SQL内层的GROUP BY因为分组字段和查询的非聚合字段完全一致,用
distinct()实现更简洁,执行效果和GROUP BY完全相同。
内容的提问来源于stack exchange,提问作者vashi
相关产品推荐
相关产品推荐

