在Databricks中基于特定分组逻辑向前复制缺失的store_id月度记录
PySpark处理月度销售数据缺失门店记录补全方案
核心逻辑
- 定位上月存在但当月未出现的
store_id及其对应的type/category/brand/segment/product_hierarchy_type唯一分组组合 - 将该组合对应的上月完整记录复制,更新
visit_month和visit_date为当月值 - 合并原始当月数据与补全后的记录,完成数据补全
代码实现
1. 数据预处理与月份筛选
假设历史销售数据DataFrame为sales_df,先统一月份格式并提取上月、当月数据:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 从visit_date提取标准格式的visit_month(yyyy-MM) sales_df = sales_df.withColumn("visit_month", F.date_format(F.to_date("visit_date"), "yyyy-MM")) # 获取所有月份并排序,取最新两个月份(当月、上月) months = sorted([row.visit_month for row in sales_df.select("visit_month").distinct().collect()]) current_month = months[-1] previous_month = months[-2] # 筛选上月、当月数据子集 prev_month_df = sales_df.filter(F.col("visit_month") == previous_month) current_month_df = sales_df.filter(F.col("visit_month") == current_month)
2. 识别缺失的门店-分组组合
找出上月存在但当月未覆盖的store_id+分组字段组合:
# 提取上月的门店-分组唯一组合 prev_store_groups = prev_month_df.select( "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ).distinct() # 提取当月的门店-分组唯一组合 current_store_groups = current_month_df.select( "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ).distinct() # 左反连接得到缺失的组合 missing_store_groups = prev_store_groups.join( current_store_groups, on=["store_id", "type", "category", "brand", "segment", "product_hierarchy_type"], how="left_anti" )
3. 复制上月记录并更新月份信息
关联上月数据获取完整记录,更新visit_month和visit_date为当月值(示例保留上月日份,如上月2023-10-15改为当月2023-11-15):
# 获取需要补全的完整记录 records_to_copy = missing_store_groups.join( prev_month_df, on=["store_id", "type", "category", "brand", "segment", "product_hierarchy_type"], how="inner" ) # 更新月份和日期字段 records_to_copy = records_to_copy.withColumn( "visit_month", F.lit(current_month) ).withColumn( "visit_date", F.date_add( F.trunc(F.to_date(F.lit(current_month + "-01")), "month"), F.dayofmonth(F.to_date("visit_date")) - 1 ) )
4. 合并数据完成补全
将补全记录与原始当月数据合并,可选择更新回历史DataFrame:
# 合并当月原始数据与补全记录 final_current_month_df = current_month_df.unionByName(records_to_copy) # 更新全量历史数据 updated_sales_df = sales_df.filter(F.col("visit_month") != current_month).unionByName(final_current_month_df)
5. 全量历史数据迭代补全(可选)
若需要对2023年至今的所有月份依次补全,可通过循环迭代处理:
months_sorted = sorted([row.visit_month for row in sales_df.select("visit_month").distinct().collect()]) processed_df = sales_df.filter(F.col("visit_month") == months_sorted[0]) for i in range(1, len(months_sorted)): curr_month = months_sorted[i] prev_month = months_sorted[i-1] prev_processed = processed_df.filter(F.col("visit_month") == prev_month) curr_original = sales_df.filter(F.col("visit_month") == curr_month) # 重复缺失组合识别、记录复制逻辑 prev_store_groups = prev_processed.select( "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ).distinct() curr_store_groups = curr_original.select( "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ).distinct() missing_groups = prev_store_groups.join(curr_store_groups, on=[ "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ], how="left_anti") records_to_copy = missing_groups.join(prev_processed, on=[ "store_id", "type", "category", "brand", "segment", "product_hierarchy_type" ], how="inner").withColumn( "visit_month", F.lit(curr_month) ).withColumn( "visit_date", F.date_add( F.trunc(F.to_date(F.lit(curr_month + "-01")), "month"), F.dayofmonth(F.to_date("visit_date")) - 1 ) ) curr_processed = curr_original.unionByName(records_to_copy) processed_df = processed_df.unionByName(curr_processed)
注意事项
visit_date的更新逻辑可根据业务调整,比如统一用当月第一天/最后一天,而非保留上月日份- 若销售指标(如销售额、销量)需要置零而非保留上月值,可在复制记录时添加
withColumn("sales_amount", F.lit(0))这类操作 - 处理超大规模数据集时,避免用
collect()获取月份列表,可改用Spark窗口函数row_number()对月份排序后提取
内容的提问来源于stack exchange,提问作者Chinnu
相关产品推荐
相关产品推荐

