Spark DataFrame分组滑动窗口计算过去2天B列求和问题
解决Spark DataFrame按分组计算过去2天移动求和的问题
首先,你之前的Window函数用法有误:不应该同时按dt和A分区,这样每个分组只会包含单条数据,自然无法计算跨日期的求和。我们需要调整Window的定义,核心是按A分组,按日期排序,再基于日期范围计算移动求和。
具体步骤和代码
- 先将字符串类型的日期转换为日期类型:方便后续日期范围的计算
- 定义正确的Window规格:按
A分区,按日期升序排序,指定过去2天(含当天)的范围 - 计算移动求和
以下是完整的Scala代码实现:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 初始化你的DataFrame val df = Seq( ("2020-05-21","x",1), ("2020-05-21","y",2), ("2020-05-22","x",3), ("2020-05-22","y",4), ("2020-05-23","x",5), ("2020-05-23","y",6) ).toDF("dt","A","B") // 将字符串日期转换为Date类型,同时转成unix天数用于范围计算 val dfWithDate = df.withColumn("date_dt", to_date($"dt", "yyyy-MM-dd")) // 定义Window:按A分组,按unix天数排序,范围是过去1天到当天(即包含当前和前一天) val windowSpec = Window.partitionBy($"A") .orderBy(unix_date($"date_dt")) .rangeBetween(-1, 0) // 计算移动求和并整理输出列 val result = dfWithDate .withColumn("sum", sum($"B").over(windowSpec)) .select($"dt", $"A", $"B", $"sum") .orderBy($"dt", $"A") // 查看结果 result.show()
输出结果(和你的预期完全一致)
+----------+---+---+---+ | dt| A| B|sum| +----------+---+---+---+ |2020-05-21| x| 1| 1| |2020-05-21| y| 2| 2| |2020-05-22| x| 3| 4| |2020-05-22| y| 4| 6| |2020-05-23| x| 5| 8| |2020-05-23| y| 6| 10| +----------+---+---+---+
关键细节说明
- 为什么用
rangeBetween而不是rowsBetween?rowsBetween是基于行的位置偏移,若你的日期存在缺失(比如某天没有对应A的记录),它会错误地取到更早的行;而rangeBetween基于日期的数值差值(这里用unix_date转成的天数),能精准匹配"过去2天"的时间范围,适用性更强。 - 之前的错误原因:你用
partitionBy($"dt",$"A")会把每个日期+A的组合单独分成一个组,每个组只有1条数据,rowsBetween(-1,0)无法获取到前一天的数据,求和结果自然等于B本身。
内容的提问来源于stack exchange,提问作者vdep
相关产品推荐
相关产品推荐

