How does Kafka Streams keep consumer offsets, changelog writes, and output topic writes atomically consistent under EOS?
answer
- offsets committed BY the producer, in the txn
- sendOffsetsToTransaction + groupMetadata (KIP-447)
- commit markers = atomic visibility flip
- abort → records replay, no offset advance
- changelog writes also inside the txn
basics
~20 sStreams uses a single Kafka transaction per commit. The output writes, changelog (state) writes, and the consumed input offsets are all added to that transaction and committed together via sendOffsetsToTransaction + commitTransaction, so they advance as one atomic unit.
solid answer
~40 sUnder EOS, Streams treats one commit interval as a transaction on the transactional producer. During the interval it produces output and changelog records inside the transaction. At commit time it calls producer.sendOffsetsToTransaction(consumedOffsets, groupMetadata) to write the input offsets into the same transaction (offsets are stored in the __consumer_offsets topic, written by the producer), then commitTransaction() atomically commits markers across all involved partitions. The key point: offsets are committed by the producer, not the consumer's own commitSync — this is what binds reading position to data writes. If the commit fails or the task crashes, the transaction aborts: no output is visible to read_committed consumers and the input offset is never advanced, so the same records replay. This eliminates the gap where offsets and data could diverge in at-least-once mode.
go deeper
Know that offsets, output, and state are committed together as one unit.
Explain sendOffsetsToTransaction, that the producer (not consumer) commits offsets, and the abort→replay path.
Detail commit markers vs data bytes, transaction.timeout.ms interplay, and KIP-447 group-metadata fencing.
Reason about commit-marker write amplification across wide fan-out, and how changelog-in-txn guarantees state/offset consistency on restore.
## Why naive offset commit is unsafe In at-least-once mode the consumer commits offsets independently of the producer's writes. There's always a window: either output is written but offsets aren't (→ reprocess → duplicates) or offsets are committed but output failed (→ lost records). EOS closes this window by making offset advancement part of the same atomic transaction as the data writes. ## The transactional read-process-write protocol 1. **Begin**: the StreamThread's transactional producer is already initialized (`initTransactions()` once at startup, which fences zombies and recovers any pending transaction). The thread implicitly begins a transaction on first send after a commit. 2. **Process**: as records flow through the topology, output records go to sink topics and state-store updates go to **changelog topics** — all via `producer.send()` inside the open transaction. 3. **Commit offsets into the txn**: at `commit.interval.ms`, Streams calls `producer.sendOffsetsToTransaction(offsetsToCommit, consumerGroupMetadata)`. Crucially, the **producer** writes these offsets to the internal `__consumer_offsets` topic as part of the transaction — the consumer does NOT call `commitSync`. Passing `consumerGroupMetadata` (KIP-447) lets the broker verify group generation and fence zombie instances. 4. **Commit**: `producer.commitTransaction()` writes commit markers to every partition touched (output partitions, changelog partitions, and the offsets partition) atomically. Only after markers land are the records visible to `read_committed` readers. ## What atomicity means here ‘Atomic’ means the **commit markers** make all writes (data + offsets) become visible together or not at all. The data bytes are written incrementally during the transaction, but a `read_committed` consumer ignores them until the commit marker appears (and treats them as discardable if an abort marker appears instead). So consumers observe an all-or-nothing flip. ## Abort / crash path If the task crashes or the commit throws (e.g., `ProducerFencedException`, `InvalidProducerEpochException`, transaction timeout), the in-flight transaction is **aborted**. An abort marker is eventually written; `read_committed` consumers discard the aborted records. Because offsets were never committed, the consumer's fetch position on restart is the last committed offset, so the records replay and are reprocessed into a fresh transaction. The net effect: exactly once. ## Why offsets-via-producer is the linchpin The single most important design point: **offsets are committed by the transactional producer, inside the transaction**, not by the consumer separately. That binding is what makes ‘how far we read’ and ‘what we wrote’ inseparable. Candidates who miss this usually think the consumer still commits offsets — it doesn't, under EOS. ## Edge cases - **Transaction timeout**: `transaction.timeout.ms` (default 60s) must exceed the longest commit interval + processing burst, or the broker aborts a still-live transaction. Streams caps `commit.interval.ms` accordingly. - **Many partitions**: commit markers must be written to every touched partition, so very wide fan-out makes commits heavier. - **Changelog as part of txn**: because state changelogs are in the transaction, a restored state store is always consistent with committed offsets — no half-applied state.
- In EOS, does the consumer call commitSync to advance offsets?No. The transactional producer commits offsets via sendOffsetsToTransaction inside the transaction. The consumer never commits independently — that binding is exactly what makes offsets and writes atomic.
- Why does sendOffsetsToTransaction take consumer group metadata in v2?KIP-447 passes group metadata so the broker can validate the consumer's group generation and fence zombie producers using group coordination, removing the need for one transactional.id per partition.
saying these in an interview costs you the question
- Saying the consumer still commits offsets normally under EOS.
- Believing data is invisible until commit because it isn't written yet — actually it's written but skipped by read_committed until the marker.
- Ignoring that changelog/state writes are also part of the same transaction.
- Confusing 'atomic' with 'single network write' — it's about commit markers flipping visibility.