You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何利用周度数据集高效补全日度访客数据缺失值?

日度访客数据补全的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"))
)

高效优化方案

优化思路

  1. 简化周关联逻辑:直接通过日期计算每条日度数据对应的周六(Dataset B的日期标识),避免自定义周数计算的误差与复杂度。
  2. 合并窗口操作:将原有的两次窗口计算(总和、行数)合并为一次,减少Shuffle开销。
  3. 减少中间列:直接在最终列计算中判断补全条件,避免冗余标记列占用内存。

优化后代码

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 12:03:30