Skip to content
Everyone knows 3 min read · Selection 4 practice problems ↓

dropDuplicates

Remove duplicate rows or keys, and why "keep the latest row" needs a window instead.

You will learn

  • The difference between distinct and dropDuplicates
  • Which row survives when you deduplicate on some columns
  • How to keep the latest row per key, deterministically
  • How deduplication works in streaming

Read first

Comfortable with these? Read on.

TL;DR distinct() removes rows that are identical in every column. dropDuplicates(["key"]) keeps one row per key, but which row is not defined. To keep the latest, use a window with row_number.

What it does

CallRemovesSQL
distinct()Rows identical in every columnSELECT DISTINCT *
dropDuplicates()Same as distinctSELECT DISTINCT *
dropDuplicates(["user_id"])Rows with the same user_id, keeping one of themNo direct equivalent; see below

Step by step

events

user_idemailupdated_at
1a@x.com2025-01-01
1a@new.com2025-03-01
2b@x.com2025-02-01
2b@x.com2025-02-01

distinct()

user_idemailupdated_at
1a@x.com2025-01-01
1a@new.com2025-03-01
2b@x.com2025-02-01

Only the exact duplicate for user 2 is removed. User 1 still has two rows because the emails differ.

Which row survives?

dropDuplicates(["user_id"]) keeps the first row it meets for each key, and "first" depends on partitioning and task timing. On a different day, with different file sizes, it can keep a different row. It is fine for exact duplicates; it is wrong for "keep the latest".

Keep the latest row per key

Rank the rows of each key by time with a window, and keep rank 1:

PySparkSpark SQL
from pyspark.sql import functions as F
from pyspark.sql.window import Window

w = Window.partitionBy("user_id").orderBy(F.col("updated_at").desc())
result = (events
    .withColumn("row_num", F.row_number().over(w))
    .filter(F.col("row_num") == 1)
    .drop("row_num"))
SELECT user_id, email, updated_at
FROM (
  SELECT *, ROW_NUMBER() OVER (
    PARTITION BY user_id ORDER BY updated_at DESC) AS rn
  FROM events)
WHERE rn = 1

Switch to PySpark to edit and run this example in your browser.

Some engines, such as Databricks SQL and Snowflake, accept QUALIFY rn = 1 to skip the subquery. The subquery form works everywhere, including open-source Spark SQL. If two rows share the latest timestamp, add a tiebreaker to the orderBy.

Under the hood

Both distinct and dropDuplicates are aggregations: Spark groups by the chosen columns and keeps one row, with the same partial aggregation before the shuffle as a groupBy. In the plan they appear as HashAggregate nodes around an Exchange.

In Structured Streaming

On a stream, dropDuplicates has to remember every key it has seen, forever, unless you bound the state with a watermark: withWatermark("event_time", "1 hour").dropDuplicates(["event_id", "event_time"]). Spark 3.5 added dropDuplicatesWithinWatermark, which deduplicates on the key alone within the watermark delay.

Common mistakes

Using dropDuplicates to keep the latest version

The surviving row is arbitrary. Use row_number over a window ordered by time.

Deduplicating on all columns when a load timestamp differs

A column like ingested_at makes every row unique, so nothing is removed. Choose the business key.

Unbounded streaming deduplication

State grows until the job fails. Add a watermark.

Key takeaways

  • distinct removes rows identical in every column.
  • dropDuplicates(cols) keeps an arbitrary row per key.
  • Use row_number over a window to keep the latest row deterministically.
  • Streaming deduplication needs a watermark to bound its state.

Check yourself

3 questions

1. Which row does dropDuplicates(["id"]) keep for a duplicated id?

Show the answer

An arbitrary one, which can change between runs. It keeps whichever row Spark meets first, which depends on partitioning. Use a window for a defined choice.

2. A table has a column loaded_at set at ingestion. What does distinct() remove?

Show the answer

Probably nothing, since loaded_at differs between loads. distinct compares every column. A per-load timestamp makes re-loaded rows unique.

3. What does a streaming dropDuplicates need to avoid unbounded state?

Show the answer

A watermark. The watermark tells Spark when old keys can be forgotten.

Practice it

Interview problems that use dropDuplicates: write the PySpark, run it, and get graded on hidden tests.

Solve: Distinct User Actions →

Go deeper

Primary sources: DataFrame.dropDuplicates · DataFrame.distinct