How does Kafka make a multi-partition transactional write atomic? Describe the coordinator, control records, and Last Stable Offset.
answer
- coordinator = leader of __transaction_state partition
- AddPartitionsToTxn registers each partition
- data written immediately; markers on commit
- two-phase: PREPARE → WriteTxnMarkers → COMPLETE
- LSO + aborted-txn index gate read_committed
basics
~20 sA transaction coordinator (a broker) tracks the transaction in the __transaction_state log. Records are written to each partition immediately, then the coordinator writes commit or abort control records (markers) to every involved partition. read_committed consumers only read up to the Last Stable Offset, so they never see uncommitted or aborted data.
solid answer
~50 sEach transactional.id maps to a transaction coordinator — the broker leading the relevant __transaction_state partition. The producer registers each new partition it writes with an AddPartitionsToTxn request, and the coordinator persists the transaction's state and partition set. Data records are appended to the actual data partitions as they are sent (so they physically exist before commit), tagged with the producer id and epoch. On commitTransaction the coordinator performs a two-phase protocol: it logs PREPARE_COMMIT to __transaction_state, then writes a commit control record (transaction marker) into each involved partition, then logs COMPLETE_COMMIT. Abort is symmetric with abort markers. Consumers with isolation.level=read_committed only read up to the Last Stable Offset — the offset just before the first still-open transaction — and use the markers plus an aborted-transaction index to filter out aborted records. This gives all-or-nothing visibility across partitions even though writes landed incrementally.
go deeper
Know there is a coordinator and that commit/abort markers decide visibility.
Explain that data is written immediately and read_committed gates on the LSO.
Describe AddPartitionsToTxn, control records, and the aborted-transaction index for skipping aborted records.
Walk the full two-phase coordinator protocol, crash recovery via durable PREPARE state, and the throughput/latency tradeoffs.
## The actors - **Transaction coordinator**: a broker-side component. Each `transactional.id` is hashed to a partition of the internal `__transaction_state` topic; the leader of that partition is that id's coordinator. It owns the transaction's lifecycle and durable state. - **`__transaction_state`**: a compacted internal topic storing, per transactional.id, the current state (Empty, Ongoing, PrepareCommit, PrepareAbort, CompleteCommit, CompleteAbort), epoch, timeout, and the set of partitions in the transaction. - **Data partitions**: the user topic-partitions actually being written. ## Step by step 1. **initTransactions()** — producer finds its coordinator, gets a producer id + epoch, and the coordinator recovers any prior transaction. 2. **beginTransaction()** — purely client-side. 3. **First write to a partition** — before sending data there, the producer sends `AddPartitionsToTxn` to the coordinator, which records that partition in the transaction's partition set (persisted to `__transaction_state`). This is how the coordinator later knows where to write markers. 4. **send(...)** — data records are appended to the leader of each data partition **immediately**, stamped with producer id + epoch and a sequence number. They are physically in the log but not yet 'committed' from a read_committed viewpoint. 5. **commitTransaction()** — the producer asks the coordinator to commit. The coordinator runs a mini **two-phase commit**: a. Append `PREPARE_COMMIT` to `__transaction_state` (durable intent). b. Send a **`WriteTxnMarkers`** request to each partition's leader, which appends a **commit control record** (a 'transaction marker') into that partition's log. c. Once all markers are acknowledged, append `COMPLETE_COMMIT` to `__transaction_state`. Because the intent is durably logged before markers are written, a coordinator crash mid-commit is recoverable — on recovery it re-drives the markers. 6. **abortTransaction()** — identical but with `PREPARE_ABORT` and **abort** control records. ## Control records (markers) Control records are special messages in the data partition that consumers never deliver to application code; they only signal 'transaction X with this producer id/epoch committed/aborted here'. They let consumers resolve the fate of the preceding transactional data records in that partition. ## Last Stable Offset (LSO) and read_committed - The **LSO** is the highest offset such that all transactions below it are resolved (committed or aborted). It sits just before the earliest **open** transaction on that partition. - A `read_committed` consumer fetches only up to the LSO, so it never sees data from an in-flight transaction. - Brokers also maintain an **aborted transactions index** per segment; on fetch, the broker returns the list of aborted producer-id ranges so the consumer can **skip aborted records** that lie below the LSO. - A `read_uncommitted` consumer ignores all of this and reads to the high watermark, seeing aborted and uncommitted data. ## Why it is atomic across partitions The records on different partitions become *visible* only after their commit markers are written and the LSO advances past them. Since the coordinator writes markers for all involved partitions as one logical operation (driven from durably-logged intent), readers flip from 'see nothing' to 'see everything' — no partial visibility. A failure before PREPARE_COMMIT yields an abort (eventually via timeout); a failure after is re-driven to completion. This is the heart of atomic multi-partition writes. ## Costs / tradeoffs - Extra latency from coordinator round-trips and marker writes (commit is not free). - read_committed consumers can lag at the LSO when long transactions are open. - `__transaction_state` and marker traffic add load; keep transactions short and reasonably sized.
- If data records are appended to partitions before commit, why don't read_committed consumers see them?Because read_committed consumers only fetch up to the Last Stable Offset, which sits before the earliest open transaction. The records sit above the LSO until commit markers are written and the LSO advances past them.
- How does the system recover if the coordinator crashes after PREPARE_COMMIT but before all markers are written?The PREPARE_COMMIT intent is durably in __transaction_state, so on recovery the coordinator re-drives WriteTxnMarkers to any partitions missing their commit marker, then logs COMPLETE_COMMIT — the commit is idempotently completed, never partially applied.
saying these in an interview costs you the question
- Saying data is buffered on the broker and only written at commit — it is appended immediately, just gated by the LSO.
- Claiming read_committed consumers filter purely client-side without broker-provided aborted-transaction info.
- Describing it as full distributed XA/2PC across clusters; it is a single-cluster coordinator-driven protocol.
- Forgetting that AddPartitionsToTxn is what tells the coordinator which partitions need markers.