PySpark中多表关联+Case When+GroupBy计算销售占比问题
PySpark实现分组计算ON/OFF状态销售占比解决方案
需求回顾
需按ID、Region、Brand、Code、Start_dt维度,分别计算ON状态和OFF状态的销售额占对应分组总销售额的比例,已完成两张表的关联,但无法在单段代码中同时处理状态过滤与占比计算。
核心思路
通过窗口函数先计算每个分组的总销售额,再用条件聚合分别统计ON/OFF状态的销售额,最后用状态销售额除以分组总销售额得到占比。
方案1:PySpark DataFrame API实现
假设关联后的DataFrame命名为joined_df,包含字段:ID, Region, Brand, Code, Start_dt, status, sales。
完整代码示例
from pyspark.sql import Window from pyspark.sql.functions import col, sum, when # 定义窗口:按指定维度分组 window_spec = Window.partitionBy("ID", "Region", "Brand", "Code", "Start_dt") # 计算分组总销售额,同时聚合ON/OFF状态的销售额,最终计算占比 result_df = joined_df.withColumn("total_sales", sum("sales").over(window_spec)) \ .groupBy("ID", "Region", "Brand", "Code", "Start_dt") \ .agg( sum(when(col("status") == "ON", col("sales")).otherwise(0)).alias("on_sales"), sum(when(col("status") == "OFF", col("sales")).otherwise(0)).alias("off_sales"), sum("total_sales").alias("total_sales") # 分组后取总销售额(因窗口计算后每组内值相同,sum后为单值) ) \ .withColumn("on_sales_ratio", col("on_sales") / col("total_sales")) \ .withColumn("off_sales_ratio", col("off_sales") / col("total_sales")) \ .select("ID", "Region", "Brand", "Code", "Start_dt", "on_sales", "off_sales", "total_sales", "on_sales_ratio", "off_sales_ratio") # 查看结果 result_df.show()
代码说明
- 窗口函数:通过
partitionBy指定分组维度,计算每个分组的总销售额total_sales,确保每条记录都能获取到所在分组的总销售额。 - 条件聚合:用
when-otherwise分别统计ON/OFF状态的销售额,避免多次过滤导致的重复计算。 - 占比计算:直接用状态销售额除以分组总销售额,得到占比(可根据需求用
round函数保留小数位数)。
方案2:Spark SQL实现
如果更熟悉SQL语法,可直接基于关联后的临时视图执行SQL查询,逻辑与DataFrame API一致:
完整SQL示例
-- 先将关联后的DataFrame注册为临时视图 joined_df.createOrReplaceTempView("joined_sales") -- 执行查询计算占比 SELECT ID, Region, Brand, Code, Start_dt, SUM(CASE WHEN status = 'ON' THEN sales ELSE 0 END) AS on_sales, SUM(CASE WHEN status = 'OFF' THEN sales ELSE 0 END) AS off_sales, SUM(sales) AS total_sales, SUM(CASE WHEN status = 'ON' THEN sales ELSE 0 END) / SUM(sales) AS on_sales_ratio, SUM(CASE WHEN status = 'OFF' THEN sales ELSE 0 END) / SUM(sales) AS off_sales_ratio FROM joined_sales GROUP BY ID, Region, Brand, Code, Start_dt
说明
直接通过CASE WHEN实现条件聚合,同时计算分组总销售额,一步得到占比,若需处理除零问题,可添加NULLIF(SUM(sales), 0)避免报错。
常见问题修正
如果之前的代码无法同时处理ON/OFF,大概率是因为分开过滤ON和OFF后再合并,导致重复计算分组总销售额。上述两种方案均通过一次分组+条件聚合完成,避免了多次计算的冗余,同时保证逻辑统一。
内容的提问来源于stack exchange,提问作者Shiva Golla
相关产品推荐
相关产品推荐

