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
Start by naming where your duplicates come from, because each source has a different key and a different delay.
| Origin | Typical delay after the original | What identifies it |
|---|---|---|
| Producer retry after a lost acknowledgement | Milliseconds to seconds | Same event ID; a new broker offset |
| Connector or CDC restart | Seconds to minutes | Same source position (for example a log sequence number) |
| Spark reprocessing a micro-batch after failure | Seconds | Same offsets; handled by checkpoint plus idempotent sink |
| Client resend of a cached batch | Minutes to days | Same event ID if the client assigns one |
| Operator backfill or replay | Days to months | Same 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.
- 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.
- 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.
- 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
| Choice | Gain | Cost |
|---|---|---|
| Longer tier 1 horizon | Fewer duplicates reach aggregates | State memory and checkpoint size grow linearly |
| Longer tier 2 lookback | Catches the slow tail | More data read per micro-batch; higher latency |
| Content-hash key | Works without producer IDs | Can merge legitimate repeats |
| Custom stateful deduper | Per-key horizons, counts, latest-wins | More 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
- List every duplicate origin in your pipeline and write down which key identifies each one.
- Confirm producers keep the same event ID across retries by reading the retry code, then test it.
- Keep an undeduplicated landing table and run the delay-measurement job to get p99.9 and maximum delay.
- Set tier 1's horizon from p99.9, tier 2's lookback from the maximum, and schedule a tier 3 audit.
- Move deduplication before any aggregation or join, and alert on state size and on the audit's duplicate count.
- Reach for flatMapGroupsWithState or transformWithState only when you need per-key horizons or latest-wins semantics.