Walk through the two-phase atomic commit Kafka uses to commit a transaction across multiple partitions.
answer
- Phase 1 = PrepareCommit append = commit point
- Phase 2 = WriteTxnMarkers to all partitions
- CompleteCommit closes it out
- crash recovery replays __transaction_state, re-drives markers
- decision atomic in one log append, not a participant vote
basics
~20 sOn commit, the coordinator first durably writes a PrepareCommit record to __transaction_state (phase 1), then writes COMMIT markers to every touched partition and finally a CompleteCommit record (phase 2). The prepared state lets it retry markers after a crash so the commit is atomic.
solid answer
~40 sKafka commits a multi-partition transaction with a two-phase protocol anchored in the __transaction_state log. Phase 1 (prepare): on receiving EndTxn(commit), the coordinator appends a PrepareCommit record listing all participating partitions and fsyncs it; this is the point of no return. Phase 2 (complete): the coordinator sends WriteTxnMarkers requests so each partition leader appends a COMMIT control batch; once all markers are acknowledged, the coordinator appends CompleteCommit. Atomicity comes from durability of the PrepareCommit record plus idempotent, retryable marker writes: if the coordinator crashes after prepare but before all markers are written, the new coordinator replays __transaction_state, sees the prepared (incomplete) transaction, and re-sends the markers. Markers are keyed by PID/epoch so re-delivery is safe. The decision is made once and atomically in phase 1; phase 2 only propagates it.
go deeper
Know commit happens in two steps: record the decision, then write markers everywhere.
Name PrepareCommit/CompleteCommit states and WriteTxnMarkers as the propagation step.
Identify the PrepareCommit append as the atomic commit point and explain crash-recovery re-drive of markers.
Contrast with classic 2PC, reason about the LSO/availability impact of in-flight transactions, and epoch-based zombie fencing.
## Goal 'Atomic commit across partitions' means: every partition the producer wrote to during a transaction ends up *either* with all records committed *or* all aborted — never a mix — even if brokers crash mid-commit. Kafka achieves this with a two-phase commit (2PC) where the **transaction coordinator** is the single decision-maker and the internal **`__transaction_state`** compacted topic is the durable transaction log. ## Background actors - **Producer**: runs the transaction, calls `commitTransaction()`. - **Coordinator**: the broker owning the `__transaction_state` partition for this `transactional.id`. Holds the authoritative `TransactionMetadata` (state machine: Empty -> Ongoing -> PrepareCommit/PrepareAbort -> CompleteCommit/CompleteAbort). - **Partition leaders**: brokers leading the data partitions the transaction wrote to. They received `AddPartitionsToTxn` registrations during the transaction so the coordinator knows the full participant set. ## Phase 1 — Prepare (the atomic decision) 1. Producer sends `EndTxn(commit=true)` to the coordinator. 2. Coordinator transitions the in-memory state to **PrepareCommit** and **appends a PrepareCommit record to `__transaction_state`**, including the list of all participating `TopicPartition`s. This append is durably persisted (replicated, and the log is fsynced per config). 3. **This is the commit point.** Once PrepareCommit is durable, the transaction *will* commit no matter what fails next. Symmetrically, EndTxn(commit=false) -> PrepareAbort, and the transaction will abort. ## Phase 2 — Complete (propagate the decision) 4. Coordinator sends a `WriteTxnMarkers` request to each partition leader. Each leader appends a **COMMIT control batch (marker)** to its log for the given PID/epoch. The marker's presence is what makes the records visible to `read_committed` consumers. 5. The coordinator also writes the marker for the consumer-offsets partitions if the transaction included `sendOffsetsToTransaction` (so offset commits are part of the same atomic unit — the basis of read-process-write exactly-once). 6. When all markers are acknowledged, the coordinator **appends a CompleteCommit record** to `__transaction_state` and transitions to CompleteCommit. The transaction is now fully done. ## Why this is atomic under failure - **Coordinator crash after PrepareCommit, before all markers**: the new coordinator (a different broker becomes leader of the `__transaction_state` partition) replays the log on load. It finds a transaction in PrepareCommit with no CompleteCommit -> it **re-drives phase 2**, re-sending WriteTxnMarkers. Markers are idempotent per (PID, epoch, partition): a leader that already has the marker ignores the duplicate. So the commit eventually completes everywhere. - **A partition leader crash**: the WriteTxnMarkers request is retried against the new leader. - **Producer crash after EndTxn**: irrelevant — the decision is the coordinator's; it finishes phase 2 regardless. The key insight: the *decision* is committed atomically in a single log append (phase 1). Phase 2 is just durable, retryable *propagation* of an already-final decision — unlike classic distributed 2PC, participants cannot vote 'no' at this point, so there is no blocking/uncertainty window for the coordinator. This is sometimes described as a 2PC where the participants pre-agreed by accepting the data earlier. ## Edge cases - **Hanging/long-running transactions** block the Last Stable Offset (LSO) from advancing, stalling `read_committed` consumers; `transaction.max.timeout.ms` and the coordinator's abort of timed-out transactions bound this. - **Zombie producers**: a fenced (old-epoch) producer cannot inject markers because WriteTxnMarkers carries the epoch and stale epochs are rejected.
- At what exact moment is the transaction's outcome decided?When the PrepareCommit (or PrepareAbort) record is durably appended to __transaction_state. After that, phase 2 only propagates the already-final decision.
- How does a newly elected coordinator know to finish an in-flight commit?It loads/replays its __transaction_state partition on becoming leader, finds transactions in PrepareCommit/PrepareAbort without a matching Complete record, and re-sends the WriteTxnMarkers requests idempotently.
- How is this different from classic blocking 2PC?Participants don't vote in phase 2; they already accepted the data earlier, so they cannot reject the decision. There is no in-doubt blocking window — the coordinator unilaterally finalizes and retries until markers land.
saying these in an interview costs you the question
- Saying the commit point is when markers are written (it is the PrepareCommit append).
- Claiming partitions can vote to abort during phase 2.
- Forgetting that recovery re-drives markers from __transaction_state.
- Confusing this with a per-partition independent commit (it must be coordinated).