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

基于时间戳差异处理Spark DataFrame中的重复ID

Solution for Spark DataFrame Time-Based Deduplication

Alright, let's figure out how to solve this deduplication problem in Spark. The core requirement is: for each ID, retain the first record, and only keep subsequent records if their timestamp is more than 3 hours later than the last retained record of that ID. Here's a practical implementation:

Step 1: Prepare the DataFrame

First, make sure your time_stamp column is cast to a proper timestamp type so we can calculate time differences accurately.

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.TimestampType
import java.sql.Timestamp
import scala.collection.mutable.ArrayBuffer

// Assume your raw DataFrame is named `rawDf`
val typedDf = rawDf
  .withColumn("time_stamp", col("time_stamp").cast(TimestampType))
  .orderBy("ID", "time_stamp") // Sort records by ID and timestamp to process in chronological order

Step 2: Process Each ID Group

We'll use groupByKey combined with flatMapGroups to iterate through each ID's records, tracking the last retained timestamp and deciding whether to keep the current record:

val deduplicatedDf = typedDf
  .groupByKey(row => row.getAs[String]("ID"))
  .flatMapGroups { (id, rows) =>
    val retainedRecords = ArrayBuffer[(String, Timestamp)]()
    var lastKeptTime: Option[Timestamp] = None

    rows.foreach { row =>
      val currentTime = row.getAs[Timestamp]("time_stamp")
      lastKeptTime match {
        case None => 
          // Keep the first record of the ID
          retainedRecords += ((id, currentTime))
          lastKeptTime = Some(currentTime)
        case Some(lastTime) =>
          // Calculate time difference in milliseconds
          val timeDiffMs = currentTime.getTime - lastTime.getTime
          val threeHoursInMs = 3 * 60 * 60 * 1000 // 3 hours converted to milliseconds

          // Keep the record only if it's more than 3 hours after the last kept record
          if (timeDiffMs > threeHoursInMs) {
            retainedRecords += ((id, currentTime))
            lastKeptTime = Some(currentTime)
          }
      }
    }

    retainedRecords.iterator
  }
  .toDF("ID", "time_stamp")

Step 3: Verify the Result

If you run deduplicatedDf.show(), you'll get exactly the expected output:

+----+-------------------+
|  ID|         time_stamp|
+----+-------------------+
| 1AB|2015-01-23 08:23:16|
|25CD|2015-01-23 09:02:20|
| 1AB|2015-01-23 12:38:40|
| 1AB|2015-01-24 08:25:16|
|25CD|2015-01-24 09:01:20|
| 1AB|2015-01-24 12:30:40|
+----+-------------------+

Key Notes:

  • Sorting First: We sort the DataFrame by ID and timestamp to ensure we process each ID's records in chronological order, which is critical for the time gap check.
  • State Tracking: For each ID group, we use lastKeptTime to dynamically track the most recent retained timestamp, ensuring we only keep records that meet the 3-hour threshold relative to the last valid entry, not just the first entry.
  • Efficiency: This approach processes each group in a single pass, making it efficient even for large datasets.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:31:28