What does setting processing.guarantee=exactly_once_v2 in a Kafka Streams application actually guarantee, and how do you enable it?
answer
- atomic read-process-write in one txn
- offsets + changelog + output committed together
- exactly_once_v2 = KIP-447, one producer per thread
- read_committed for downstream
- effect once, not delivery once
basics
~20 sIt guarantees each input record affects the application's state and output exactly once, even if the app crashes and restarts — no duplicates and no lost updates. You enable it by setting the config processing.guarantee to exactly_once_v2.
solid answer
~40 sSetting processing.guarantee=exactly_once_v2 (EOS) tells Kafka Streams that each input record is reflected in state stores and output topics exactly once, with no duplicates and no lost records, even across crashes and rebalances. Streams achieves this by wrapping each read-process-write cycle in a Kafka transaction: the offsets it consumed, the writes to its changelog (state) topics, and the writes to output topics are all committed atomically. If a task fails mid-cycle, the transaction aborts, downstream consumers reading in read_committed mode never see the aborted writes, and the input offsets are not advanced, so the records are reprocessed cleanly. You only set the one config; Streams configures the underlying transactional producer, idempotence, and read_committed isolation automatically. The default is at_least_once.
go deeper
Know the config name exactly_once_v2 and the one-line guarantee: no duplicates, no lost updates across crashes.
Explain the read-process-write transaction binding offsets + changelog + output atomically, and that downstream needs read_committed.
Contrast v2 vs original exactly_once (per-thread vs per-task producer), broker version requirement, and the auto-configured idempotence/acks/isolation.
Reason about cluster-locality constraints, durability config (min.insync.replicas), latency tradeoffs, and when EOS is/ isn't worth it architecturally.
## The problem EOS solves A Kafka Streams app continuously does a **read-process-write** loop: it reads records from input topics, updates internal **state stores** (which are backed by **changelog topics** in Kafka for recovery), and writes results to **output topics**. Three separate things must stay consistent: (1) the **consumer offsets** marking how far it has read, (2) the writes to changelog topics, and (3) the writes to output topics. Without protection, a crash between any of these steps causes problems. If the app writes output but crashes before committing offsets, on restart it reprocesses the same input and **duplicates** the output (this is **at-least-once**, the default `at_least_once`). If it commits offsets first but crashes before writing output, you **lose** results (at-most-once). ## What exactly_once_v2 guarantees `processing.guarantee=exactly_once_v2` makes the **effect** of each input record on state and output happen **exactly once**. Note: it does NOT mean a record is physically delivered once over the network — the producer may retry sends. It means the observable end-state (state stores + output topics) is as if each record were processed once. ## How it works under the hood Streams wraps each commit interval's work in a single **Kafka transaction** using a transactional producer: 1. Output records and changelog records are produced within the transaction. 2. The **consumed input offsets** are added to the same transaction via `sendOffsetsToTransaction`. 3. `commitTransaction` atomically commits all of the above. Because offsets and data writes share one transaction, they cannot diverge. If the task crashes, the transaction is **aborted**: a transaction marker is written, and downstream consumers configured with `isolation.level=read_committed` skip aborted records. Streams sets read_committed automatically for its own consumers. ## v2 vs the original (deprecated `exactly_once`) The original `exactly_once` (KIP-98/129) used **one transactional producer per input partition (task)**. With many partitions this created many producers, many connections, and slow rebalances. **exactly_once_v2** (KIP-447) uses **one producer per StreamThread (per instance)** and relies on the broker to fence zombies via consumer-group metadata, so it scales far better. v2 requires brokers ≥ 2.5. The older `exactly_once` is deprecated; new apps should use `exactly_once_v2`. ## Enabling it Set a single property: ``` processing.guarantee=exactly_once_v2 ``` Streams then auto-configures: transactional producer with `enable.idempotence=true`, `acks=all`, `isolation.level=read_committed` on consumers, and a default `commit.interval.ms` of 100 ms (vs 30000 ms for at-least-once) to keep latency low. You generally do not set the transactional.id yourself — Streams derives it. ## Edge cases / requirements - Source and target topics, plus changelogs, must live in the **same Kafka cluster** — transactions are single-cluster. - Output topics should have replication and `min.insync.replicas` high enough that `acks=all` is durable. - Downstream consumers that want to see exactly-once results must read with `read_committed`; otherwise they see uncommitted/aborted records. - EOS adds latency (transaction markers + read_committed end-offset lag), traded for correctness.
- Does exactly-once mean a record is sent over the network only once?No. The producer may physically retry sends; idempotence dedups them. EOS guarantees the observable effect on state and output topics is as if processed once, not that bytes traverse the wire once.
- What must downstream consumers do to actually observe exactly-once results?Read with isolation.level=read_committed so they skip records from aborted transactions and only see committed output.
saying these in an interview costs you the question
- Saying EOS means a record is physically delivered/sent only once over the network.
- Claiming it works across multiple Kafka clusters or with mirrored topics.
- Thinking you must manually manage the transactional producer or transactional.id in Streams.
- Forgetting that downstream consumers need read_committed to see exactly-once results.