Spark Scala中Window.rowsBetween的使用及窗口范围控制问题
解决Spark行窗口统计的问题
你之前的代码用了rangeBetween,这其实是个容易踩的常见坑——rangeBetween是基于排序列的数值范围(比如你的time列值的区间),而不是行数。如果要按行数来定义窗口,得用rowsBetween才行!
下面分两种常见场景给你具体的实现方案:
一、滑动行窗口(比如每个窗口包含5行数据)
如果是要做滑动的行窗口(比如当前行+前4行,共5行;或者当前行+后4行),直接用rowsBetween定义窗口范围:
首先导入必要的包:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._
然后定义窗口并计算统计值:
// 定义滑动窗口:当前行及前4行(总共5行) // 要是需要当前行到后4行,就改成 rowsBetween(0, 4) val slidingWinSpec = Window.orderBy("time").rowsBetween(-4, 0) // 计算窗口内的class总数、class1数量、其他类数量 val dfWithSlidingStats = df .withColumn("total_class", count("class").over(slidingWinSpec)) .withColumn("class1_count", sum(when(col("class") === 1, 1).otherwise(0)).over(slidingWinSpec)) .withColumn("other_class_count", sum(when(col("class") !== 1, 1).otherwise(0)).over(slidingWinSpec))
这里的rowsBetween(-4, 0)参数含义:
-4:相对于当前行往前数4行0:当前行
所以整个窗口包含当前行和它前面的4行,刚好5行数据。
二、固定行范围窗口(比如第1行到第10行)
如果是要统计指定行区间(比如第1到第10行)的统计值,需要先给数据添加行号,再基于行号过滤或定义窗口:
方案1:直接过滤行号后聚合(只输出统计结果)
// 先给数据添加从1开始的行号 val dfWithRowNum = df.withColumn("row_num", row_number().over(Window.orderBy("time"))) // 过滤出第1到第10行,然后聚合统计 val fixedRangeStats = dfWithRowNum .filter(col("row_num").between(1, 10)) .agg( count("class").alias("total_class"), sum(when(col("class") === 1, 1).otherwise(0)).alias("class1_count"), sum(when(col("class") !== 1, 1).otherwise(0)).alias("other_class_count") )
方案2:保留所有行,标记固定窗口的统计值
如果需要把固定窗口的统计值附加到每一行上,可以这么写:
val dfWithRowNum = df.withColumn("row_num", row_number().over(Window.orderBy("time"))) // 定义覆盖第1到第10行的固定窗口 val fixedWinSpec = Window .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) .where(col("row_num").between(1, 10)) val dfWithFixedStats = dfWithRowNum .withColumn("total_class_fixed", count("class").over(fixedWinSpec)) .withColumn("class1_count_fixed", sum(when(col("class") === 1, 1).otherwise(0)).over(fixedWinSpec)) .withColumn("other_class_count_fixed", sum(when(col("class") !== 1, 1).otherwise(0)).over(fixedWinSpec))
关键提醒
- 一定要区分
rowsBetween和rangeBetween:前者是基于行数偏移,后者是基于排序列的数值范围(比如你的time如果是时间戳,rangeBetween(0,5)会匹配time在当前行time到time+5之间的行,不是5行) - 窗口必须先通过
orderBy指定排序规则,否则行的顺序是不确定的,窗口统计结果也会不可靠
内容的提问来源于stack exchange,提问作者Foaad Mohamad Haddod
相关产品推荐
相关产品推荐

