如何利用周度数据集高效补全日度访客数据缺失值?
日度访客数据补全的Spark高效实现优化
问题背景
现有两份城市访客数据集:
- Dataset A:日度数据,部分日期访客数为0(视为缺失)
- Dataset B:周度汇总数据,日期均为每周周六
需求:当Dataset A中存在周日至周六连续7天访客数均为0时,用Dataset B对应周六的周度访客数除以7补全该周每日数据。
数据集示例
Dataset A(原始日度数据)
City | Date | visitors A | 9/1/22 | 10 A | 9/2/22 | 12 A | 9/3/22 | 20 A | 9/4/22 | 0 A | 9/5/22 | 0 A | 9/6/22 | 0 A | 9/7/22 | 0 A | 9/8/22 | 0 A | 9/9/22 | 0 A | 9/10/22 | 0 A | 9/11/22 | 19
Dataset B(周度汇总数据)
City | Date | visitors_weekly A | 9/3/22 | 30 A | 9/10/22 | 50 A | 9/17/22 | 25
补全后Dataset A示例
City | Date | visitors A | 9/1/22 | 10 A | 9/2/22 | 12 A | 9/3/22 | 20 A | 9/4/22 | 7.143 A | 9/5/22 | 7.143 A | 9/6/22 | 7.143 A | 9/7/22 | 7.143 A | 9/8/22 | 7.143 A | 9/9/22 | 7.143 A | 9/10/22 | 7.143 A | 9/11/22 | 19
当前实现代码
# step 1: find the week of the year we're in, factoring in when we loop into next year (the otherwise condition) # We only deal with data 365d out maximum # adj_weekOfYear adjusts it so that we start with today as the first week of the year instead of Jan adj_weekOfYear = F.weekofyear(F.col("date")+1) - F.weekofyear(F.current_date()+1) dataset_A = dataset_A.withColumn( "week_of_year", F.when( (adj_weekOfYear >= 0) & (F.year(F.col("date")) == F.year(F.current_date())), adj_weekOfYear ).otherwise(adj_weekOfYear + 52) ) dataset_B = dataset_B.withColumn( "week_of_year", F.when( (adj_weekOfYear >= 0) & (F.year(F.col("date")) == F.year(F.current_date())), adj_weekOfYear ).otherwise(adj_weekOfYear + 52) ) df = dataset_A.join( dataset_B.drop("date"), ["city", "week_of_year"], 'left' ).withColumn( # in dataset A, if weeks_daily_visitors is null then we can tell we're missing that week's info "weeks_daily_visitors", F.sum(F.col("visitors")).over(Window.partitionBy("city", "week_of_year")) ).withColumn( # edge case for beg and end of data, want to make sure we have a full week sun-sat before filling in with dataset B "num_rows_for_week", F.count("date").over(Window.partitionBy("city", "week_of_year")) ).withColumn( "missing_weeks_daily_visitors", F.when( # if we dont have all 7 days of the week, cant fill in F.col("num_rows_for_week") < 7, F.lit("No") ).when( (F.col("weeks_daily_visitors").isNull() | (F.col("weeks_daily_visitors") < 1)), # if we're missing the week's data in A we should mark to fill in with B F.lit("Yes") ).otherwise(F.lit("No")) ).withColumn( "visitors", F.when( F.col("missing_weeks_daily_demand") == 'Yes', F.col("visitors_weekly") / 7 ).otherwise(F.col("demand")) )
高效优化方案
优化思路
- 简化周关联逻辑:直接通过日期计算每条日度数据对应的周六(Dataset B的日期标识),避免自定义周数计算的误差与复杂度。
- 合并窗口操作:将原有的两次窗口计算(总和、行数)合并为一次,减少Shuffle开销。
- 减少中间列:直接在最终列计算中判断补全条件,避免冗余标记列占用内存。
优化后代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 为Dataset A的每条数据计算对应周的周六(Dataset B的日期格式) # Spark中dayofweek返回1=周日,7=周六,因此计算需要补的天数:6 - dayofweek(date) +1 dataset_A = dataset_A.withColumn( "week_end_date", F.date_add(F.col("Date"), 6 - F.dayofweek(F.col("Date")) + 1) ) # 2. 重命名Dataset B的日期列为关联键 dataset_B = dataset_B.withColumnRenamed("Date", "week_end_date") # 3. 按城市和周结束日期关联两个数据集 df = dataset_A.join( dataset_B.select("City", "week_end_date", "visitors_weekly"), ["City", "week_end_date"], "left" ) # 4. 定义窗口:按城市和周结束日期分区 week_window = Window.partitionBy("City", "week_end_date") # 5. 一次窗口操作计算周访客总和与数据行数,同时完成补全逻辑 df = df.withColumn( "week_total_visitors", F.sum("visitors").over(week_window) ).withColumn( "week_row_count", F.count("Date").over(week_window) ).withColumn( "visitors", F.when( # 满足条件:该周有7条数据,且总访客数为0 (F.col("week_row_count") == 7) & (F.col("week_total_visitors") == 0), F.round(F.col("visitors_weekly") / 7, 3) # 保留3位小数与示例一致 ).otherwise(F.col("visitors")) ).drop("week_end_date", "week_total_visitors", "week_row_count", "visitors_weekly") # 可选:按日期排序(如果需要) df = df.orderBy("City", "Date")
优化点说明
- 避免周数计算误差:通过实际日期推导周六,不依赖当前日期或自定义周数,逻辑更稳定,不会因跨年、周起始规则变化出错。
- 降低Shuffle开销:原代码两次窗口操作会触发两次Shuffle,优化后合并为一次,减少集群资源消耗。
- 减少内存占用:去掉了
missing_weeks_daily_visitors等中间标记列,直接在目标列计算中判断条件,降低数据处理过程中的内存压力。 - 更直观的逻辑:关联条件基于实际日期,可读性更强,便于后续维护。
内容的提问来源于stack exchange,提问作者user90823745
相关产品推荐
相关产品推荐

