What you will be able to do
- Remove rows that are identical in every column with distinct() or dropDuplicates()
- Deduplicate on a chosen subset of columns and predict which rows and columns remain
- Limit deduplication state in a streaming query with withWatermark and dropDuplicatesWithinWatermark
Key concept
Duplicate comparison columns — Spark treats two rows as duplicates when they match on a set of comparison columns. By default that set is every column; you can also name a subset. The deduplication methods differ mainly in which columns they compare and in how long a streaming query remembers the rows it has already seen.
1.Exact duplicates with distinct()
Duplicate rows get into a pipeline in ordinary ways: a load is retried, two extracts overlap, or an upstream system sends the same record again. The simplest tool for removing them is distinct(). It takes no arguments and "Returns a new DataFrame containing the distinct rows in this DataFrame." Two rows count as duplicates only when they match in every column. Like other DataFrame transformations, it returns a new DataFrame and leaves the original as it was.
df = spark.createDataFrame(
[(14, "Tom"), (23, "Alice"), (23, "Alice")], ["age", "name"])
df.distinct().show()
# +---+-----+
# |age| name|
# +---+-----+
# | 14| Tom|
# | 23|Alice|
# +---+-----+On the same DataFrame, the reference shows df.distinct().count() returning 2, so you can compare counts before and after to check whether a dataset contains duplicates. The limitation is that every column has to match. If one column differs between two copies of the same record, such as a load timestamp, distinct() keeps both rows.
Two. The ingested_at values differ, so the rows don't match in every column, and distinct() only removes rows that match in every column. You need a way to compare only some of the columns, which is what dropDuplicates() provides.
Checkpoint 1 of 5· Check yourself
df contains the rows (14, 'Tom'), (23, 'Alice') and (23, 'Alice'). What does df.distinct().count() return?
distinct() compares all columns and takes no arguments. The two (23, 'Alice') rows are identical, so two distinct rows are left.
“Returns a new DataFrame containing the distinct rows in this DataFrame.”Source: docs.databricks.com
Sources1
2.Choosing the comparison columns with dropDuplicates()
dropDuplicates(subset=None) lets you choose the columns that are compared. Its description reads: "Return a new DataFrame with duplicate rows removed, optionally only considering certain columns." The single parameter, subset, is a list of column names. When you leave it out, all columns are compared, which is the same rule distinct() uses.
df = spark.createDataFrame([
Row(name='Alice', age=5, height=80),
Row(name='Alice', age=5, height=80),
Row(name='Alice', age=10, height=80)
])
df.dropDuplicates().show()
# +-----+---+------+
# | name|age|height|
# +-----+---+------+
# |Alice| 5| 80|
# |Alice| 10| 80|
# +-----+---+------+df.dropDuplicates(['name', 'height']).show()
# +-----+---+------+
# | name|age|height|
# +-----+---+------+
# |Alice| 5| 80|
# +-----+---+------+There are two things to see in this output. First, the subset changes which rows count as duplicates, but it doesn't remove any columns: age is still there. Second, the row with age=10 is gone, even though no other row had that age, because it matched the others on name and height. The reference example happens to keep the age=5 row. It doesn't say which row is kept in general, so don't write logic that depends on a particular row surviving. In the order example above, dropDuplicates(['order_id']) would treat the two loads as one record and ignore ingested_at.
Checkpoint 2 of 5· Fill the gap
Which method completes this call so that rows are deduplicated using only the name and height columns?
df. ? (['name', 'height']).show()dropDuplicates accepts a subset of columns to compare. distinct() takes no columns, dropna removes rows with nulls, and dropDuplicatesWithinWatermark only works on streaming DataFrames.
Source: docs.databricks.comCheckpoint 3 of 5· Exam question
You have a DataFrame `orders_df` with columns `order_id`, `customer_id`, `amount`. A teammate wants an exact drop of duplicate rows considering every column and writes `orders_df.distinct()`. You need to rewrite this using `dropDuplicates()` so the result is functionally identical. Which call is equivalent?
Correct answer: A — orders_df.dropDuplicates()
- A. Calling dropDuplicates() with no subset compares every column, exactly matching the full-row comparison that distinct() performs.
- B. Restricting the subset to order_id only removes rows that share an order_id, which is narrower than comparing all columns and can drop rows distinct() would keep.
- C. Narrowing the comparison to customer_id alone collapses multiple genuinely different orders from the same customer, producing far fewer rows than a full-row distinct().
- D. na.drop() removes rows containing null values; it performs no deduplication of identical rows at all.
- E. Projecting down to a single column before calling distinct() discards the other columns, so the result cannot match the schema or row set of a full-row distinct().
Sources2
3.Deduplicating streams: state and watermarks
On a batch DataFrame the work is done once: "For a static batch DataFrame, it just drops duplicate rows." A stream is different, because a duplicate can arrive in any later micro-batch. To catch it, Spark has to remember what it has already seen: "For a streaming DataFrame, it will keep all data across triggers as intermediate state to drop duplicates rows." distinct() behaves the same way, and the watermarks guide warns: "Without a watermark, state grows indefinitely and can cause memory issues."
The fix is to call withWatermark on an event-time column before deduplicating. The watermark limits how late a duplicate can arrive, so Spark can discard old state. It also has a cost: "data older than watermark will be dropped to avoid any possibility of duplicates." Choosing a watermark is therefore a trade-off between how much state you keep and how much late data you lose.
In real event streams, copies of an event are often not identical. A retried event carries the same guid but a different arrival time. dropDuplicatesWithinWatermark is designed for this case. You "Use dropDuplicatesWithinWatermark to remove duplicates on any field, even when fields differ across duplicate records, such as event time or arrival time." It has two hard requirements: "This only works with streaming DataFrame, and watermark for the input DataFrame must be set via withWatermark."
streamingDf = spark.readStream. ...
# deduplicate using guid column with watermark based on eventTime column
(streamingDf
.withWatermark("eventTime", "10 hours")
.dropDuplicatesWithinWatermark(["guid"])
)Duplicates that arrive within the watermark threshold are always removed. Duplicates that arrive outside it might be removed, but that isn't guaranteed. So size the threshold based on how far apart your duplicates can be. You can't assume deduplication happens for you elsewhere either: "Structured Streaming guarantees exactly-once processing but doesn't deduplicate records from data sources."
| Method | Columns compared | Where it works | Watermark |
|---|---|---|---|
| distinct() | All columns only | Batch and streaming | Optional in streaming, but without it state grows indefinitely |
| dropDuplicates(subset) | All columns, or the subset you name | Batch and streaming | Optional in streaming; limits state and drops late data |
| dropDuplicatesWithinWatermark(subset) | All columns, or the subset you name | Streaming only | Required, set via withWatermark |
Checkpoint 4 of 5· Check yourself
Duplicate events that share a guid can arrive up to 3 hours apart. With dropDuplicatesWithinWatermark(['guid']), which watermark guarantees that every duplicate is dropped?
The guarantee only holds when the threshold is greater than the largest time gap between duplicates. Here that gap is 3 hours, so 4 hours is the only option that qualifies. The method also refuses to run without a watermark.
“To guarantee queries drop all duplicates, set the watermark threshold to be greater than the maximum timestamp difference between duplicate events.”Source: docs.databricks.com
Checkpoint 5 of 5· Exam question
An `events_df` DataFrame has columns `device_id`, `event_time`, `payload`. Duplicate device readings should be deduplicated so that the **earliest** `event_time` per `device_id` is retained, and the result must be deterministic across repeated runs on unchanged data. Which code correctly implements this?
Correct answer: A — ```python from pyspark.sql import Window from pyspark.sql import functions as F w = Window.partitionBy("device_id").orderBy("event_time") result_df = (events_df .withColumn("rn", F.row_number().over(w)) .filter(F.col("rn") == 1) .drop("rn")) ```
- A. A window partitioned by device_id and ordered by event_time with row_number() deterministically ranks readings per device, so filtering rn == 1 always keeps the same earliest row regardless of partitioning or re-execution.
- B. Even after orderBy, dropDuplicates does not guarantee which row within a duplicate group is retained, because Spark does not carry the sort order through the deduplication shuffle, so the kept row can vary between runs.
- C. Deduplicating on both device_id and event_time only removes rows that are exact duplicates on that pair; it does not collapse multiple distinct event_time values down to the earliest one per device.
- D. Aggregating with min(event_time) returns only device_id and the minimum timestamp, dropping the payload column and any other fields from the original row.
- E. distinct() compares the full row across all columns, so it will not deduplicate at the device_id grain at all if payload or event_time differ between duplicate readings.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.distinct() accepts column names, so you can use it to deduplicate on specific columns.Why is that wrong?
distinct() takes no arguments and always compares every column. To deduplicate on specific columns, use dropDuplicates() or dropDuplicatesWithinWatermark().
Covered in Exact duplicates with distinct()
2.dropDuplicates(['name', 'height']) returns only the name and height columns.Why is that wrong?
subset only decides which columns are compared. Every column stays in the result, as the reference output shows with age.
Covered in Choosing the comparison columns with dropDuplicates()
3.dropDuplicatesWithinWatermark is a general-purpose dedup method that also works on batch DataFrames.Why is that wrong?
It only works on streaming DataFrames, and only after a watermark has been set with withWatermark. For batch data, use dropDuplicates.
Covered in Deduplicating streams: state and watermarks
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Returns a new DataFrame containing the distinct rows in this DataFrame.”
↩︎ Exact duplicates with distinct() - 2.
“Return a new DataFrame with duplicate rows removed, optionally only considering certain columns.”
↩︎ Choosing the comparison columns with dropDuplicates()“For a static batch DataFrame, it just drops duplicate rows.”
↩︎ Deduplicating streams: state and watermarks“For a streaming DataFrame, it will keep all data across triggers as intermediate state to drop duplicates rows.”
↩︎ Deduplicating streams: state and watermarks“data older than watermark will be dropped to avoid any possibility of duplicates.”
↩︎ Deduplicating streams: state and watermarks“List of columns to use for duplicate comparison (default All columns).”
↩︎ Key concept“List of columns to use for duplicate comparison (default All columns).”
↩︎ Exam trap 2 - 3.
“Without a watermark, state grows indefinitely and can cause memory issues.”
↩︎ Deduplicating streams: state and watermarks“Use dropDuplicatesWithinWatermark to remove duplicates on any field, even when fields differ across duplicate records, such as event time or arrival time.”
↩︎ Deduplicating streams: state and watermarks“Structured Streaming guarantees exactly-once processing but doesn't deduplicate records from data sources.”
↩︎ Deduplicating streams: state and watermarks“To deduplicate specific columns instead of all columns, use dropDuplicates() or dropDuplicatesWithinWatermark() instead of distinct.”
↩︎ Exam trap 1“To guarantee queries drop all duplicates, set the watermark threshold to be greater than the maximum timestamp difference between duplicate events.”
↩︎ Checkpoint - 4.https://docs.databricks.com/aws/en/pyspark/reference/classes/dataframe/dropDuplicatesWithinWatermarkOfficial docs
“This only works with streaming DataFrame, and watermark for the input DataFrame must be set via withWatermark.”
↩︎ Deduplicating streams: state and watermarks“This only works with streaming DataFrame, and watermark for the input DataFrame must be set via withWatermark.”
↩︎ Exam trap 3