CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 5 · Lesson 28/32

    Streaming dropDuplicates: Deduplication With and Without a Watermark

    Perform Streaming Deduplication in Structured Streaming, both with and without watermark usage.

    10 min read
    3.12% of exam
    6 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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?

    Sources12

    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:

    A static DataFrame with an exact duplicate row (Alice, 5, 80)python
    df = spark.createDataFrame([
        Row(name='Alice', age=5, height=80),
        Row(name='Alice', age=5, height=80),
        Row(name='Alice', age=10, height=80)
    ])
    Passing a subset compares only name and height, so the age=10 row also counts as a duplicatepython
    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.

    How dropDuplicates behaves on static and streaming input
    InputWhat dropDuplicates doesState kept
    Static batch DataFrameDrops duplicate rows in one passNone across triggers
    Streaming DataFrame, no watermarkDrops rows matching any key seen in any earlier triggerAll data across triggers, growing without bound
    Streaming DataFrame with withWatermarkDrops duplicates within the lateness limit; data older than the watermark is droppedLimited 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?

    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:

    Watermark applied before distinct on a streampython
    (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?

    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?

    Sources435

    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?

    Sources6

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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.

      Covered in dropDuplicates on a stream without a watermark

    3. 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. 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. 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. 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. 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. 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. 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

    Continue to page 2 of 2

    dropDuplicatesWithinWatermark: Deduplicating Streams by Unique ID

    Spotted a mistake, or was something unclear? Tell us.