Databricks中基于特定分组逻辑补全上月缺失store_id记录
解决PySpark中月度销售数据缺失门店分组记录的补全问题
需求回顾
我们有两个Spark DataFrame:
current_month_data:当月销售数据,包含store_id、visit_month、visit_date、type、category、brand、segment、amount及其他销售指标字段previous_month_data:上月销售数据,字段与当月一致
需要补全的逻辑:上月中存在的(store_id, type, category, brand, segment, product_hierarchy_type)唯一组合,若在本月完全缺失,则将该组合对应的上月完整记录复制到本月数据中,最终输出包含所有字段的完整DataFrame。
实现步骤
1. 定义分组标识字段
先明确用来判断“缺失”的核心字段组合:
# 定义分组键:这些字段的唯一组合用来判断是否需要补全 group_keys = ["store_id", "type", "category", "brand", "segment", "product_hierarchy_type"]
2. 获取本月已存在的分组组合
提取当月所有已有的分组,避免重复导入:
current_existing_groups = current_month_data.select(group_keys).distinct()
3. 筛选上月中本月缺失的分组记录
用左反连接(left_anti)从上月数据中找出所有不在本月分组里的记录:
# 左反连接会返回上月数据中,分组组合未出现在本月的所有记录 missing_records_from_prev = previous_month_data.join( current_existing_groups, on=group_keys, how="left_anti" )
4. 调整补全记录的月份字段
因为补全的记录要归入本月,需要更新visit_month字段(如果visit_date需要同步调整,可根据业务规则修改):
from pyspark.sql.functions import lit # 从当月数据中提取当前月份(避免硬编码) current_month_value = current_month_data.select("visit_month").distinct().collect()[0][0] # 更新缺失记录的月份为当前月 updated_missing_records = missing_records_from_prev.withColumn("visit_month", lit(current_month_value)) # 可选:如果需要将visit_date改为当月第一天,可添加以下代码 # from pyspark.sql.functions import to_date # updated_missing_records = updated_missing_records.withColumn( # "visit_date", # to_date(lit(f"{current_month_value}-01"), "yyyy-MM-dd") # )
5. 合并数据得到最终结果
将当月原始数据和补全的记录合并:
# unionByName确保按字段名合并,避免因列顺序不一致出问题 final_complete_data = current_month_data.unionByName(updated_missing_records, allowMissingColumns=False)
验证逻辑(可选)
可以通过以下方式验证补全是否正确:
- 对比补全前后的分组数:
final_complete_data.select(group_keys).distinct().count()应等于current_existing_groups.count() + missing_records_from_prev.select(group_keys).distinct().count() - 抽查特定缺失的
store_id对应的分组记录是否出现在最终数据中
内容的提问来源于stack exchange,提问作者Chinnu
相关产品推荐
相关产品推荐

