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

将指定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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 22:27:04