Duplicates in a stream are not an edge case; they are the price of at-least-once delivery, which is what nearly every producer, broker and connector gives you. A payment event is retried after a timeout that actually succeeded, a CDC connector restarts and replays its last few seconds, a mobile client resends a batch, an operator backfills a day of history. Each produces a row that is the same event as one you already processed.

Spark Structured Streaming has an operator for this, and Spark streaming deduplication operators covers it in detail: dropDuplicates with and without a watermark, dropDuplicatesWithinWatermark, state sizing, and an idempotent Kafka-to-Delta pipeline. This page is about the design decisions around that operator, which matter more than the call itself: what the key should be, how long to remember it, how to split that memory between the stream and the sink, what to do when the built-in operator is not enough, and where deduplication goes in a pipeline.

Where duplicates come from

Where duplicates enter a streaming pipeline, and which tier removes themProducerretries, resendsBroker / CDCat-least-onceSpark replaybatch reprocessedBackfills, clientsre-ingested historyTier 1: in-stream statecovers p99.9 of delay, hoursTier 2: sink-side checkMERGE on key, daysTier 3: batch auditcount dups, weeksmostidempotent writelate tailHorizon splitstate cost grows with horizon; sink cost grows with lookbackMeasure the duplicate-delay distribution first; it decides every tier's horizon.
Duplicates enter at four points. Most arrive within seconds or minutes and are cheapest to remove in streaming state; the long tail is removed at the sink or by a periodic audit.

Start by naming where your duplicates come from, because each source has a different key and a different delay.

OriginTypical delay after the originalWhat identifies it
Producer retry after a lost acknowledgementMilliseconds to secondsSame event ID; a new broker offset
Connector or CDC restartSeconds to minutesSame source position (for example a log sequence number)
Spark reprocessing a micro-batch after failureSecondsSame offsets; handled by checkpoint plus idempotent sink
Client resend of a cached batchMinutes to daysSame event ID if the client assigns one
Operator backfill or replayDays to monthsSame event ID; large volume at once

The third row is different in kind. When a Spark query fails mid-batch and restarts, it re-reads the same offsets from its checkpoint and recomputes the same batch. That is not a data duplicate; it is a write duplicate, and the fix is an idempotent sink keyed on the batch, not a dedup operator. The other rows are genuine duplicate records with different offsets, and only a key-based check removes them.

Choosing the dedup key

The key decides what counts as a duplicate, and each choice has a failure mode.

  • Producer-assigned event ID. The best key: a UUID or ULID generated once, when the event is created, and carried through every retry. It catches producer retries, client resends and backfills. It fails only when producers regenerate the ID on retry, which is a common bug; check the retry path, not just the happy path.
  • Source position. Topic, partition and offset, or a CDC log sequence number. It catches connector replays exactly but misses producer retries, because a retried send gets a new offset. Use it for CDC streams, where the database position is the identity.
  • Content hash. A hash of the business fields. It works when nothing assigns IDs, but it merges legitimate repeats: two identical 5-dollar coffee purchases a minute apart become one. Include a field that distinguishes real repeats, such as the client's own sequence number, and exclude fields that change on retry, such as an ingestion timestamp or a retry counter.
  • Natural business key plus version. Order ID with an order version. Useful when the stream carries updates and you want to drop repeated versions while keeping new ones.

Building a content key safely

Hash content deterministically. Serialise the chosen fields in a fixed order with fixed formats before hashing; hashing a JSON string directly breaks when two producers order keys differently or one writes 1.0 where the other writes 1.

from pyspark.sql import functions as F

KEY_FIELDS = ["account_id", "merchant_id", "amount_minor", "currency", "client_seq"]

def with_content_key(df):
    parts = [F.coalesce(F.col(f).cast("string"), F.lit("\u0000")) for f in KEY_FIELDS]
    return df.withColumn("dedup_key", F.sha2(F.concat_ws("\u001f", *parts), 256))

The null sentinel and the unit-separator delimiter prevent two different rows from producing the same string, for example ('ab', 'c') and ('a', 'bc').

Measuring the duplicate horizon

The second decision is the horizon: how long after the original a duplicate can still be removed. Do not guess it. Measure it from history with a batch job that finds, for every repeated key, the delay between first and later sightings.

from pyspark.sql import Window, functions as F

events = spark.read.format("delta").load("/lake/raw_events")      # undeduplicated landing table
w = Window.partitionBy("event_id").orderBy("ingest_ts")
delays = (events
    .withColumn("first_ts", F.first("ingest_ts").over(w))
    .withColumn("rn", F.row_number().over(w))
    .where("rn > 1")
    .select((F.col("ingest_ts").cast("long") - F.col("first_ts").cast("long")).alias("delay_s")))

delays.select(
    F.count("*").alias("duplicates"),
    F.expr("percentile_approx(delay_s, array(0.5, 0.99, 0.999, 0.9999))").alias("p50_p99_p999_p9999"),
    F.max("delay_s").alias("max_s"),
).show(truncate=False)

Run it on raw, undeduplicated data, which means you need a landing table that keeps everything. Measure by ingestion time, because that is what drives how long the stream must remember a key; event time measures something else. Typical results are lopsided: most duplicates arrive within seconds, a thin tail within hours, and a few, from backfills and client resends, days later.

Three tiers instead of one window

No single mechanism covers that whole distribution cheaply. In-stream state costs memory and checkpoint size for every key for the whole horizon. A sink-side check costs a join on every batch, growing with the lookback. So split the horizon into tiers.

  1. Tier 1, in-stream. Use a watermarked dedup operator with a horizon covering about the 99.9th percentile of delay, often one to six hours. This removes nearly all duplicates before they reach downstream aggregations, where they would inflate counts.
  2. Tier 2, at the sink. Use an insert-only MERGE on the key against a lookback of days, as shown in the operator article. It catches the tail tier 1 missed and makes replays harmless.
  3. Tier 3, audit. Run a daily batch job that counts keys appearing more than once over weeks. It does not fix anything in real time; it tells you whether the first two tiers are sized correctly, and repairs the table when they are not.

Worked example: sizing a clickstream

A worked example. A clickstream of 50,000 events per second has a producer-assigned 16-byte event ID. The measurement job shows 0.4% duplicates, a 99.9th-percentile delay of 40 minutes and a maximum of 3 days.

Tier 1 with a 1-hour horizon holds about 50,000 x 3,600 = 180 million keys. In a RocksDB state store, budget on the order of 100 bytes per key once key bytes, the stored timestamp, row overhead and RocksDB amplification are counted. That is about 18 GB of state, spread across shuffle partitions. With 200 partitions that is about 90 MB each, which is comfortable. A 3-day tier 1 would need 72 times as much, about 1.3 TB, which is why the tail goes to the sink.

Tier 2 with a 3-day lookback joins each micro-batch against three days of the target table. With the table clustered or partitioned by event date and the MERGE condition bounded by date, each batch reads three partitions. Tier 3 runs nightly over 30 days and reports whether any duplicates escaped. A nonzero count means a delay outside both horizons, or a producer that changes IDs.

The per-key byte figure is an assumption for planning. Measure the real value from the state operator's memory metrics once the query runs, and resize from that.

When the built-in operator is not enough

The built-in operators keep the first occurrence and drop the rest, using a watermark-based horizon. Sometimes you need more: a horizon per key (for example 7 days for payments and 1 hour for clicks in the same stream), keeping the latest version rather than the first, or emitting a count of suppressed duplicates. Then write the state yourself. In Scala, flatMapGroupsWithState with an event-time timeout does it.

import org.apache.spark.sql.streaming.{GroupState, GroupStateTimeout, OutputMode}

case class Event(eventId: String, kind: String, ts: java.sql.Timestamp, payload: String)
case class Seen(firstTsMs: Long, dupCount: Long)

def horizonMs(kind: String): Long =
  if (kind == "payment") 7L * 24 * 3600 * 1000 else 3600L * 1000

def dedup(id: String, rows: Iterator[Event], state: GroupState[Seen]): Iterator[Event] = {
  if (state.hasTimedOut) { state.remove(); Iterator.empty }        // horizon passed: forget the key
  else {
    val batch = rows.toSeq.sortBy(_.ts.getTime)
    if (state.exists) {                                               // already emitted earlier
      state.update(state.get.copy(dupCount = state.get.dupCount + batch.size))
      Iterator.empty
    } else {
      val first = batch.head
      state.update(Seen(first.ts.getTime, batch.size - 1))
      state.setTimeoutTimestamp(first.ts.getTime + horizonMs(first.kind))
      Iterator(first)
    }
  }
}

val deduped = events
  .withWatermark("ts", "10 minutes")
  .groupByKey(_.eventId)
  .flatMapGroupsWithState(OutputMode.Append, GroupStateTimeout.EventTimeTimeout)(dedup)

The timeout fires when the watermark passes the timestamp you set, so the watermark delay still matters: it bounds how late a row can be and still be compared. Two points are easy to miss. Event-time timeouts require a watermark on the input. And because the timeout is event-time based, a key whose events stop arriving is only cleared when the watermark advances, so an idle source keeps state alive.

Spark 4.0 adds transformWithState, available as transformWithStateInPandas in Python, as the successor to this API. Its state variables can carry a processing-time TTL, so expiry no longer has to be coded by hand. If you are on 4.0, prefer it for new code; the design reasoning above is unchanged.

Where dedup goes in the pipeline

Deduplicate early, before anything that counts or joins. A duplicate that reaches a windowed aggregation inflates a count, and the dedup operator downstream cannot see it any more because the row has been aggregated away. In a stream-stream join, a duplicate on one side produces duplicate join output, so deduplicate each input before the join.

Keep the key cheap to shuffle. The dedup operator repartitions by key, so put the key and the event-time column first, and drop large payload columns you do not need before the shuffle if you can join them back later.

Treat the key and horizon as part of the checkpoint's contract. Changing the dedup columns changes the state schema, and the query cannot restart from the old checkpoint. Plan such changes as a migration: start a new query with a new checkpoint from an offset slightly before the old query stopped, and let tier 2 absorb the overlap.

Failure modes

  • Producers regenerate IDs on retry. Every tier is blind; tier 3's audit shows near-zero duplicates while downstream counts are high. Fix the producer, or fall back to a content key.
  • Content key merges real events. Revenue drops slightly after enabling deduplication. Add a distinguishing field and re-measure.
  • State grows without bound. A dropDuplicates call without a watermark keeps every key forever. Always pair it with a watermark, and alert on the state operator's numRowsTotal trend.
  • Backfill floods tier 2. A replay of a month falls outside the MERGE lookback and inserts a month of duplicates. Route backfills through a batch job with a full-history anti-join instead of through the stream.
  • A quiet input stalls the watermark. The watermark only advances when new rows arrive, so a stream that goes quiet keeps its state. In a query with several inputs, the default policy takes the minimum of their watermarks, so one quiet input holds back eviction for all; change spark.sql.streaming.multipleWatermarkPolicy to max only after understanding that it drops more late data.

Trade-offs

ChoiceGainCost
Longer tier 1 horizonFewer duplicates reach aggregatesState memory and checkpoint size grow linearly
Longer tier 2 lookbackCatches the slow tailMore data read per micro-batch; higher latency
Content-hash keyWorks without producer IDsCan merge legitimate repeats
Custom stateful deduperPer-key horizons, counts, latest-winsMore code; you own correctness and schema evolution

For the watermark mechanics behind tier 1, read Spark streaming watermarks, and for the Kafka offsets that make replays deterministic, the Kafka source. Sink choices are compared in streaming sinks.

What to do next

  1. List every duplicate origin in your pipeline and write down which key identifies each one.
  2. Confirm producers keep the same event ID across retries by reading the retry code, then test it.
  3. Keep an undeduplicated landing table and run the delay-measurement job to get p99.9 and maximum delay.
  4. Set tier 1's horizon from p99.9, tier 2's lookback from the maximum, and schedule a tier 3 audit.
  5. Move deduplication before any aggregation or join, and alert on state size and on the audit's duplicate count.
  6. Reach for flatMapGroupsWithState or transformWithState only when you need per-key horizons or latest-wins semantics.
Key takeaway: Streaming deduplication is a design problem before it is an operator call. Choose a key that survives retries, measure how late duplicates actually arrive, cover the bulk in watermarked state and the tail with an idempotent sink check and a periodic audit, deduplicate before aggregations and joins, and write custom state only when you need per-key horizons.