skip to content

Walk through the two-phase atomic commit Kafka uses to commit a transaction across multiple partitions.

level: seniorimportance: must knowfreq 50%

answer

  1. Phase 1 = PrepareCommit append = commit point
  2. Phase 2 = WriteTxnMarkers to all partitions
  3. CompleteCommit closes it out
  4. crash recovery replays __transaction_state, re-drives markers
  5. decision atomic in one log append, not a participant vote

basics

~20 s

On 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 s

Kafka 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

for a junior

Know commit happens in two steps: record the decision, then write markers everywhere.

for a middle

Name PrepareCommit/CompleteCommit states and WriteTxnMarkers as the propagation step.

for a senior

Identify the PrepareCommit append as the atomic commit point and explain crash-recovery re-drive of markers.

for a principal

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).

context