Spark Scala(2.3.2)如何将DataFrame的日期区间展开为新DataFrame
Spark 2.3.2 日期区间展开实现方案
针对Spark 2.3.2版本无内置日期序列生成函数的情况,可通过以下步骤将包含事件、起始/结束日期的DataFrame展开为每日一条记录:
步骤说明
- 日期类型转换:将字符串格式的日期转为Spark Date类型,便于后续日期计算
- 计算区间天数:算出起始到结束日期的总天数(包含首尾日期)
- 生成日期序列并展开:通过数字序列映射生成区间内的每一天,再炸开为单行记录
代码实现
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.DateType // 1. 创建示例原始DataFrame(模拟你的输入数据) val rawDF = spark.createDataFrame(Seq( ("event1", "01/01/2023", "04/01/2023"), ("event2", "15/02/2023", "17/02/2023") )).toDF("EVENT", "INITIAL_DATE", "END_DATE") // 2. 转换日期字符串为Date类型,计算区间天数 val dateCalcDF = rawDF .withColumn("initial_date", to_date(col("INITIAL_DATE"), "dd/MM/yyyy")) .withColumn("end_date", to_date(col("END_DATE"), "dd/MM/yyyy")) .withColumn("days_diff", datediff(col("end_date"), col("initial_date")) + 1) // 包含首尾日期,需+1 // 3. 生成数字序列并映射为每日日期,最后展开 val resultDF = dateCalcDF .withColumn("day_offset", explode(expr("sequence(0, days_diff - 1)"))) // 生成从0到天数差-1的数字序列 .withColumn("DATE", date_add(col("initial_date"), col("day_offset"))) .select("EVENT", "DATE") .withColumn("DATE", date_format(col("DATE"), "dd/MM/yyyy")) // 转回原始日期字符串格式 // 查看结果 resultDF.show()
注意事项
- 若数据量较大,优先使用内置函数方案(避免UDF带来的性能损耗)
- 确保日期格式匹配,示例中为
dd/MM/yyyy,若你的日期格式不同,需修改to_date和date_format中的格式参数 - 若遇到
sequence函数兼容问题,可改用自定义UDF生成日期数组:
备选方案(自定义UDF)
// 自定义UDF:根据起始日期和天数生成日期数组 val generateDateRange = udf((start: java.sql.Date, days: Int) => { (0 until days).map(offset => new java.sql.Date(start.getTime + offset * 86400000L)) }) val resultDFWithUDF = dateCalcDF .withColumn("date_array", generateDateRange(col("initial_date"), col("days_diff"))) .selectExpr("EVENT", "explode(date_array) as DATE") .withColumn("DATE", date_format(col("DATE"), "dd/MM/yyyy")) resultDFWithUDF.show()
内容的提问来源于stack exchange,提问作者user2083726
相关产品推荐
相关产品推荐

