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
相关产品推荐
相关产品推荐

