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
相关产品推荐
相关产品推荐

