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

PySpark DataFrame如何新增累计值列与分组最大时间列

问题原因

你当前代码中val_sum(累计值)和max_time(分组最大时间)共用了同一个窗口定义w,这个窗口通过rowsBetween(Window.unboundedPreceding, Window.currentRow)限制了计算范围仅为分组内从第一行到当前行,所以累计求和的结果符合预期,但最大值计算只会取到当前行的最大值,无法覆盖全部分组。

调整方案

单独定义一个无行范围限制的分区窗口,专门用于计算全分组的最大时间,修改后的代码如下:

from pyspark.sql import SparkSession
from pyspark.sql import Window

import pyspark.sql.functions as F

data = [{"Category": 'A', "ID": 1, "Value": 121.44, "Truth": True, "time": 1},
        {"Category": 'B', "ID": 2, "Value": 300.01, "Truth": False, "time": 2},
        {"Category": 'C', "ID": 3, "Value": 10.99, "Truth": None, "time": 3},
        {"Category": 'C', "ID": 4, "Value": 33.87, "Truth": True, "time": 4},
        {"Category": 'D', "ID": 4, "Value": 33.87, "Truth": True, "time": 5},
        {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True, "time": 6},
        {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True, "time": 7},
        {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True, "time": 8}
        ]

spark = SparkSession.builder.getOrCreate()
df = spark.createDataFrame(data)

# 累计求和用的窗口:范围是分组内从首行到当前行
w_sum = Window.partitionBy(F.col("Category")).rowsBetween(Window.unboundedPreceding, Window.currentRow)
# 分组最大时间用的窗口:不限制行范围,默认覆盖整个分区
w_max = Window.partitionBy(F.col("Category"))

df.filter(df["ID"] == 4)\
  .withColumn("val_sum", F.sum(F.col("Value")).over(w_sum))\
  .withColumn("max_time", F.max(F.col("time")).over(w_max))\
  .show()

运行上述代码即可得到你预期的输出结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 05:45:10