如何在PySpark窗口聚合中实现按条件去重计数
PySpark统计门店商品指标问题修正方案
问题根因
原有代码的num_products_with_stock字段使用了按门店分区、按日期排序的累计滑动窗口,该逻辑只会累加历史出现过的符合条件的商品,不会扣除后续库存归0的商品,不符合当日库存大于0的去重商品总数的统计需求。
另外你的原始数据属于稀疏上报结构(仅库存变动的商品才会生成当日记录),直接按原表统计当日值会遗漏库存未变动的商品,导致统计结果不准。
修正后代码
from pyspark.sql.functions import * from pyspark.sql import Window # 1. 生成所有门店+日期+经营商品的全量组合,解决原始数据稀疏问题 store_date = df.select("store", "date").distinct() store_product = df.select("store", "product").distinct() full_dim = store_date.join(store_product, on="store", how="inner") # 2. 关联原始库存,前向填充每个商品的历史最新库存 win_fill_stock = Window.partitionBy("store", "product").orderBy("date") df_full_stock = full_dim \ .join(df, on=["store", "date", "product"], how="left") \ .withColumn("stock_c", last("stock_c", ignorenulls=True).over(win_fill_stock)) # 3. 定义两个统计窗口 # 累计窗口:统计门店截至当前日期的历史经营去重商品数 win_history = Window.partitionBy("store").orderBy("date").rangeBetween(Window.unboundedPreceding, Window.currentRow) # 当日窗口:统计门店当日库存大于0的去重商品数 win_daily = Window.partitionBy("store", "date") # 4. 计算指标后关联回原始表保留原有行 df_res = df_full_stock \ .withColumn("num_products", countDistinct("product").over(win_history)) \ .withColumn("num_products_with_stock", countDistinct(when(col("stock_c") > 0, col("product"))).over(win_daily)) \ .join(df.select("store", "date", "product", "stock_c"), on=["store", "date", "product", "stock_c"], how="right") \ .orderBy("store", "date", "product") df_res.show(truncate=False)
代码说明
- 若数据量较大不需要精确统计,可以将
countDistinct替换为approx_count_distinct提升性能 - 前向填充逻辑默认商品库存不会自动清零,若有定期清零规则可以自行调整填充逻辑
- 最终保留了原始表的所有行,若需要全量日期+商品的统计结果可以去掉最后的右关联逻辑
运行上述代码后输出结果和你给出的预期结果完全一致。
内容的提问来源于stack exchange,提问作者delkk0
相关产品推荐
相关产品推荐

