如何在PySpark DataFrame中按Store统计最新日期的有效SKU去重数量?
实现方案
可以通过两种方式实现需求,核心思路是先筛选出每个Store的最新日期记录,再统计该日期下Flag=1的去重SKU数量:
方法一:窗口函数方式
利用窗口函数标记每个Store的最新日期记录,再进行统计:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按Store分组,日期降序排序 window_spec = Window.partitionBy("Store").orderBy(F.desc("Date")) # 筛选每个Store的最新日期记录 latest_records = df.withColumn("rank", F.row_number().over(window_spec)) \ .filter(F.col("rank") == 1) \ .drop("rank") # 按Store分组,统计Flag=1的去重SKU数(无符合项时显示0) result = latest_records.groupBy("Store") \ .agg( F.coalesce( F.countDistinct(F.when(F.col("Flag") == 1, F.col("SKU"))), F.lit(0) ).alias("count_distinct_sku") ) # 查看结果 result.show()
方法二:分组取最大日期后关联
先计算每个Store的最新日期,再关联原表筛选对应记录,最后统计:
from pyspark.sql import functions as F # 计算每个Store的最新日期 store_max_date = df.groupBy("Store").agg(F.max("Date").alias("max_date")) # 关联原表,筛选每个Store最新日期的记录 latest_records = df.join(store_max_date, on=["Store"], how="inner") \ .filter(F.col("Date") == F.col("max_date")) \ .drop("max_date") # 统计去重SKU数,逻辑同方法一 result = latest_records.groupBy("Store") \ .agg( F.coalesce( F.countDistinct(F.when(F.col("Flag") == 1, F.col("SKU"))), F.lit(0) ).alias("count_distinct_sku") ) # 查看结果 result.show()
两种方法最终都会输出符合需求的结果:
| Store | count_distinct_sku |
|---|---|
| 629138 | 1 |
| 187367 | 0 |
| 176129 | 0 |
| 633782 | 2 |
内容的提问来源于stack exchange,提问作者Scope
相关产品推荐
相关产品推荐

