如何在PySpark中计算列累计和并新增列,含百分比特殊计算逻辑
问题原因排查
你的现有代码存在三个核心问题:
- 排序逻辑错误:你在窗口排序时使用了字符串类型的
Month字段,没有按转换后的时间戳排序,直接导致百分比计算的顺序不符合要求 - 未适配百分比数据类型:带
%的TrueValue是字符串格式,直接求和会得到错误结果,也没有实现百分比类数据的累计求均值逻辑 - 没有区分普通数值和百分比的不同计算规则
完整解决方案代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 预处理字段 df = df.withColumn("month_in_timestamp", F.to_timestamp(df.Month, 'M/d/yyyy')) \ # 将TrueValue转换为可计算的数值,同时识别百分比格式 .withColumn("numeric_value", F.when(F.col("TrueValue").endswith("%"), F.regexp_replace(F.col("TrueValue"), "%", "").cast("double")) .otherwise(F.col("TrueValue").cast("double"))) \ # 标记当前Brand+Sector分组是否为百分比类型(同组数据类型统一) .withColumn("is_percent", F.max(F.when(F.col("TrueValue").endswith("%"), True).otherwise(False)) .over(Window.partitionBy("Brand", "Sector"))) # 2. 定义窗口:按品牌+领域分组,按日期升序排序,取当前行及之前所有行作为计算范围 windowval = Window.partitionBy('Brand','Sector').orderBy('month_in_timestamp') \ .rangeBetween(Window.unboundedPreceding, 0) # 3. 按规则计算TotalSumValue df = df.withColumn("cum_sum", F.sum("numeric_value").over(windowval)) \ .withColumn("cum_count", F.count("numeric_value").over(windowval)) \ .withColumn("TotalSumValue", F.when(F.col("is_percent") == True, # 百分比类型:累计和除以累计有效月份数,保留1位小数后加% F.concat(F.round(F.col("cum_sum") / F.col("cum_count"), 1), F.lit("%"))) .otherwise(F.col("cum_sum"))) # 可选:删除中间生成的辅助字段 df = df.drop("month_in_timestamp", "numeric_value", "is_percent", "cum_sum", "cum_count")
该代码完全匹配你给出的预期输出结果,空值行会自动继承上一行的累计值,百分比计算逻辑完全符合示例要求。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

