What you will be able to do
- Explain why a Structured Streaming query can still emit duplicate records despite its exactly-once guarantee
- Describe how dropDuplicates and distinct behave on a streaming DataFrame compared with a static one
- Bound deduplication state by applying withWatermark before dropDuplicates or distinct
- Recognise which changes to a streaming deduplication break recovery from the checkpoint after a restart
Key concept
Deduplication state — To recognise a duplicate in a stream, Spark has to remember every key it has already seen, and it carries those keys from one micro-batch to the next as state. With no watermark that memory never shrinks. A watermark lets Spark discard keys once they are too old to have duplicates arriving.
1.Exactly-once processing does not mean duplicate-free data
Candidates often get this wrong. Structured Streaming guarantees that each record it reads is processed exactly once. That guarantee says nothing about whether the source sent the same event more than once. Duplicates are normal with many sources. For example, many message queues have at-least-once guarantees, so duplicate records should be expected when reading from one of these message queues.
Removing those duplicates is your job, and the tools for it are distinct, dropDuplicates and dropDuplicatesWithinWatermark. All of them share one property that shapes everything else in this lesson: they are stateful. A stateless query only tracks which rows have been processed from the source to the sink. A stateful query must also keep intermediate state and update it incrementally. The Databricks documentation lists distinct and dropDuplicates next to streaming aggregations and stream-stream joins as stateful operations.
Checkpoint 1 of 5· Check yourself
Which of these does the Databricks documentation list as a stateful operation in Structured Streaming?
dropDuplicates must remember data between micro-batches, so it is listed alongside distinct, streaming aggregation and stream-stream joins as stateful.
“Stateful operations include streaming aggregation, distinct, dropDuplicates, stream-stream joins, and custom stateful applications.”Source: docs.databricks.com
2.dropDuplicates on a stream without a watermark
dropDuplicates(subset=None) returns a new DataFrame with duplicate rows removed. With no argument, it compares all columns. Pass a list of column names to decide duplicates on those columns only. On a static DataFrame, that is all it does. The reference example builds three rows, two of which are identical:
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(['name', 'height']).show()
# +-----+---+------+
# | name|age|height|
# +-----+---+------+
# |Alice| 5| 80|
# +-----+---+------+On a streaming DataFrame you call the same method, but it behaves differently. A static DataFrame is deduplicated in a single pass. A stream never ends, so a duplicate of a row from the first micro-batch could show up a week later. To catch it, Spark keeps all data across triggers as intermediate state. With no watermark, nothing tells Spark when an old key can be forgotten, so the state keeps growing for as long as the query runs. distinct has the same issue, because it tracks every unique record in state.
| Input | What dropDuplicates does | State kept |
|---|---|---|
| Static batch DataFrame | Drops duplicate rows in one pass | None across triggers |
| Streaming DataFrame, no watermark | Drops rows matching any key seen in any earlier trigger | All data across triggers, growing without bound |
| Streaming DataFrame with withWatermark | Drops duplicates within the lateness limit; data older than the watermark is dropped | Limited by the watermark |
Checkpoint 2 of 5· Check yourself
A long-running streaming query calls dropDuplicates(["orderId"]) and has no watermark. What happens to its deduplication state over time?
On a streaming DataFrame, dropDuplicates keeps data across triggers so it can recognise duplicates later. With no watermark, nothing ever removes that state.
“For a streaming DataFrame, it will keep all data across triggers as intermediate state to drop duplicates rows.”Source: docs.databricks.com
Sources3
3.Bounding the state with withWatermark
You bound that state with a watermark. withWatermark(eventTime, delayThreshold) takes the name of the event-time column and an interval string such as "1 minute" or "5 hours". It only works in Structured Streaming. Spark computes the current watermark as the MAX(eventTime) seen across all partitions in the query, minus the delay threshold. Coordinating that value across partitions is costly, so Spark only guarantees that the watermark it actually uses is at least delayThreshold behind the real event time.
With a watermark in place, dropDuplicates changes in two ways. The dropDuplicates reference says: "You can use withWatermark to limit how late the duplicate data can be and the system will accordingly limit the state." It also adds that "data older than watermark will be dropped to avoid any possibility of duplicates." So the bounded state has a cost: a record that arrives after the watermark has passed it is discarded, not checked against the remembered keys.
The same pattern works for distinct. Apply the watermark first, then deduplicate:
(streamingDf
.withWatermark("eventTime", "1 hour")
.distinct()
)In this example, the query removes duplicate records that arrive within 1 hour of the latest observed eventTime. After that threshold passes, it drops the deduplication state. distinct compares every column. To deduplicate on specific columns, the documentation points you to dropDuplicates() or dropDuplicatesWithinWatermark().
Checkpoint 3 of 5· Exam question
A clickstream pipeline ingests events with columns `click_id`, `event_time`, and `page_url`. The business guarantees that a duplicate `click_id` never arrives more than 20 minutes after the original event. An engineer wants to deduplicate on `click_id` while bounding the size of the dedup state store using this known delay. Which code correctly implements this?
Correct answer: A — ```python clicks_df \ .withWatermark("event_time", "20 minutes") \ .dropDuplicates(["click_id", "event_time"]) ```
- A. This is correct: calling `withWatermark` on the event-time column before `dropDuplicates`, and including that same event-time column in the dedup key list, lets the engine bound state using the watermark and expire keys once the watermark passes their event time by more than the declared delay.
- B. Calling `withWatermark` after `dropDuplicates` does not bound the dedup state, because the watermark must be established before the stateful operation so the engine can use it to decide when state is safe to evict. The state store still grows without bound here.
- C. The watermark must be defined on an event-time column, not on the `click_id` identifier column, and `click_id` needs to appear in the dedup key list for identifier-based deduplication to work as intended. This swaps the two columns' roles incorrectly.
- D. Filtering on `current_timestamp()` compares against processing time, not the watermark mechanism, and does not tell the stateful dedup operator when it is safe to drop old keys, so the state store still keeps every `click_id` indefinitely.
- E. Filtering after `dropDuplicates` narrows the output rows but does not give the deduplication operator any watermark information, so the underlying dedup state continues to grow without a bound regardless of the later filter.
Checkpoint 4 of 5· Check yourself
A stream has withWatermark("eventTime", "1 hour") followed by dropDuplicates. The largest eventTime seen so far is 14:00. A record with eventTime 12:30 arrives. What is the documented behaviour?
The watermark is about 13:00, and 12:30 is behind it. Spark may already have removed the state for that period, so it drops the record rather than risk emitting a duplicate.
“In addition, data older than watermark will be dropped to avoid any possibility of duplicates.”Source: docs.databricks.com
4.Deduplication state and query restarts
Deduplication state is not held only in memory. Structured Streaming automatically checkpoints state data to fault-tolerant storage and restores it after a restart. That restore assumes the state schema is the same before and after the restart. Streaming deduplication is one of the stateful operations whose schema must not change between restarts. The checkpoint documentation uses sdf.dropDuplicates("a") as its example and states that any change in the number or type of grouping keys is not allowed.
In practice, the dedup columns are fixed when the query first starts against its checkpoint. If you later add a column to the subset, Spark cannot recover the state it already has.
On Databricks there is one more constraint. dropDuplicates() and dropDuplicatesWithinWatermark() can fail to restart because of a state schema compatibility check when the compute access mode changes. Changing between dedicated and no-isolation is allowed, and so is changing between standard and serverless. Other combinations should be avoided.
Checkpoint 5 of 5· Check yourself
A running query deduplicates with dropDuplicates(["a"]). An engineer changes it to dropDuplicates(["a", "b"]) and restarts it on the same checkpoint. According to the checkpoint documentation, what is true?
Deduplication state is stored with a fixed schema. Changing the number or type of keys stops Spark from restoring that state from the checkpoint.
“Streaming deduplication: For example, sdf.dropDuplicates("a"). Any change in number or type of grouping keys or aggregates is not allowed.”Source: docs.databricks.com
Both are stateful operators, and their state goes through a schema compatibility check on restart. A stateless query has no operator state to check.
Sources6
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Because Structured Streaming is exactly-once, a duplicate event sent by the source will not appear in the output.Why is that wrong?
Exactly-once describes how the engine processes the records it reads. If the source delivers the same event twice, both copies are processed unless you deduplicate them yourself.
Covered in Exactly-once processing does not mean duplicate-free data
2.dropDuplicates on a stream without a watermark is safe, because Spark cleans up its deduplication state on its own.Why is that wrong?
Without a watermark, Spark keeps all data across triggers as state, and that state grows without limit. Only a watermark lets Spark remove old state.
3.You can change the dedup columns of a running streaming query and restart it on the same checkpoint.Why is that wrong?
Deduplication is a stateful operation with a fixed state schema. Changing the number or type of its keys is not allowed between restarts.
Covered in Deduplication state and query restarts
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“because many message queues have at-least once guarantees, duplicate records should be expected when reading from one of these message queues”
↩︎ Exactly-once processing does not mean duplicate-free data - 2.
“Stateful operations include streaming aggregation, distinct, dropDuplicates, stream-stream joins, and custom stateful applications.”
↩︎ Exactly-once processing does not mean duplicate-free data - 3.
“For a static batch DataFrame, it just drops duplicate rows.”
↩︎ dropDuplicates on a stream without a watermark“You can use withWatermark to limit how late the duplicate data can be and the system will accordingly limit the state.”
↩︎ Bounding the state with withWatermark“For a streaming DataFrame, it will keep all data across triggers as intermediate state to drop duplicates rows.”
↩︎ Checkpoint“In addition, data older than watermark will be dropped to avoid any possibility of duplicates.”
↩︎ Checkpoint - 4.
“The current watermark is computed by looking at the MAX(eventTime) seen across all of the partitions in the query minus a user specified delayThreshold.”
↩︎ Bounding the state with withWatermark - 5.
“the streaming query removes duplicate records that arrive within 1 hour of the latest observed eventTime.”
↩︎ Bounding the state with withWatermark“Without a watermark, state grows indefinitely and can cause memory issues.”
↩︎ Key concept“Structured Streaming guarantees exactly-once processing but doesn't deduplicate records from data sources.”
↩︎ Exam trap 1“The distinct operation tracks every unique record in state. Without a watermark, state grows indefinitely and can cause memory issues.”
↩︎ Exam trap 2“Structured Streaming guarantees exactly-once processing but doesn't deduplicate records from data sources.”
↩︎ Prediction - 6.
“The stateful operators dropDuplicates() and dropDuplicatesWithinWatermark() can fail to restart due to a state schema compatibility check when changing between compute access modes.”
↩︎ Deduplication state and query restarts“Any change in number or type of grouping keys or aggregates is not allowed.”
↩︎ Exam trap 3“Streaming deduplication: For example, sdf.dropDuplicates("a"). Any change in number or type of grouping keys or aggregates is not allowed.”
↩︎ Checkpoint