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

如何基于历史连续行过滤DataFrame,保留单调递增列?

如何在Spark DataFrame中过滤出value列单调递增的记录(基于历史最大值)

你遇到的需求是要从按timestamp排序的DataFrame中,保留那些value值超过所有此前保留行最大值的记录,最终得到单调递增的序列——这个需求可以通过Spark的窗口函数完美解决,不需要复杂的自定义状态逻辑。

核心思路

因为你的DataFrame已经按timestamp排序,我们可以用窗口函数计算从第一行到当前行的全局累积最大值,然后只保留那些value值等于这个累积最大值的行。这些行就是每次刷新历史最大值的节点,自然构成了单调递增的序列。

完整实现代码

首先导入必要的Spark函数和窗口类:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.max
import org.apache.spark.sql.functions.col

然后处理你的输入数据:

val input = Seq( (1,4), (2,2), (3,3), (4,5), (5, 1), (6, 9), (7, 6) ).toDF("timestamp", "value")
input.show()

接下来定义全局排序的窗口,并计算累积最大值:

// 定义按timestamp排序的全局窗口
val globalWindow = Window.orderBy("timestamp")

// 添加累积最大值列
val withRunningMax = input.withColumn("running_max", max("value").over(globalWindow))
withRunningMax.show()

这一步的输出会显示每行的累积最大值:

+---------+-----+-----------+
|timestamp|value|running_max|
+---------+-----+-----------+
|        1|    4|          4|
|        2|    2|          4|
|        3|    3|          4|
|        4|    5|          5|
|        5|    1|          5|
|        6|    9|          9|
|        7|    6|          9|
+---------+-----+-----------+

最后筛选出value等于累积最大值的行,并移除辅助列:

val result = withRunningMax.filter(col("value") === col("running_max")).drop("running_max")
result.show()

得到的结果完全符合你的预期:

+---------+-----+
|timestamp|value|
+---------+-----+
|        1|    4|
|        4|    5|
|        6|    9|
+---------+-----+

额外说明

  • 这个方法是通用的,不管你的DataFrame有多大,只要确保数据按timestamp全局排序,窗口函数就能正确计算累积最大值。
  • 如果需要严格递增且不保留重复的最大值(比如序列4,4,5中只保留第一个4和5),可以在筛选后再按value分组,取每个最大值对应的第一条记录:
    val strictIncrementResult = result.withColumn("row_num", row_number().over(Window.partitionBy("value").orderBy("timestamp")))
                                      .filter(col("row_num") === 1)
                                      .drop("row_num")
    
  • 注意如果你的数据是分布式存储的,全局窗口可能会有性能开销,但对于大多数场景来说,这种实现已经足够高效且简洁。

内容的提问来源于stack exchange,提问作者gagan singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:01:58