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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:57:52