CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 3 · Lesson 13/32

    Deduplicating Spark DataFrames: distinct, dropDuplicates and Watermarks

    Perform data deduplication and validation operations on DataFrames.

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

    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.

    distinct() collapses the two identical (23, "Alice") rows into onepython
    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.

    Checkpoint 1 of 5· Check yourself

    df contains the rows (14, 'Tom'), (23, 'Alice') and (23, 'Alice'). What does df.distinct().count() return?

    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.

    With no subset, only the fully identical rows are collapsed. The age=10 row is kept.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)
    ])
    
    df.dropDuplicates().show()
    # +-----+---+------+
    # | name|age|height|
    # +-----+---+------+
    # |Alice|  5|    80|
    # |Alice| 10|    80|
    # +-----+---+------+
    Comparing only name and height: all three rows count as duplicatespython
    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()

    Checkpoint 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?

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

    Deduplicate a stream on guid, with state bounded by a 10-hour watermark on eventTimepython
    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."

    The three deduplication methods compared
    MethodColumns comparedWhere it worksWatermark
    distinct()All columns onlyBatch and streamingOptional in streaming, but without it state grows indefinitely
    dropDuplicates(subset)All columns, or the subset you nameBatch and streamingOptional in streaming; limits state and drops late data
    dropDuplicatesWithinWatermark(subset)All columns, or the subset you nameStreaming onlyRequired, 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?

    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?

    Sources234

    Exam traps

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

    1. 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. 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. 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. 1.
      “Returns a new DataFrame containing the distinct rows in this DataFrame.”
      ↩︎ Exact duplicates with distinct()
    2. 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. 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. 4.
      “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

    Continue to page 2 of 2

    Validating Spark DataFrames: dropna, fillna, replace and Quarantine Filters

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