基于时间戳差异处理Spark DataFrame中的重复ID
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
lastKeptTimeto 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

