Almost every streaming pipeline sees duplicates. Producers retry when an acknowledgement is lost, upstream services resend after a timeout, mobile clients replay a queue when they reconnect, and your own query re-processes a micro-batch after a crash. If the output feeds billing, fraud scores or counts, each duplicate is a wrong answer. Spark Structured Streaming gives you a stateful operator that drops repeated keys, but it only works well if you understand what it remembers, for how long, and what it cannot see.
This article explains deduplication from first principles: the two different kinds of duplicate, how the operator stores keys in a state store, the three APIs and their exact semantics, how to size state with real numbers, and how to make the sink idempotent so a replayed batch does not undo the work. Watermark mechanics in general are covered in structured streaming watermarks, and state stores and checkpoints in Structured Streaming state.
Two kinds of duplicate
It helps to separate duplicates by where they are created, because each kind needs a different tool.
Source duplicates are already in the input. The same logical event appears twice in a Kafka topic, usually with the same business identifier such as an order id or an idempotency key, because a producer retried a send that had in fact succeeded. Spark reads both copies as distinct records with different offsets. Only a stateful operator that remembers which keys it has seen can remove them.
Replay duplicates are created by Spark itself. A micro-batch is planned with a fixed offset range, written to the sink, and only then committed in the checkpoint. If the driver dies after the sink write but before the commit, the restarted query runs the same batch again with the same batch id. The dedup operator does not help here, because its state is restored to the version before the batch. The sink must recognise the repeat and ignore it.
A correct pipeline handles both: a dedup operator for source duplicates and an idempotent sink for replays. The general pattern of idempotency keys is described in idempotency in system design.
How the streaming dedup operator works
When you call a dedup API on a streaming DataFrame, the planner inserts a stateful operator. Before it, Spark shuffles rows by a hash of the dedup key, so every copy of a key lands in the same state partition. The number of state partitions is the value of spark.sql.shuffle.partitions when the query first started, and it is fixed in the checkpoint from then on.
Inside each task the logic is simple. For each input row, look up the key in the state store. If it is absent, emit the row and put the key into state. If it is present, drop the row. At the end of the batch the task removes state entries that the watermark says can no longer receive duplicates, and the state store writes a new version that the checkpoint references. The output is the first copy Spark sees in processing order, which is not necessarily the copy with the earliest event time.
The three APIs
| API | What the key is | When state is removed | Good for |
|---|---|---|---|
dropDuplicates(['id']) with no watermark | id | Never | Small, bounded key spaces only |
withWatermark(...).dropDuplicates(['id', 'ts']) | id plus event time | When the stored event time falls below the watermark | Duplicates that carry an identical event time |
withWatermark(...).dropDuplicatesWithinWatermark(['id']) | id only | When the watermark reaches the key's event time plus the delay, about twice the delay after the newest data | Duplicates whose event times differ slightly |
They differ in what counts as the same event and in how long a key is remembered.
dropDuplicates without a watermark
Without a watermark there is no bound on how late a duplicate may arrive, so the Spark documentation states that the query stores data from all past records as state. Every key ever seen stays in the store. State grows linearly with the number of distinct keys, and so do checkpoint size, recovery time and memory pressure.
This is acceptable only when the key space itself is bounded: deduplicating a stream of device registrations where there are a few hundred thousand devices, for example. For an event stream with a fresh id per event, it is a slow memory leak that ends in executor failures weeks after launch.
dropDuplicates with a watermark
The documented pattern is to define a watermark on the event-time column and include that same column in the dedup key. The watermark is the maximum event time Spark has seen minus the delay you specify, computed at the end of each batch and applied in the next. State entries whose event time is below the watermark are evicted, because by definition no new rows that old will be accepted. Input rows older than the watermark are treated as too late and dropped, and the drop is counted in the query progress.
The catch is the key. Because the event time is part of it, two copies are duplicates only if their event times are identical. That holds when the producer stamps the time once and every retry resends the same payload. It fails when each retry stamps a new time, for instance when the timestamp is assigned at write time by a non-idempotent writer. Then ('order-42', 10:00:01.120) and ('order-42', 10:00:01.480) are different keys and both are emitted.
deduped = (events
.withWatermark("event_ts", "30 minutes")
.dropDuplicates(["event_id", "event_ts"]))
dropDuplicatesWithinWatermark
Spark 3.5.0 added dropDuplicatesWithinWatermark for exactly that case. The key is just the identifier, and the event time is used only to decide how long to remember it. The API documentation states the guarantee precisely: events are deduplicated as long as the time distance between the earliest and latest copies is smaller than the delay threshold of the watermark. It requires a watermark, it works only on streaming DataFrames, and data older than the watermark is still dropped as too late.
The Spark guide advises setting the delay longer than the maximum timestamp difference among duplicated events. If producer retries can be spread over ten minutes, a two-minute delay leaves a gap through which duplicates pass. Choosing the delay is therefore a question about your producers, not about Spark: measure the largest gap between copies of the same id in historical data and add a margin.
deduped = (events
.withWatermark("event_ts", "2 hours")
.dropDuplicatesWithinWatermark(["event_id"]))
Worked example: sizing state for a payments stream
Suppose a payments topic carries 50,000 events per second, each with a unique payment_id assigned by the client. Analysis of a week of raw data shows that 0.4% of ids appear more than once, and that 99.99% of duplicate pairs are within 40 minutes of each other, with a long tail from offline mobile clients up to several hours. You choose dropDuplicatesWithinWatermark with a two-hour delay.
Retention is longer than the delay suggests. The operator stores each key with an expiry of its event time plus the delay and evicts it once the watermark reaches that expiry; since the watermark itself trails the newest event time by the delay, a key lives for about twice the delay, here four hours of event time. The store holds 50,000 x 14,400 = 720 million keys. If each entry costs on the order of 100 bytes in RocksDB once you count the key, the expiry timestamp and storage overhead, that is about 72 GB of state. Spread over 200 state partitions it is around 360 MB per partition, which is comfortable for RocksDB on local disk and uncomfortable for the default in-memory store, which keeps state on the JVM heap. Treat the 100-byte figure as an assumption to replace with a measurement from a staging run, by reading the state metrics after a few hours.
The few duplicates older than two hours will pass through. Rather than stretch the watermark to cover them, which would multiply state, you catch them at the sink with a MERGE against the target table, described below. This split is the main design move: a short, cheap in-stream window for the common case and a durable table lookup for the tail.
The full pipeline in code
The example below reads from Kafka, parses JSON, deduplicates within a watermark, and writes each micro-batch to Delta with the idempotent-write options. The state store is switched to RocksDB, which keeps state off the JVM heap.
from pyspark.sql import SparkSession, functions as F, types as T
spark = (SparkSession.builder
.config("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
.config("spark.sql.shuffle.partitions", "200") # fixed by the checkpoint after first run
.getOrCreate())
schema = T.StructType([
T.StructField("payment_id", T.StringType()),
T.StructField("account_id", T.StringType()),
T.StructField("amount_minor", T.LongType()),
T.StructField("event_ts", T.TimestampType()),
])
raw = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "payments")
.option("startingOffsets", "latest")
.load())
events = (raw
.select(F.from_json(F.col("value").cast("string"), schema).alias("e"))
.select("e.*")
.where(F.col("payment_id").isNotNull()))
deduped = (events
.withWatermark("event_ts", "2 hours")
.dropDuplicatesWithinWatermark(["payment_id"]))
APP_ID = "payments-dedup-v1" # stable across restarts; change only with a new checkpoint
def write_batch(batch_df, batch_id):
(batch_df.write.format("delta")
.mode("append")
.option("txnAppId", APP_ID)
.option("txnVersion", batch_id) # a replayed batch_id is skipped by Delta
.save("/lake/payments_clean"))
query = (deduped.writeStream
.foreachBatch(write_batch)
.option("checkpointLocation", "/chk/payments_clean")
.trigger(processingTime="30 seconds")
.start())Two details matter. First, txnAppId and txnVersion are Delta write options for idempotent writes inside foreachBatch: Delta records the pair and ignores a later write with the same application id and a version it has already committed. Second, the application id must be stable across restarts of the same checkpoint. Generating it randomly at start-up silently disables the protection. Delta's behaviour more broadly is covered in Delta Lake on Spark.
Long-horizon deduplication with MERGE
For duplicates that arrive after the in-stream window, make the sink itself reject keys it already holds. Inside foreachBatch you can replace the append with a MERGE that inserts only unmatched ids.
from delta.tables import DeltaTable
def merge_batch(batch_df, batch_id):
target = DeltaTable.forPath(spark, "/lake/payments_clean")
(target.alias("t")
.merge(batch_df.alias("s"),
"t.payment_id = s.payment_id AND t.event_ts >= current_timestamp() - INTERVAL 7 DAYS")
.whenNotMatchedInsertAll()
.execute())A MERGE that inserts only unmatched rows is naturally idempotent for replays too, because a replayed batch finds every id already present. The cost is a join against recent target data on every batch, which raises latency and compute. Partition the target by a date derived from event_ts, or cluster it on that column, and keep the lookback as short as the duplicate tail allows. A seven-day lookback over a table partitioned by day reads seven partitions, not two years.
Monitoring the operator
Every micro-batch produces a progress report, available from query.lastProgress or a StreamingQueryListener. Its stateOperators entry for the dedup operator includes numRowsTotal (keys currently in state), numRowsUpdated (keys added this batch), memoryUsedBytes and numRowsDroppedByWatermark. The eventTime block shows the current watermark.
- State size.
numRowsTotalshould plateau at roughly arrival rate times retention. A line that keeps rising means eviction is not happening: no watermark, a watermark column missing from the key in the plain API, or a watermark that is not advancing. - Late drops. A non-zero
numRowsDroppedByWatermarkmeans real events were discarded as too late, not just duplicates. Alert on it as a rate.
Failure modes
- Unbounded state. Using
dropDuplicateswithout a watermark on a high-cardinality key. The query runs for weeks, then fails on memory or takes an hour to recover. - Event time in the key, different per retry. The plain API with a watermark lets through copies whose timestamps differ. Use
dropDuplicatesWithinWatermarkor fix the producer to stamp time once. - Stalled watermark. The watermark follows the maximum event time seen across a source, so it stops moving when the whole source goes quiet, or when a query has several watermarked inputs and the default min policy waits for the slowest. State then stops being evicted. Check the watermark in progress reports, not just throughput.
- Clock-skewed producers. One client sending timestamps a year in the future drags the watermark forward and every honest event becomes late. Filter impossible timestamps before
withWatermark. - Changing the query shape. The dedup key, the operator and the number of state partitions are bound to the checkpoint. Changing them requires a new checkpoint, which starts with empty state; plan a backfill or an overlap window during which the MERGE sink catches what the empty state misses.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| Longer watermark delay | Catches duplicates spread further apart | State grows linearly with the delay; late real events wait longer |
| Shorter watermark delay | Small state, fast recovery | Duplicates beyond the window pass through; more events counted as late |
| dropDuplicatesWithinWatermark vs event time in key | Handles retries with new timestamps | Requires Spark 3.5.0 or later |
| MERGE sink | Unlimited horizon, replay-safe | Join cost on every batch; latency |
What to do next
- Pick the business identifier that defines 'the same event' and confirm it is assigned once, before any retry, by the producer.
- Measure, from a week of raw data, the duplicate rate and the largest time gap between copies of the same id.
- Choose
dropDuplicatesWithinWatermarkif retries can change the timestamp; otherwisedropDuplicateswith the event-time column in the key. Never use either without a watermark on an unbounded key space. - Set the watermark delay above the measured gap for the bulk of duplicates, and estimate state as arrival rate times retention times bytes per key, where retention is about the delay for
dropDuplicatesand about twice the delay fordropDuplicatesWithinWatermark. - Switch to the RocksDB state store if the estimate does not fit comfortably in executor heap.
- Make the sink idempotent:
txnAppIdplustxnVersionfor Delta appends, or a MERGE that inserts only unmatched ids with a bounded lookback. - Alert on rising
numRowsTotal, non-zeronumRowsDroppedByWatermarkand a watermark that stops advancing. - Kill the driver mid-batch in staging and confirm the output has no duplicates after restart.