Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Deduplicating with window vs dropDuplicates

PySpark · DataFrame API in Practice

Deduplicating with window vs dropDuplicates

Mediumpyspark-83
deduplicationrow-numberdropduplicatesstreaming

Question

When should you use row_number() for deduplication instead of dropDuplicates?

Solution

dropDuplicates keeps one row per key and has no say in which one. A row_number() window lets you choose, for example the latest by timestamp, and gives the same answer on every run.

dropDuplicates

df.dropDuplicates(["order_id"])

For each order_id, Spark keeps some row. Which one depends on partition layout and the order rows arrive in, so it can differ from run to run. That is fine if the duplicates are exact copies. It is wrong if the copies differ, for example an older status and a newer status, because you may keep the stale one.

Deterministic: row_number

w = Window.partitionBy("order_id").orderBy(F.col("updated_at").desc(), F.col("ingested_at").desc())

latest = (df.withColumn("rn", F.row_number().over(w))
            .filter("rn = 1")
            .drop("rn"))

You control the winner: the latest updated_at, with ingested_at as a tie-breaker. If two rows are fully tied, add one more column to the order so the result is stable. The cost is a shuffle and a sort, the same as dropDuplicates, so there is no big price for the control.

In Structured Streaming

dropDuplicates on a stream has to remember the keys it has already seen. Without a bound, that state grows forever. Use a watermark together with an event-time column in the key, so Spark can drop old state:

(stream.withWatermark("event_ts", "1 hour")
       .dropDuplicates(["event_id", "event_ts"]))

Spark 3.5 added dropDuplicatesWithinWatermark, which removes duplicates within the watermark delay without needing the timestamp to be identical, which fits cases where duplicates arrive with slightly different event times.

A row_number window is not supported on a plain stream the same way. For "latest per key" in streaming, write each micro-batch with foreachBatch and MERGE into the target.

Rule of thumb

Exact duplicates, and you do not care which: dropDuplicates. Versions of the same record: row_number with an explicit order. Streaming: watermarked dedup or MERGE.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext