Spark中如何按id分组合并重叠时间区间得到无交集连续区间
Spark按ID分区合并重叠日期区间实现方案
实现思路
- 按id分区,所有区间按开始日期升序排序,保证可以从上到下依次判断区间重叠关系
- 用窗口函数取当前行的上一行结束日期,判断当前行开始日期是否大于上一行结束日期:如果是说明两个区间不重叠,新增一个分组标记,否则和上一个区间同组
- 对分组标记做累加,相同累加值的区间属于同一个需要合并的组
- 按id和分组累加值分组,取组内最小开始日期作为新的区间开始,最大结束日期作为新的区间结束,得到无重叠的合并结果
代码实现
Scala DataFrame API 实现
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 替换为你的原始DataFrame val rawDf = ??? // 日期类型转换,如果原始字段已经是Date/Timestamp类型可跳过此步 val dateFormatDf = rawDf .withColumn("start", to_date(col("start"))) .withColumn("end", to_date(col("end"))) // 定义按id分区、按start升序排序的窗口 val idPartitionWindow = Window.partitionBy("id").orderBy("start") // 生成分组标记 val groupedDf = dateFormatDf // 取同id下上一行的结束日期 .withColumn("last_end", lag("end", 1).over(idPartitionWindow)) // 不重叠则标记为1,重叠标记为0 .withColumn("split_flag", when(col("last_end").isNull || col("start") > col("last_end"), 1).otherwise(0)) // 累加标记得到最终分组ID .withColumn("group_id", sum("split_flag").over(idPartitionWindow.rowsBetween(Window.unboundedPreceding, 0))) // 合并同组区间得到最终结果 val resultDf = groupedDf .groupBy("id", "group_id") .agg( min("start").alias("start"), max("end").alias("end") ) .drop("group_id")
Spark SQL 实现
WITH get_last_end AS ( SELECT id, start, end, LAG(end, 1) OVER(PARTITION BY id ORDER BY start) AS last_end FROM 你的原始表名 ), gen_group_id AS ( SELECT id, start, end, SUM(CASE WHEN last_end IS NULL OR start > last_end THEN 1 ELSE 0 END) OVER( PARTITION BY id ORDER BY start ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS group_id FROM get_last_end ) SELECT id, MIN(start) AS start, MAX(end) AS end FROM gen_group_id GROUP BY id, group_id
注意事项
- 如果业务规则里认为日期连续也算重叠(比如上一个区间结束是2023-01-05,当前区间开始是2023-01-06需要合并),将判断条件中的
start > last_end修改为start > date_add(last_end, 1)即可 - 该逻辑天然支持多id并行处理,每个id的区间合并完全独立,不会互相干扰
- 支持任意数量的重叠区间嵌套场景,合并后不会出现区间重叠或者覆盖遗漏的问题
内容的提问来源于stack exchange,提问作者Miguel
相关产品推荐
相关产品推荐

