PySpark DataFrame如何用窗口函数实现按条件计算前序行累计和
PySpark 窗口函数实现跨行滑动计算SUM列
实现逻辑说明
无需编写逐行循环逻辑,通过PySpark内置滑动窗口函数即可实现等价计算,性能远高于逐行遍历,核心逻辑如下:
- 首先定义滑动窗口范围:按业务要求的行序(示例中为
date字段升序,和样例数据的时间顺序一致),窗口范围固定取当前行之前的3行,不包含当前行,和pandas中取index-3:index-1的范围完全对应。 - 在窗口上预计算两个聚合值:前3行的
local字段之和、前3行的amount字段之和。 - 按规则判断取值:当前行
amount >= 8.99且前3行local之和等于300时,取前3行amount的和;其余所有情况(包括前序不足3行的场景)统一取值256。
完整实现代码
# 导入依赖 from pyspark.sql import Window import pyspark.sql.functions as F # 定义滑动窗口:全局按date升序,取当前行之前3行(不包含当前行) # 如果需要按单卡维度统计,只需要加上.partitionBy("card_uid")即可 win_spec = Window.orderBy("date").rowsBetween(-3, -1) # 新增SUM计算列 result_df = df.withColumn( "SUM", F.when( (F.col("amount") >= 8.99) & (F.sum("local").over(win_spec) == 300), F.round(F.sum("amount").over(win_spec), 2) # 保留2位小数避免浮点误差 ).otherwise(256) )
逻辑匹配说明
用给出的示例数据验证:
- 排序后前5行(前5条数据)要么前序不足3行,要么前3行存在
local=0的记录,local和凑不够300,因此SUM列均返回256。 - 从第6行开始,前3行的
local值均为100,求和刚好为300,且当前行amount=8.99满足阈值,因此前3行amount求和为8.99*3=26.97,和给出的目标结果完全一致。
注意事项
- 窗口定义用
rowsBetween而非rangeBetween,是为了严格按行位置取数,和pandas按行索引遍历的逻辑完全对齐,避免date字段重复值导致窗口范围计算错误。 - 计算amount和时加
round保留2位小数,是为了避免浮点数计算精度问题导致结果出现类似26.970000001的异常值。 - 如果业务要求按单个卡片维度独立统计(不跨卡片计算前3行),只需要把窗口定义修改为按
card_uid分区即可:
win_spec = Window.partitionBy("card_uid").orderBy("date").rowsBetween(-3, -1)
内容的提问来源于stack exchange,提问作者Babbara
相关产品推荐
相关产品推荐

