PySpark如何按Brand、Sector分组按月排序计算TrueValue列累计和
PySpark 按月正确计算分组累计和的解决方案
排序错误的核心原因是Month字段为字符串类型,默认按字典序排序,会将首字符为1的10/1/2021、11/1/2021等排在首字符为2的2/1/2021之前,导致窗口排序不符合预期。
解决步骤
- 完成两个必要的字段类型转换:
- 将字符串格式的
Month转为日期类型,保证可以按自然月份顺序排序 - 将字符串类型的
TrueValue转为数值类型,保证求和计算正常
- 窗口定义时使用转换后的日期字段排序,再计算累计和即可。
完整代码示例
from pyspark.sql import functions as F from pyspark.sql.window import Window # 字段类型转换 df1 = df1.withColumn('TrueValue', F.col('TrueValue').cast('double')) # 匹配日期格式:月/日/年,若为日/月/年可改为d/M/yyyy df1 = df1.withColumn('sort_month', F.to_date(F.col('Month'), 'M/d/yyyy')) # 定义窗口按转换后的日期排序 windowval = (Window.partitionBy('Brand','Sector').orderBy('sort_month') .rangeBetween(Window.unboundedPreceding, 0)) df1 = df1.withColumn('TotalSumValue', F.sum('TrueValue').over(windowval)) # 删除临时排序字段(可选) df1 = df1.drop('sort_month')
简化写法(无临时字段)
可以直接在窗口的orderBy方法中完成日期转换,不需要额外新增字段:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 转换TrueValue为数值类型 df1 = df1.withColumn('TrueValue', F.col('TrueValue').cast('double')) # 直接在orderBy中做日期转换 windowval = (Window.partitionBy('Brand','Sector') .orderBy(F.to_date(F.col('Month'), 'M/d/yyyy')) .rangeBetween(Window.unboundedPreceding, 0)) df1 = df1.withColumn('TotalSumValue', F.sum('TrueValue').over(windowval))
注意事项
上述代码中日期格式串M/d/yyyy匹配样例给出的日期格式,可根据实际数据的日期格式调整。空值的TrueValue在求和时会被自动忽略,最终累计结果和预期输出完全一致。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

