如何从事实表正确获取上季度销售值并解决聚合值不一致问题
问题:获取上季度聚合销售值并统一同季度period的数值
我尝试两种方法获取上季度销售值:一种用lag函数,另一种是通过聚合+关联的方式,但遇到了问题:按product_id、year、quarter聚合sales后,同一年度同一季度的不同period对应的聚合值不一致,期望同一季度的所有period显示相同的上季度聚合销售值。
原始DataFrame结构
+----------+-----------+-----+------+----+-------+------------+---------+ |product_id|customer_id|sales|period|year|quarter|prev_quarter|prev_year| +----------+-----------+-----+------+----+-------+------------+---------+ | 24| 2| 288| 8|2022| 3| 2| 2022| | 14| 5| 527| 1|2022| 1| 4| 2021| | 32| 3| 763| 6|2022| 2| 1| 2022| | 25| 5| 175| 4|2022| 2| 1| 2022| | 36| 5| 840| 1|2023| 1| 4| 2022| | 37| 3| 471| 11|2021| 4| 3| 2021| | 4| 1| 990| 7|2023| 3| 2| 2023| | 42| 1| 225| 10|2023| 4| 3| 2023| | 24| 5| 951| 4|2020| 2| 1| 2020| | 25| 1| 531| 10|2023| 4| 3| 2023| | 2| 4| 802| 7|2021| 3| 2| 2021| | 34| 2| 950| 12|2023| 4| 3| 2023| | 13| 5| 916| 3|2020| 1| 4| 2019| | 28| 3| 432| 2|2020| 1| 4| 2019| +----------+-----------+-----+------+----+-------+------------+---------+
尝试的代码及问题
我先创建聚合表计算各product_id在各季度的销售总额,再与原表关联获取上季度值:
df_kpi = df df_kpi = df_kpi.groupBy("year", "quarter", "product_id").agg(sum(col("sales")).alias('sales_total')) df_joined = df.alias('df').join( df_kpi.alias('kpi'), (col("df.product_id") == col('kpi.product_id')) & (col('df.prev_quarter') == col('kpi.quarter')) & (col('df.prev_year') == col("kpi.year")), "left" ).select("df.year", "df.prev_year","df.quarter", "df.sales", "df.period","df.prev_quarter", "df.product_id", "df.customer_id", "kpi.sales_total")
但得到的结果中,同year同quarter的不同period对应的sales_total值不一致,比如2020年Q2的三个period分别显示4550、6675、2555,而期望同一季度的所有period都显示相同的上季度聚合值。
解决方案
问题的核心是确保同一product_id、year、quarter的所有行,能统一获取到对应的上季度聚合销售值。可以通过聚合表+窗口函数填充的方式解决:
1. 生成正确的聚合表
先按product_id、year、quarter聚合计算销售总额,确保每个维度的聚合值唯一:
from pyspark.sql import functions as F # 聚合各product_id在每季度的销售总额 df_kpi = df.groupBy("product_id", "year", "quarter").agg(F.sum("sales").alias("sales_total"))
2. 关联原表与聚合表
将原表与聚合表关联,匹配上季度(prev_year+prev_quarter)的销售总额:
df_joined = df.alias("df").join( df_kpi.alias("kpi"), (F.col("df.product_id") == F.col("kpi.product_id")) & (F.col("df.prev_year") == F.col("kpi.year")) & (F.col("df.prev_quarter") == F.col("kpi.quarter")), how="left" ).select( "df.year", "df.prev_year", "df.quarter", "df.sales", "df.period", "df.prev_quarter", "df.product_id", "df.customer_id", F.col("kpi.sales_total").alias("prev_quarter_sales_total") )
3. 窗口函数统一同季度值
如果同一季度内存在部分行匹配到聚合值、部分未匹配的情况,用窗口函数将同一product_id、year、quarter内的聚合值统一填充:
from pyspark.sql.window import Window # 定义窗口:按product_id、year、quarter分组 window_spec = Window.partitionBy("product_id", "year", "quarter") # 填充同一季度内的上季度销售总额,确保所有行值一致 df_final = df_joined.withColumn( "prev_quarter_sales_total", F.first("prev_quarter_sales_total", ignorenulls=True).over(window_spec) )
处理后,同一product_id、year、quarter的所有period行,都会显示相同的上季度聚合销售值,符合预期。
内容的提问来源于stack exchange,提问作者Tomás Jullier
相关产品推荐
相关产品推荐

