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

Spark Scala实现:补全日期区间内所有周日期并转换DataFrame

Spark Scala: Fill Missing Weekly Dates per User

Got it, let's work through this problem step by step. The core task is to generate rows for every weekly date within each user's date range, while aggregating duplicate date amounts and filling in 0 for dates with no data.

Step 1: Setup & Data Preparation

First, let's start by importing necessary Spark functions and defining our input DataFrame. We'll convert string dates to proper DateType to simplify date calculations.

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.DateType
import org.apache.spark.sql.Window

val spark = SparkSession.builder()
  .appName("WeeklyDateFill")
  .master("local[*]")
  .getOrCreate()

// Sample input DataFrame
val inputDF = spark.createDataFrame(Seq(
  ("Jhon", "4/6/2018", 100),
  ("Jhon", "4/6/2018", 200),
  ("Jhon", "4/13/2018", 300),
  ("Jhon", "4/20/2018", 500),
  ("Lee", "5/4/2018", 100),
  ("Lee", "4/4/2018", 200),
  ("Lee", "5/4/2018", 300),
  ("Lee", "4/11/2018", 700)
)).toDF("name", "date", "amount")

// Convert string date to DateType (format: MM/dd/yyyy)
val formattedDF = inputDF.withColumn("date", to_date(col("date"), "MM/dd/yyyy"))

Step 2: Aggregate Duplicate Dates

Since the input has multiple entries for the same user and date, we'll first sum the amounts for those duplicates to get a clean, single row per user-date pair:

val aggregatedDF = formattedDF
  .groupBy("name", "date")
  .agg(sum("amount").as("amount"))

Step 3: Generate Weekly Date Ranges per User

Next, we'll calculate the earliest and latest date for each user, then generate a sequence of weekly dates between those two points. We use Spark's sequence function with a 7-day interval to create the weekly date list:

// Define window to get min/max date per user
val userWindow = Window.partitionBy("name")

val dateRangeDF = aggregatedDF
  .withColumn("min_date", min("date").over(userWindow))
  .withColumn("max_date", max("date").over(userWindow))
  // Generate weekly date sequence from min to max date
  .withColumn("weekly_dates", sequence(col("min_date"), col("max_date"), expr("interval 7 days")))
  // Explode the sequence to get individual rows for each weekly date
  .select("name", explode("weekly_dates").as("date"))
  .dropDuplicates() // Remove duplicate dates (in case of overlapping ranges)

Step 4: Join & Fill Missing Amounts

Now we'll join the generated weekly dates with our aggregated data, and fill in 0 for any dates that don't have corresponding amount data using coalesce:

val finalDF = dateRangeDF
  .join(aggregatedDF, Seq("name", "date"), "left_outer")
  .withColumn("amount", coalesce(col("amount"), lit(0)))
  .orderBy("name", "date")

finalDF.show()

Expected Output

When you run the code above, you'll get this result:

+----+----------+------+
|name|      date|amount|
+----+----------+------+
|Jhon|2018-04-06|   300|
|Jhon|2018-04-13|   300|
|Jhon|2018-04-20|   500|
|Lee |2018-04-04|   200|
|Lee |2018-04-11|   700|
|Lee |2018-04-18|     0|
|Lee |2018-05-04|   400|
+----+----------+------+

Key Notes

  • The sequence function requires Spark 2.4+ – if you're on an older version, you can use a UDF to generate the weekly dates instead.
  • coalesce ensures we replace null amounts (from the left join) with 0, which is exactly what we need for missing dates.
  • The dropDuplicates step ensures we don't get duplicate date rows for a single user.

内容的提问来源于stack exchange,提问作者veeta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:07:34