End-to-End Exactly-Once — Replayable Sources, 2PC & Idempotent Sinks
Exactly-once in streaming never means each record is physically processed once. After a crash, records are replayed. It means the effect of each record on state and on externally visible output happens once. That property is end-to-end only if three links hold: a replayable source with offsets you can rewind, checkpointed state consistent with those offsets, and a sink that is idempotent or transactional and commits only what a completed checkpoint covers. Break any link and you are back to at-least-once or at-most-once.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What are the three requirements for end-to-end exactly-once?
Answer
A replayable source, consistent snapshots of state and source positions, and an idempotent or transactional sink that publishes only output covered by a completed checkpoint.
L2
Explain a two-phase-commit sink in a stream processor.
Answer
Each checkpoint period writes into a transaction. On the barrier the sink pre-commits, flushes, makes the data durable but invisible, and opens a new transaction. When the checkpoint completes, the sink commits. On failure, uncommitted transactions abort and the source replays from the checkpoint.
L3
Why does exactly-once to Kafka add latency?
Answer
Output is invisible to read_committed consumers until the transaction commits, which happens when the checkpoint completes. Latency is about the checkpoint interval plus commit time.
L4
When do you use an idempotent sink, and when a transactional one?
Answer
Idempotent when results have a natural key and overwrites are fine: aggregates and materialized views. Transactional when output is append-only, or downstream must never see intermediate or duplicate records.
L5
How do Kafka transactions fence a zombie producer?
Answer
A transactional id is registered with the transaction coordinator, which bumps an epoch on each new producer session. Writes or commits from an older epoch are rejected.
L6
A job calls a payment API inside a map function. What happens on failure?
Answer
The call is replayed, so the payment may run twice. Pass an idempotency key derived from the event id, or emit a command to a topic consumed by an idempotent worker.
L7
How does Spark Structured Streaming get exactly-once to a database?
Answer
Replayable source offsets and state are tracked per micro-batch in the checkpoint. In foreachBatch, write the batch and its batch id in one database transaction and skip batches already recorded.
Failure modes
HTTP inside the operator
Recovery replays the call. The engine cannot roll back a POST.
read_uncommitted consumers
They observe aborted output, so the transaction work was wasted.
Dedupe TTL shorter than replay
A retry that outlives the seen-set is accepted again.
Transaction timeout during an outage
The broker aborts transactions older than transaction.timeout.ms. The timeout must exceed the checkpoint interval plus the outage you will sit through.
Misconceptions
Exactly-once means the record is processed once.
The record is replayed. The effect on state and on visible output happens once.
An upsert of a counter increment is exactly-once.
An upsert of a deterministic total is fine. INSERT plus one on every delivery is not.
Kafka transactions cover an email or a charge.
They cover Kafka output. Side effects need their own idempotency key or an outbox.
Interviewer traps
Declaring exactly-once because checkpoints are on.
Ask what the sink does with a replay of the uncheckpointed tail.
Reusing one transactional id across subtasks.
Prefix it with the subtask id or you fence yourself.
Design scenario
Same prompt for every reader.
Requirements
Replayable source, checkpointed offsets and state, a transactional Kafka sink with read_committed consumers, and an idempotency key on the payment call.
Traffic / scale
Thousands of orders per second. Checkpoint about every minute.
Latency
Billing readers can wait for the checkpoint commit. The payment call cannot.
Consistency
Committed Kafka output matches the completed checkpoint. Aborted transactions stay invisible.
Availability
A zombie producer from the crashed attempt is fenced by transactional id and epoch.
Failure assumptions
- The process dies after writing and before the next checkpoint completes.
- Downstream consumers might forget read_committed.
Constraints
- Do not call the payment API inside map with no idempotency key.
- Do not treat a database increment as an upsert.
Prompt
Orders must land in Kafka for billing and also notify a payment HTTP API. A crash after four records, with a checkpoint after two, must not double-charge or double-count.
API
What is the idempotency key on the payment request?
Data
Which records are inside the open transaction at the crash?
Architecture
Which hop is two-phase commit, and which hop is an idempotent worker?
What a reader is allowed to see
Prefer
Commit with the checkpoint, or overwrite by key
A two-phase sink hides the open transaction until the checkpoint completes. An upsert converges because the write is keyed and deterministic.
- Recovery aborts the dangling transaction.
- Replay starts at the offset stored in the checkpoint.
- read_committed consumers never add the aborted rows.
Alternative
Append, then hope
The rows written after the checkpoint and before the crash are appended again on replay. The sum is wrong and the log looks fine.
- An HTTP POST inside map is this shape.
- An INSERT that increments a counter is this shape.
- A dedupe set shorter than the replay horizon accepts the duplicate later.
Three links, or it is not end to end
The processor view. The log's own acks and offsets are a different lesson.
- 1
Rewindable source
Kafka offsets, Kinesis sequence numbers, or files. A webhook is not a source until you land it in a log. - 2
Snapshot matches those offsets
The barrier cut from the state lesson. Auto-commit on its own breaks the link. - 3
Sink publishes only that cut
Upsert by a natural key, or commit a transaction when the checkpoint completes. - 4
Crash aborts the open attempt
Uncommitted output stays invisible. Replay rewrites it under a new transaction or the same key.
Three guarantees
| Guarantee | Mechanism | On failure | Cost | Acceptable for |
|---|---|---|---|---|
| At-most-once | Commit offsets before processing. No replay | Loss | Lowest | Sampling, best-effort telemetry |
| At-least-once | Process, then commit. Replay from the last commit | Duplicates | Low | Idempotent consumers, counters with dedupe downstream |
| Exactly-once, effectively | Replay plus consistent snapshots plus an idempotent or transactional sink | No loss, no visible duplicates | Commit waits for the checkpoint | Money, inventory, billing, aggregates that feed decisions |
The producer and consumer view of these words is Kafka delivery semantics. This page is the stream processor: how the engine and the sink cooperate. CDC sinks that dedupe on a key are exactly-once CDC. A side effect that must commit with a database row is the transactional outbox inside event-driven architecture. The database protocol named two-phase commit, with a coordinator and participants, is two-phase commit. A stream sink uses the same commit shape and a different recovery story.
A two-phase sink
Decisions
- 1
1. Rewindable source offsets
- next2. State matches that cut
- 2
2. State matches that cut
- next3. Sink writes open txn
- 3
3. Sink writes open txn
- next4. Barrier, sink pre-commits
- 4
4. Barrier, sink pre-commits
- next5. Tasks ack the checkpoint
- 5
5. Tasks ack the checkpoint
- next6. Checkpoint complete?
- ?
6. Checkpoint complete?
- yes7a. Commit the txn
- crash7b. Abort txn and replay
- 7
7a. Commit the txn
- 8
7b. Abort txn and replay
Lesson map
End-to-End Exactly-Once — Replayable Sources, 2PC & Idempotent Sinks
Exactly-once in streaming never means each record is physically processed once. After a crash, records are replayed. It means the effect of each record on state and on externally visible output happens once. That property is end-to-end only if three links hold: a replayable source with offsets you can rewind, checkpointed state consistent with those offsets, and a sink that is idempotent or transactional and commits only what a completed checkpoint covers. Break any link and you are back to at-least-once or at-most-once.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB off["1. Offsets and state in checkpoint n"] txn["2. Open txn, pre-commit on barrier"] commit["3. Ack, then commit T_n"] abort["4. Crash: abort T_n and replay"] off -->|1. Offsets and state in checkpoint n| txn txn -->|2. Open txn, pre-commit on barrier to| commit txn -->|2. Open txn, pre-commit on barrier to| abort
Flink's two-phase sink, and the newer sink committer API, does this: begin a transaction, pre-commit on the barrier, commit on checkpoint complete, abort on failure.
Readers see output only after commit, so end-to-end latency includes the checkpoint interval. A one-minute checkpoint means read_committed consumers see results up to a minute late.
Commit must be recoverable. If the job dies after the checkpoint is complete and before the commit call, recovery must commit that pre-committed transaction again. The commit itself is idempotent.
Press Run. Snippets must be self-contained — no network, files, or native modules.
AppendSink rows=7 sum=220 DUPLICATES
UpsertSink rows=5 sum=150 OK
TwoPhaseSink rows=5 sum=150 OKThe append sink double-wrote o3 and o4. They were written, the job crashed, and replay started at offset 2. The upsert converged because the writes are keyed. The two-phase sink never showed the uncommitted writes and aborted them on recovery.
Sink strategies
| Strategy | How | Pros | Cons | Typical sinks |
|---|---|---|---|---|
| Idempotent upsert | A deterministic key, such as key plus window start. Overwrite | Simple, low latency, no coordination | Needs a natural key. Intermediate values are visible. Non-deterministic logic breaks it | Postgres ON CONFLICT, Cassandra, Elasticsearch by id, Redis SET |
| Transactional two-phase commit | Write in a transaction. Commit when the checkpoint completes | Works for append-only output | Latency equals the checkpoint interval. Transaction timeouts. The sink must support transactions | Kafka transactions, some JDBC, file sinks that rename |
| Atomic file commit | Write temp files. Publish on checkpoint | Good for lakes | Small files. Commit latency | Flink FileSink, Iceberg, Delta, Hudi |
| Batch-id idempotence | Store the last committed batch id with the data in one transaction | Natural for micro-batch | The sink must store metadata | Spark foreachBatch writing the batch id in the same transaction |
| Downstream dedupe | A unique event id. The consumer drops repeats | Works with any at-least-once pipe | The dedupe state must be bounded | Webhooks, notifications, legacy APIs |
Kafka transactions and a TTL seen-set
When source and sink are both Kafka, consume-transform-produce can be one transaction: output records and consumer offsets commit together via sendOffsetsToTransaction. The transactional.id epoch fences a zombie instance. Downstream must use isolation.level=read_committed or it will read aborted output. Kafka Streams wraps this as processing.guarantee=exactly_once_v2. Flink's Kafka sink with exactly-once delivery uses one transaction per checkpoint.
When the sink cannot join a transaction, dedupe by a stable event id inside a bounded window.
Press Run. Snippets must be self-contained — no network, files, or native modules.
read_uncommitted: 6 records sum=170
read_committed : 4 records sum=100
dedupe(ttl=60s) kept: e1@0, e2@5, e3@30, e1@90 (e1@90 re-accepted: retry outlived TTL)
seen-set size at end: 1e1 at time 90 was accepted again because the retry outlived the 60 second TTL. Dedupe is only as good as its window. Size it to the maximum retry horizon, or make the effect itself idempotent.
Where duplicates sneak in
| Where | Why | Mitigation |
|---|---|---|
| Side effects in operators | Replayed on recovery. Not part of any transaction | An idempotent sink with a request id, or an outbox topic |
| Non-deterministic logic | now(), random, or unordered iteration changes the output on replay | Event time, seeded randomness, deterministic keys |
| Kafka transaction timeout | The broker aborts transactions older than transaction.timeout.ms, capped by transaction.max.timeout.ms (15 minutes by default) | Timeout greater than the checkpoint interval plus the outage you expect |
| read_uncommitted consumers | They read aborted data | Enforce read_committed |
| Source is not replayable | Webhooks, UDP, MQTT QoS 0 | Land in Kafka or Kinesis first, then process |
| Spark custom sinks | foreach with no idempotence | foreachBatch plus the batch id plus an upsert or a transaction |
Pitfalls
From the Python output, say which orders the append sink wrote twice and which checkpoint offset the replay used. Then say which of those rows a read_committed consumer would have seen if the sink had been the two-phase one.
Interview Q&A
What are the three requirements for end-to-end exactly-once?
Answer
A replayable source, consistent snapshots of state and source positions, and an idempotent or transactional sink that publishes only output covered by a completed checkpoint.
Explain a stream processor's two-phase-commit sink.
Answer
Each checkpoint period writes into a transaction. On the barrier the sink pre-commits, makes the data durable but invisible, and opens a new transaction. When the coordinator reports the checkpoint complete, the sink commits. On failure, uncommitted transactions abort and the source replays from the checkpoint.
Why does exactly-once to Kafka add latency?
Answer
Output is invisible to read_committed consumers until the transaction commits, which happens at checkpoint completion. Latency is about the checkpoint interval plus commit time.
Idempotent sink versus transactional sink. When each?
Answer
Idempotent when results have a natural key and overwrites are fine, such as aggregates and materialized views. Lower latency, simpler. Transactional when output is append-only, or downstream must never see intermediate or duplicate records.
How do Kafka transactions prevent zombie producers?
Answer
A transactional id is registered with the transaction coordinator, which bumps an epoch on each new producer session. Writes or commits from an older epoch are rejected.
A job calls a payment API inside a map function. What happens on failure?
Answer
The call is replayed after recovery, so payments may run twice. Pass an idempotency key derived from the event id, or emit a command to a topic consumed by an idempotent payment worker.
How does Spark Structured Streaming get exactly-once to a database?
Answer
Replayable source offsets and state are tracked per micro-batch in the checkpoint location. With foreachBatch, write the batch and its batch id in one database transaction and skip batches already recorded.