How does a Spark Structured Streaming query achieve end-to-end exactly-once output?
answer
- the guarantee is about the sink, not the records
- something is written before the batch runs
- a batch always covers the same input range
- three parties must cooperate, not one
- one of the built-in sinks is only at-least-once
basics
~20 sSpark writes each micro-batch's input offset range to a write-ahead log under checkpointLocation before processing it, so a restart replays exactly that batch. Exactly-once then holds if the source is replayable by offset and the sink is idempotent or transactional.
solid answer
~50 sExactly-once in Structured Streaming is about the **effect on the sink**, not about a record being processed once. Before running micro-batch *N*, Spark records the deterministic offset range for that batch in `offsets/N` inside `checkpointLocation` — a write-ahead log — and writes `commits/N` only after the batch's sink write succeeds. On restart, a batch with an offsets entry and no commit entry is replayed over the identical input range. That makes the guarantee conditional on two things Spark does not control: the **source** must be replayable by offset (Kafka, files, Delta — a socket is not), and the **sink** must be idempotent or transactional. The built-in file sink qualifies through its `_spark_metadata` log, which makes file commits atomic and lets a re-run batch skip files it already wrote. The Kafka sink is at-least-once, so duplicates are possible after a failure. With `foreachBatch` you get the `batchId` and implement idempotency yourself.
code
python · 10 linesdef write_idempotently(micro_batch_df, batch_id):
(micro_batch_df.write
.format("delta").mode("overwrite")
.option("replaceWhere", f"batch_id = {batch_id}")
.save("/tables/events"))
(stream.writeStream
.foreachBatch(write_idempotently)
.option("checkpointLocation", "/ckpt/events")
.start())go deeper
Know that checkpointLocation is required, that it stores where the query got to, and that it is how a restarted query resumes instead of starting over. Never point two queries at the same directory.
Explain the offsets write-ahead log and the commit log: the input range is recorded before the batch runs and the commit after it succeeds, which is what makes a replayed batch cover identical input.
Show that you check the sink. Be ready to say which of your sinks is genuinely idempotent, why the built-in Kafka sink is at-least-once, and how you would use batchId inside foreachBatch to close the gap.
Own the end-to-end contract across teams: which stages are idempotent, where deduplication happens, what the recovery procedure is when a checkpoint must be rebuilt, and whether source retention is long enough to support the replay window you promise.
## What "exactly-once" means here The phrase is routinely misheard. Spark does **not** promise that a record is read once or that a task runs once — tasks are retried, whole micro-batches are re-executed after a driver failure, and a record can be processed several times. What Structured Streaming offers is *exactly-once semantics with respect to the sink*: the observable output is as if every input record contributed exactly once. Getting there is a joint effort between the engine, the source and the sink, and the engine only owns one third of it. ## The checkpoint directory `option("checkpointLocation", path)` is mandatory for any production query, and the directory it names is the query's durable identity. Its contents: - **`offsets/`** — the write-ahead log. Before micro-batch *N* processes anything, Spark writes `offsets/N` describing the exact input range for that batch (for Kafka, a per-partition start/end offset map). This is written *ahead* of processing, which is the whole point. - **`commits/`** — one file per batch, written only after the batch's output has been successfully committed to the sink. - **`metadata`** — the query's persistent ID. - **`state/`** — the state store data for stateful operators, partitioned and versioned per batch. - **`sources/`** — source-specific bookkeeping such as the file source's seen-files log. On restart, Spark reads the last `offsets` entry. If there is no matching `commits` entry, that batch is re-planned over the *identical* offset range and re-executed. Determinism of the range is what makes the replay safe: batch 412 always covers the same records, so a sink that can recognise or overwrite batch 412's output ends up in the same state whether the batch ran once or three times. Because the checkpoint is the query's identity, two queries must never share a directory, and deleting it restarts the query from `startingOffsets` with empty state — which for a stateful pipeline means recomputing or losing history. ## What the source must provide Replay only works if the source can be asked again for a byte-identical range. Kafka can (offsets are stable per partition). The file source can (its log tracks which files were seen). Delta and other versioned table sources can. A socket source or a rate source cannot, and Spark documents them as test-only for exactly this reason. One consequence that surprises people: for the Kafka source, offsets live in **Spark's checkpoint**, not in Kafka's `__consumer_offsets`. Spark assigns partitions directly rather than joining a consumer group, so Kafka-side consumer-group lag tooling shows nothing for a Structured Streaming query. `startingOffsets` only applies when the checkpoint is empty; on any subsequent start the checkpoint wins. ## What the sink must provide The sink is where the guarantee is usually lost. **File sink — exactly-once.** It maintains a `_spark_metadata` directory alongside the output listing the files committed by each batch. A reader that goes through Spark honours that log, so files written by a batch that failed midway are ignored, and a replayed batch does not double-count. The catch is that the log is Spark's, not the filesystem's: an external tool that globs the directory will see orphaned files from failed attempts. **Kafka sink — at-least-once.** The built-in sink does not wrap its writes in a Kafka transaction, so a replayed batch republishes its records. Deduplicate downstream on a business key, or accept duplicates. **`foreachBatch` — whatever you build.** It hands you an ordinary batch DataFrame plus the `batchId`. The standard patterns are a `MERGE` keyed on business identity, an overwrite of a partition deterministically derived from the batch, or a small table recording the highest committed `batchId` that the write consults transactionally. Note that `foreachBatch` provides at-least-once by default — the function can be invoked more than once for the same `batchId`, and idempotency is entirely your responsibility. **`foreach`** operates per record with `open(partitionId, epochId)`/`process`/`close`, and the same reasoning applies at row granularity. ## Stateful queries State is checkpointed per batch alongside the offsets, so a replayed batch also rewinds state to the version that existed before it ran. That is why a stateful query's correctness after failure depends on the state store checkpoint being intact, and why restoring an old checkpoint against a source whose retention has already expired the corresponding offsets produces a query that cannot start. ## Things that quietly break it Sharing a checkpoint directory between two queries. Deleting the checkpoint to "reset" a query and then reprocessing into a non-idempotent sink. Performing side effects inside a `map` or UDF — those run per task attempt and are re-executed on retry with no batch-level protection. Writing to an external system from `foreachBatch` without using `batchId`. And assuming continuous processing mode carries the same guarantee: it is at-least-once and experimental. ## How to answer the question in an interview Name the three parties — replayable source, offset write-ahead log plus commit log in the checkpoint, idempotent or transactional sink — and then say which of your actual sinks satisfies the third. That last sentence is what separates a memorised answer from an operational one.
- Why does the Kafka sink give only at-least-once?The built-in Kafka sink writes records with a plain producer and does not wrap a micro-batch's writes in a Kafka transaction, so there is no way to atomically discard the output of a batch that failed after partially writing. A replay of that batch republishes its records. Deduplicate downstream on a business key, or use `foreachBatch` and implement transactional publishing yourself.
- Where are Kafka offsets stored for a Structured Streaming query, and why does that surprise operators?In Spark's own checkpoint directory, under `offsets/`, not in Kafka's `__consumer_offsets`. Spark assigns partitions directly rather than joining a consumer group, so Kafka-side lag tooling reports nothing for the query. Monitor lag from `lastProgress` instead. It also means `startingOffsets` applies only when the checkpoint is empty — on any later restart the checkpoint takes precedence.
- How do you make a foreachBatch write idempotent?Use the `batchId` Spark passes in. Either `MERGE` on a business key so a replay updates rather than inserts, overwrite a target partition derived deterministically from the batch, or keep a small committed-batches table that the write consults and updates in the same transaction, skipping any batchId already recorded. Without one of these, `foreachBatch` is at-least-once because the function can be invoked twice for the same batchId.
- What breaks if two queries share one checkpointLocation?The checkpoint is the query's durable identity — its offset log, commit log and state all live there. Two queries writing the same directory interleave their offset and commit entries, so each one recovers onto the other's bookkeeping. The result is skipped or reprocessed input and corrupt state, usually surfacing as an unrecoverable failure on restart. Give every query its own directory.
saying these in an interview costs you the question
- Says exactly-once means each record is processed only once
- Claims the guarantee holds regardless of the sink
- Thinks Kafka's consumer group tracks the query's offsets
- Deletes the checkpoint to restart a query and expects no reprocessing
- Puts external side effects inside a map or UDF and calls it exactly-once