如何在Spark SQL/PySpark中合并重叠日期范围生成无重复区间
合并重叠/包含的日期区间(Spark SQL & PySpark实现)
核心思路
针对同一用户下is_relevant=true的日期区间,通过以下步骤生成无重叠的合并区间:
- 过滤出有效行(
is_relevant=true) - 按用户分组,对区间起始日期升序排序
- 用窗口函数标记不重叠的新区间组:当前区间的起始日期大于前一个区间的结束日期时,标记为新组
- 按用户和组ID聚合,取组内最小起始日期和最大结束日期,得到合并后的区间
Spark SQL实现
1. 创建示例临时表
CREATE OR REPLACE TEMP VIEW user_date_ranges AS SELECT 1 AS user, true AS is_relevant, date('2024-01-01') AS from_date, date('2024-02-01') AS to_date, 'month of january' AS description UNION ALL SELECT 1 AS user, true AS is_relevant, date('2024-01-02') AS from_date, date('2024-01-03') AS to_date, 'subset of january' AS description UNION ALL SELECT 1 AS user, true AS is_relevant, date('2024-01-15') AS from_date, date('2024-02-15') AS to_date, 'mid-jan to mid-feb' AS description UNION ALL SELECT 1 AS user, true AS is_relevant, date('2024-03-01') AS from_date, date('2024-04-01') AS to_date, 'distinct date range' AS description;
2. 核心查询语句
WITH ranked_ranges AS ( SELECT user, from_date, to_date, -- 累加标记新的区间组:当前区间与前一个不重叠时加1 SUM(CASE WHEN from_date > LAG(to_date) OVER (PARTITION BY user ORDER BY from_date) THEN 1 ELSE 0 END) OVER (PARTITION BY user ORDER BY from_date) AS group_id FROM user_date_ranges WHERE is_relevant = true ) SELECT user, MIN(from_date) AS merged_from_date, MAX(to_date) AS merged_to_date FROM ranked_ranges GROUP BY user, group_id ORDER BY user, merged_from_date;
查询结果
| user | merged_from_date | merged_to_date |
|---|---|---|
| 1 | 2024-01-01 | 2024-02-15 |
| 1 | 2024-03-01 | 2024-04-01 |
PySpark实现
完整代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark会话 spark = SparkSession.builder.appName("MergeOverlappingDateRanges").getOrCreate() # 构造示例数据 sample_data = [ (1, True, "2024-01-01", "2024-02-01", "month of january"), (1, True, "2024-01-02", "2024-01-03", "subset of january"), (1, True, "2024-01-15", "2024-02-15", "mid-jan to mid-feb"), (1, True, "2024-03-01", "2024-04-01", "distinct date range") ] # 创建DataFrame并转换日期类型 df = spark.createDataFrame(sample_data, ["user", "is_relevant", "from_date", "to_date", "description"]) df = df.withColumn("from_date", F.to_date("from_date")) \ .withColumn("to_date", F.to_date("to_date")) # 定义窗口规则:按user分组,按from_date升序排序 window_spec = Window.partitionBy("user").orderBy("from_date") # 标记新的区间组并生成group_id ranked_df = df.filter(F.col("is_relevant") == True) \ .withColumn("prev_to_date", F.lag("to_date").over(window_spec)) \ .withColumn("is_new_group", F.when(F.col("from_date") > F.col("prev_to_date"), 1).otherwise(0)) \ .withColumn("group_id", F.sum("is_new_group").over(window_spec)) # 聚合得到合并后的无重叠区间 merged_df = ranked_df.groupBy("user", "group_id") \ .agg( F.min("from_date").alias("merged_from_date"), F.max("to_date").alias("merged_to_date") ) \ .orderBy("user", "merged_from_date") # 打印结果 merged_df.show()
输出结果
+----+------------------+----------------+ |user|merged_from_date |merged_to_date | +----+------------------+----------------+ |1 |2024-01-01 |2024-02-15 | |1 |2024-03-01 |2024-04-01 | +----+------------------+----------------+
内容的提问来源于stack exchange,提问作者WhatAmIDoing
相关产品推荐
相关产品推荐

