如何基于历史连续行过滤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
相关产品推荐
相关产品推荐

