Pyspark Dataframe实现列历史最高值向下填充直至遇到更高值的方法
PySpark 实现amount列滚动历史最大值填充方案
你要实现的滚动历史最大值效果,可以直接用PySpark内置的窗口函数完成,核心逻辑是定义一个覆盖从首行到当前行的窗口,取窗口内amount的最大值即可,具体实现如下:
前置前提
你需要确保DataFrame存在可确定数据先后顺序的排序字段,比如事件时间、自增行号等,没有明确排序规则的情况下无法定义“历史”的先后顺序。
完整实现代码
# 导入依赖包 from pyspark.sql.window import Window from pyspark.sql.functions import max, col # 1. 定义窗口规则 # orderBy后面替换为你实际用来排序的字段,比如event_time、row_id # rowsBetween指定窗口范围是从第一行到当前行 window_spec = Window.orderBy("替换为你的排序字段").rowsBetween(Window.unboundedPreceding, 0) # 2. 计算滚动最大值,直接覆盖原amount列,也可以自定义新列名 df_result = df.withColumn("amount", max(col("amount")).over(window_spec)) # 查看结果 df_result.show()
扩展场景(可选)
如果你需要按维度分组计算各分组内的历史最大值,比如不同用户分别计算各自的amount历史最高值,只需在窗口定义中增加partitionBy参数指定分组字段即可:
# 示例:按user_id分组,按event_time排序计算各用户的滚动最大值 window_spec = Window.partitionBy("user_id").orderBy("event_time").rowsBetween(Window.unboundedPreceding, 0)
效果示例
假设原始数据排序后amount取值为[10,5,15,12,20,18],运行代码后输出的amount结果为[10,10,15,15,20,20],和需求的预期效果完全匹配。
内容的提问来源于stack exchange,提问作者chandu
相关产品推荐
相关产品推荐

