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

Spark按日期分组并补全缺失ID零值的30天数据集转换方案咨询

Spark数据集按日期分组并补全缺失ID零值(覆盖30天)

这是个典型的时间序列数据补全需求,要实现覆盖30天、每个日期下所有ID都有记录(缺失则填0)的目标,咱们可以分五步来做,用Spark的DataFrame API就能高效完成:

步骤1:准备基础数据与日期范围

首先加载示例数据,把日期列转成标准Date类型,然后确定要覆盖的30天范围——从原数据的最小日期开始,往后推29天(包含起始日共30天):

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._
import java.sql.Date
import java.time.{LocalDate, Period}

// 加载示例数据并转换日期格式
val genre = sc.parallelize(List(
  ("id1", "2016-05-01", "action",0), 
  ("id1", "2016-05-03", "horror",1), 
  ("id2", "2016-05-03", "art",0), 
  ("id2", "2016-05-04", "action",0)
)).toDF("id","date","genre","score")
val genreDF = genre.withColumn("date", to_date(col("date")))

// 计算30天周期的起始和结束日期
val startDate = genreDF.select(min("date")).first().getAs[Date](0).toLocalDate
val endDate = startDate.plus(Period.ofDays(29))

// 生成完整的30天日期序列
val dateRange = (0 to 29).map(days => Date.valueOf(startDate.plusDays(days)))
val dateDF = spark.createDataFrame(dateRange.map(d => Tuple1(d))).toDF("date")

步骤2:提取所有唯一ID

拿到数据里的所有唯一ID,后续用来生成日期与ID的全量组合:

val uniqueIds = genreDF.select("id").distinct().collect().map(_.getAs[String](0))
val idDF = spark.createDataFrame(uniqueIds.map(id => Tuple1(id))).toDF("id")

步骤3:生成日期与ID的笛卡尔积

这一步是补全缺失记录的核心,通过笛卡尔积得到每个日期下所有ID的完整组合,确保没有遗漏:

val dateIdCross = dateDF.crossJoin(idDF)

步骤4:左连接原数据并填充缺失值

把全量的(date, id)组合和原数据左连接,然后将缺失的genre和score字段填充为0:

val filledDF = dateIdCross
  .join(genreDF, Seq("date", "id"), "left_outer")
  .withColumn("genre", when(col("genre").isNull, lit("0")).otherwise(col("genre")))
  .withColumn("score", when(col("score").isNull, lit(0)).otherwise(col("score")))

步骤5:按日期分组聚合

最后按日期分组,把每个日期下的(id, genre, score)聚合成数组,得到你想要的格式:

val resultDF = filledDF
  .groupBy("date")
  .agg(collect_list(struct("id", "genre", "score")).alias("grouped"))
  .orderBy("date")

// 查看结果
resultDF.show(truncate = false)

结果验证

运行后得到的结果和你期望的一致(仅展示前4天示例):

+----------+----------------------------------------+
|date      |grouped                                 |
+----------+----------------------------------------+
|2016-05-01|[[id1, action, 0], [id2, 0, 0]]         |
|2016-05-02|[[id1, 0, 0], [id2, 0, 0]]              |
|2016-05-03|[[id1, horror, 1], [id2, art, 0]]       |
|2016-05-04|[[id1, 0, 0], [id2, action, 0]]         |
+----------+----------------------------------------+

优化小技巧

  • 如果ID数量较少,可以用broadcast(idDF)优化笛卡尔积的shuffle过程,提升性能;
  • Spark 2.4+版本可以用sequence函数更简洁地生成日期序列:
    val dateDF = spark.sql(s"SELECT sequence(to_date('${startDate.toString}'), to_date('${endDate.toString}'), interval 1 day) as date")
      .select(explode(col("date")).alias("date"))
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:51:23