PySpark中基于时间戳识别重叠记录并移除旧重叠项
PySpark 时间周期重叠数据去重方案
问题说明
针对同一station_id,基于start_time和end_time识别时间重叠的记录,仅保留partition_date最新的行,移除分区日期较旧的重叠行;无重叠的记录直接保留,最终实现同一station_id的任意时间段仅存在一条记录。
样例数据
data = [ (1, "2024-01-28T05:00:00Z", "2024-01-28T06:00:00Z", "1/24/24"), (1, "2024-01-28T05:30:00Z", "2024-01-28T07:00:00Z", "1/25/24"), (1, "2024-01-28T06:00:00Z", "2024-01-28T09:00:00Z", "1/24/24"), (1, "2024-01-28T07:00:00Z", "2024-01-28T10:30:00Z", "1/25/24"), (3, "2024-01-28T12:00:00Z", "2024-01-28T13:00:00Z", "1/26/24"), ] columns = ["station_id", "start_time", "end_time", "partition_date"]
期望输出
output = [ (1, "2024-01-28T05:30:00Z", "2024-01-28T07:00:00Z", "1/25/24"), (1, "2024-01-28T07:00:00Z", "2024-01-28T10:30:00Z", "1/25/24"), (3, "2024-01-28T12:00:00Z", "2024-01-28T13:00:00Z", "1/26/24"), ]
现有方案的问题
你当前的自连接方案存在以下缺陷:
- 时间重叠判断不完整:使用的OR条件未覆盖所有重叠场景(如完全包含的情况),标准的时间重叠判断应为
s1 < e2 AND s2 < e1。 - 过滤逻辑错误:
station_id2 is null or (partition_date > partition_date2)会保留旧分区中与更旧分区重叠的记录,但这类记录实际上应该被最新分区的重叠记录过滤。 - 未优先按分区日期筛选:没有先判断对方分区日期是否更晚,导致旧分区记录可能误保留。
正确实现方案
核心思路:判断每条记录是否存在同station_id下、分区日期更新且时间重叠的记录,若存在则删除该记录,否则保留。利用PySpark的left_anti连接可以高效实现这个逻辑,它会自动保留左表中无匹配右表的记录。
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("TimeOverlapDeduplication").getOrCreate() # 加载样例数据 data = [ (1, "2024-01-28T05:00:00Z", "2024-01-28T06:00:00Z", "1/24/24"), (1, "2024-01-28T05:30:00Z", "2024-01-28T07:00:00Z", "1/25/24"), (1, "2024-01-28T06:00:00Z", "2024-01-28T09:00:00Z", "1/24/24"), (1, "2024-01-28T07:00:00Z", "2024-01-28T10:30:00Z", "1/25/24"), (3, "2024-01-28T12:00:00Z", "2024-01-28T13:00:00Z", "1/26/24"), ] columns = ["station_id", "start_time", "end_time", "partition_date"] df = spark.createDataFrame(data, columns) # 将partition_date转换为日期类型,方便比较 df = df.withColumn("partition_date", F.to_date(F.col("partition_date"), "M/d/yy")) # 创建待连接的副本DataFrame,重命名列 df_dup = df.select( F.col("station_id").alias("station_id2"), F.col("start_time").alias("start_time2"), F.col("end_time").alias("end_time2"), F.col("partition_date").alias("partition_date2") ) # 定义连接条件:同station_id、对方分区更新、时间区间重叠 join_condition = ( df["station_id"] == df_dup["station_id2"] & df_dup["partition_date2"] > df["partition_date"] & (df["start_time"] < df_dup["end_time2"]) & (df_dup["start_time2"] < df["end_time"]) ) # 左反连接:保留无匹配的记录(即无更新分区重叠的记录) result_df = df.join(df_dup, join_condition, "left_anti") # 将partition_date转换回原字符串格式 result_df = result_df.withColumn("partition_date", F.date_format(F.col("partition_date"), "M/d/yy")) # 查看结果 result_df.show(truncate=False)
代码说明
- 日期转换:将
partition_date转为日期类型,避免字符串比较的误差。 - 自连接条件:
- 限制同一
station_id - 仅匹配分区日期更新的记录
- 使用标准时间重叠判断逻辑:两个区间
[s1,e1]和[s2,e2]重叠当且仅当s1 < e2且s2 < e1
- 限制同一
- 左反连接:自动过滤掉存在更新分区重叠记录的行,仅保留符合要求的记录。
结果验证
运行代码后输出结果与期望完全一致,同时覆盖了所有边缘场景:
- 完全包含的重叠记录
- 交叉重叠的记录
- 无重叠的独立记录
内容的提问来源于stack exchange,提问作者Syed Sharjeelullah
相关产品推荐
相关产品推荐

