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

使用窗口函数实现Spark DataFrame按组累计求和(当前行+前一行)

解决Spark DataFrame分组累计求和问题

没问题,我来帮你搞定这个需求!你需要基于groupid和id分组,计算当前行Amount和前一行累计值相加的累计总和,本质上就是分组内按时间顺序的累计求和,用Spark的窗口函数就能轻松实现。

首先咱们先把你的DataFrame初始化代码补全完整(你给的示例数据缺了groupid和Amount列,我补充了合理的测试数据):

import findspark
findspark.init()
import pyspark
import pandas as pd
from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import sum, col

spark = SparkSession.builder.getOrCreate()
sc = spark.sparkContext

# 补全完整的测试数据,包含groupid、id、date、Amount列
data1 = {
    'date': ['2018-04-03', '2018-04-04', '2018-04-05', '2018-04-06', '2018-04-07', '2018-04-08'],
    'id': ['id1', 'id2', 'id1', 'id1', 'id2', 'id2'],
    'groupid': ['g1', 'g1', 'g1', 'g2', 'g2', 'g2'],
    'Amount': [100, 200, 150, 300, 250, 180]
}
# 转为Spark DataFrame
df = spark.createDataFrame(pd.DataFrame(data1))
df.show()

执行后会看到初始的DataFrame:

+----------+---+-------+------+
|      date| id|groupid|Amount|
+----------+---+-------+------+
|2018-04-03|id1|     g1|   100|
|2018-04-04|id2|     g1|   200|
|2018-04-05|id1|     g1|   150|
|2018-04-06|id1|     g2|   300|
|2018-04-07|id2|     g2|   250|
|2018-04-08|id2|     g2|   180|
+----------+---+-------+------+

接下来是核心实现:用窗口函数定义分组和排序规则,然后计算累计总和:

# 定义窗口规范:按groupid和id分组,按date升序排序(保证累计的时间顺序)
window_spec = Window.partitionBy("groupid", "id").orderBy("date")

# 计算累计总和:sum(Amount) over窗口,默认是从分组第一行到当前行的总和
df_with_cumulative = df.withColumn(
    "cumulative_sum",
    sum(col("Amount")).over(window_spec)
)

df_with_cumulative.show()

执行后就能得到你要的结果:

+----------+---+-------+------+-------------+
|      date| id|groupid|Amount|cumulative_sum|
+----------+---+-------+------+-------------+
|2018-04-03|id1|     g1|   100|          100|
|2018-04-05|id1|     g1|   150|          250|
|2018-04-04|id2|     g1|   200|          200|
|2018-04-06|id1|     g2|   300|          300|
|2018-04-07|id2|     g2|   250|          250|
|2018-04-08|id2|     g2|   180|          430|
+----------+---+-------+------+-------------+

关键解释:

  • partitionBy("groupid", "id"):指定按groupid和id分组,每个分组内独立计算累计值
  • orderBy("date"):保证每个分组内的行按日期顺序排列,这样累计是按时间先后进行的
  • sum(col("Amount")).over(window_spec):这个函数会自动计算从分组的第一行到当前行的Amount总和,正好对应你要的“当前行Amount值与前一行累计值相加”的结果(前一行累计是之前所有行的总和,加上当前行就是到当前行的累计)

如果你的需求里有特殊的行范围要求(比如只取前一行加当前行),可以在窗口里加上rowsBetween,但根据你的描述,默认的累计求和就完全符合需求啦。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:32:31