You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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"),
]

现有方案的问题

你当前的自连接方案存在以下缺陷:

  1. 时间重叠判断不完整:使用的OR条件未覆盖所有重叠场景(如完全包含的情况),标准的时间重叠判断应为s1 < e2 AND s2 < e1。
  2. 过滤逻辑错误:station_id2 is null or (partition_date > partition_date2)会保留旧分区中与更旧分区重叠的记录,但这类记录实际上应该被最新分区的重叠记录过滤。
  3. 未优先按分区日期筛选:没有先判断对方分区日期是否更晚,导致旧分区记录可能误保留。

正确实现方案

核心思路:判断每条记录是否存在同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)

代码说明

  1. 日期转换:将partition_date转为日期类型,避免字符串比较的误差。
  2. 自连接条件:
    • 限制同一station_id
    • 仅匹配分区日期更新的记录
    • 使用标准时间重叠判断逻辑:两个区间[s1,e1]和[s2,e2]重叠当且仅当s1 < e2且s2 < e1
  3. 左反连接:自动过滤掉存在更新分区重叠记录的行,仅保留符合要求的记录。

结果验证

运行代码后输出结果与期望完全一致,同时覆盖了所有边缘场景:

  • 完全包含的重叠记录
  • 交叉重叠的记录
  • 无重叠的独立记录

内容的提问来源于stack exchange,提问作者Syed Sharjeelullah

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 12:50:25