What you will be able to do
- Apply dropDuplicatesWithinWatermark to a streaming DataFrame with the required watermark
- Explain how it differs from distinct and dropDuplicates when duplicate records have different event times
- Choose a watermark threshold that guarantees duplicates are removed
- Pick deduplication columns from the source's uniqueness contract
1.The dropDuplicatesWithinWatermark method and its precondition
dropDuplicatesWithinWatermark deduplicates a stream on chosen columns and keeps the state bounded. Databricks documents it for Databricks Runtime 13.3 LTS and above. Its signature looks like dropDuplicates: an optional list of columns to compare, which defaults to all columns.
dropDuplicatesWithinWatermark(subset: Optional[List[str]] = None)The method has two hard requirements. The DataFrame must be streaming, and withWatermark must have been called on it first. The watermark sets the threshold that the method's guarantee depends on, which a later section covers. The standard example removes duplicates by a guid column and uses eventTime as the watermark column:
Checkpoint 1 of 5· Fill the gap
This stream must remove duplicates by guid, including duplicates whose eventTime differs, and the method used must be one that requires a watermark. Which call fills the blank?
(streamingDf
.withWatermark("eventTime", "10 hours")
. ? (["guid"])
)dropDuplicatesWithinWatermark is the method that requires a watermark and removes duplicates by guid even when other fields, such as eventTime, are different.
Source: docs.databricks.comSources1
2.How it compares with distinct and dropDuplicates
The main difference shows up when two copies of an event are not identical. A retried delivery usually has the same business ID but can have a different arrival time or event time. The Databricks watermark guide says dropDuplicatesWithinWatermark removes duplicates on any field, "even when fields differ across duplicate records, such as event time or arrival time." So you name the identifying column, such as guid. The timestamps can differ between copies, and the records are still treated as duplicates.
Its handling of state also differs. For the plain operators, a watermark is optional, and without one the state grows without limit. dropDuplicatesWithinWatermark is watermark-aware by design, and Databricks describes it as not requiring unbounded state to detect duplicates.
| Operator | Columns compared | Watermark | State without a watermark |
|---|---|---|---|
| distinct() | All columns | Optional; apply withWatermark first to limit state | Grows indefinitely |
| dropDuplicates(subset) | subset, or all columns by default | Optional; limits how late duplicates can be | All data kept across triggers |
| dropDuplicatesWithinWatermark(subset) | subset, or all columns by default; other fields may differ | Required, set via withWatermark | Not applicable: the method fails without one |
Checkpoint 2 of 5· Exam question
Which two statements accurately describe how `dropDuplicates` behaves when used on a streaming DataFrame with a watermark defined on the event-time column included in the dedup key list? (Choose 2 answers)(Select 2)
Correct answers: A, B — Once the watermark advances past a record's event time by more than the declared delay, the state kept for that record becomes eligible for removal from the state store.; The watermark delay should be chosen based on how late duplicate events are expected to arrive relative to the original event, since a shorter delay risks missing valid duplicates.
- A. This is correct: the engine uses watermark progress to identify state that is no longer expected to see new duplicates and removes it, which is exactly how watermark-based dedup bounds the state store size over time.
- B. This is correct: the watermark delay directly controls the tradeoff between state size and correctness, so it should be set long enough to cover realistic duplicate arrival lateness or genuine duplicates arriving after the watermark passes will be treated as new records.
- C. The watermark column must be an event-time column of a timestamp type, not an arbitrary string, and Spark tracks lateness using actual timestamp values rather than string or ingestion-order comparisons.
- D. Attaching a watermark does not force downstream aggregations into complete mode; output mode is chosen independently by the query, and watermarks are commonly paired with append or update mode instead.
- E. The dedup operator relies on the declared watermark based on event time in the data, not on wall-clock processing time, which is precisely what allows it to correctly handle data that arrives out of order or with processing delays.
3.Choosing the watermark threshold
The guarantee is stated in terms of time distance. The method reference describes the semantics as: "Events are deduplicated as long as the time distance of earliest and latest events are smaller than the delay threshold of watermark." If two copies of an event are closer together than the threshold, the second is always dropped. If they are further apart, the watermark guide says queries "might also deduplicate records that arrive outside the threshold, but this isn't guaranteed."
This gives a sizing rule. Work out the largest gap you expect between duplicates of the same event, and set the threshold larger than that. A larger threshold keeps more state, so you are trading memory for certainty. Keep the late-data rule in mind as well: data that arrives older than the watermark is dropped.
Checkpoint 3 of 5· Check yourself
Your source can redeliver an event up to 3 hours after the original delivery. Which watermark threshold for dropDuplicatesWithinWatermark guarantees that the redelivered copy is removed?
Duplicates are only guaranteed to be removed when the threshold is larger than the maximum timestamp difference between them. Anything removed outside the threshold is not guaranteed.
“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
The order in which data arrives also matters. Out-of-order data can push the watermark forward too early, and older records that arrive afterwards are treated as late and dropped. When reading a Delta table on Databricks, the withEventTimeOrder option processes the initial snapshot in order of the watermark timestamp. That option is supported only with Python:
clicksDedupDf = (
spark.readStream
.option("withEventTimeOrder", "true")
.table("rawClicks")
.withWatermark("clickTimestamp", "5 seconds")
.dropDuplicatesWithinWatermark(["userId", "clickAdId"]))4.Choosing the columns that identify an event
The threshold controls how long Spark remembers a key. The subset controls what counts as the same event, and an error there drops real data without any warning. Databricks advises deduplicating on the columns that uniquely identify an event. Those columns might be a single ID or several columns together. In one example, a click sequence number is only unique within its session, so session_id and click_seq_num together identify a click.
Choose those columns from the source's uniqueness contract, meaning what the producer guarantees to be unique, and not from what happens to look unique in sample data. Columns that can legitimately repeat will discard real events. A user who clicks the same ad twice is a common example: if you deduplicate on the user and the ad, the second click is silently dropped.
Checkpoint 4 of 5· Match them up
Match each choice to its effect on streaming deduplication
Tap a term, then the definition that fits it.
The threshold decides how long duplicates are remembered. The column subset decides which records count as the same event. Getting either one wrong causes a different kind of data error.
“Columns that can legitimately repeat discard real events when you treat them as the identity.”Source: docs.databricks.com
Checkpoint 5 of 5· Exam question
A payments pipeline receives events with columns `payment_id`, `submitted_at`, and `amount`. A known issue with an upstream non-idempotent writer means the same logical payment can be re-emitted with a slightly different `submitted_at` timestamp on retry, so `submitted_at` cannot safely be part of the dedup key. The team still wants watermark-based state cleanup using `submitted_at`. Which approach fits this requirement?
Correct answer: A — ```python payments_df \ .withWatermark("submitted_at", "10 hours") \ .dropDuplicatesWithinWatermark(["payment_id"]) ```
- A. This is correct: `dropDuplicatesWithinWatermark` is designed exactly for this case, where the event-time column cannot be part of the identifier because retries change it. It deduplicates purely on `payment_id` while still using the watermark on `submitted_at` to bound and clean up state.
- B. Including `submitted_at` in the key list for `dropDuplicates` means two retries of the same payment with different timestamps are treated as distinct records and both pass through, which defeats the stated goal since `submitted_at` cannot reliably identify duplicates here.
- C. Establishing the watermark after `dropDuplicates` does not give the dedup operator any lateness information to use for state cleanup, so this does not bound state regardless of which columns are chosen as the key.
- D. Grouping and aggregating by `payment_id` computes an aggregate value rather than filtering out duplicate rows, and it also places the watermark call after a stateful operation, so it does not perform row-level deduplication as required.
- E. Keying `dropDuplicates` on `submitted_at` alone deduplicates by timestamp rather than by payment identity, so two different payments that happen to share a timestamp would be incorrectly collapsed into one record.
Sources2
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.dropDuplicatesWithinWatermark can be used without withWatermark, or on a static DataFrame, like dropDuplicates.Why is that wrong?
It only works on streaming DataFrames, and a watermark must be set on the input with withWatermark.
Covered in The dropDuplicatesWithinWatermark method and its precondition
2.With dropDuplicatesWithinWatermark, duplicates that arrive further apart than the watermark threshold are still guaranteed to be removed.Why is that wrong?
Removal is only guaranteed inside the threshold. Outside it, a duplicate might be removed, but nothing guarantees it.
Covered in Choosing the watermark threshold
3.Any set of columns whose values look distinct in sample data is a safe deduplication key.Why is that wrong?
The key has to come from what the source guarantees to be unique. Columns that can legitimately repeat will drop real events.
Covered in Choosing the columns that identify an event
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“In Databricks Runtime 13.3 LTS or above, you can use a unique identifier to deduplicate records within a watermark threshold.”
↩︎ The dropDuplicatesWithinWatermark method and its precondition“You must specify a watermark to use the dropDuplicatesWithinWatermark method”
↩︎ The dropDuplicatesWithinWatermark method and its precondition“Use dropDuplicatesWithinWatermark to remove duplicates on any field, even when fields differ across duplicate records, such as event time or arrival time.”
↩︎ How it compares with distinct and dropDuplicates“Queries might also deduplicate records that arrive outside the threshold, but this isn't guaranteed.”
↩︎ Exam trap 2“To guarantee queries drop all duplicates, set the watermark threshold to be greater than the maximum timestamp difference between duplicate events.”
↩︎ Checkpoint - 2.
“which is watermark-aware and doesn't require unbounded state to detect duplicates”
↩︎ How it compares with distinct and dropDuplicates“A user clicking the same ad twice is a common example: deduplicating on the user and the ad silently drops the second click.”
↩︎ Choosing the columns that identify an event“Choose those columns from the source's uniqueness contract, not from what looks distinct in sample data.”
↩︎ Exam trap 3“Columns that can legitimately repeat discard real events when you treat them as the identity.”
↩︎ Checkpoint - 3.https://docs.databricks.com/aws/en/pyspark/reference/classes/dataframe/dropDuplicatesWithinWatermarkOfficial docs
“Events are deduplicated as long as the time distance of earliest and latest events are smaller than the delay threshold of watermark.”
↩︎ Choosing the watermark threshold“Note: too late data older than watermark will be dropped.”
↩︎ Choosing the watermark threshold“This only works with streaming DataFrame, and watermark for the input DataFrame must be set via withWatermark.”
↩︎ Exam trap 1“This only works with streaming DataFrame, and watermark for the input DataFrame must be set via withWatermark.”
↩︎ Prediction - 4.
“out-of-order data causes the watermark value to jump ahead incorrectly.”
↩︎ Choosing the watermark threshold