Spark DataFrame按品牌计算累计和、累计平均值返回全null问题
Spark DataFrame 分组累计计算问题修复方案
问题原因
withColumn语法错误:withColumn仅接受2个参数(列名、列表达式),你传入了2个独立的when表达式,语法逻辑错误导致所有返回值为null,需要将两个判断逻辑合并为一个链式列表达式。- 数据类型不兼容:原始
TrueValue字段中brand3的取值带%后缀为字符串格式,无法直接参与数值计算,需要先清洗转换为数值类型。
修复后完整代码
1. 依赖导入
from pyspark.sql import functions as F from pyspark.sql.window import Window
2. 数据预处理
# 清洗TrueValue字段:移除%后缀,转换为double类型,空值保留null df = df.withColumn('TrueValue', F.regexp_replace(F.col('TrueValue'), '%', '').cast('double')) # 生成标准时间戳列(如你已生成可跳过该步) df = df.withColumn('month_in_timestamp', F.to_date(F.col('Month'), 'M/d/yyyy').cast('timestamp'))
3. 窗口定义(原逻辑正确无需修改)
windowval = (Window.partitionBy('Location','Brand') .orderBy('month_in_timestamp') .rangeBetween(Window.unboundedPreceding, 0))
4. 累计值计算
df = df.withColumn('TotalSumValue', F.when(F.col('Brand').isin('brand1', 'brand2'), F.sum('TrueValue').over(windowval)) .when(F.col('Brand') == 'brand3', F.avg('TrueValue').over(windowval)))
逻辑说明
- Spark内置的
sum、avg函数会自动忽略窗口范围内的null值,因此空行的累计值会自动继承上一行的计算结果,和预期输出一致。 - brand3的平均值计算结果会自动四舍五入,符合你给出的预期输出。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

