skip to content

How does the __transaction_state log let a transaction coordinator survive broker failover, and what does it store?

level: seniorimportance: should knowfreq 40%

answer

  1. __transaction_state: internal, compacted, replicated
  2. Value = TransactionMetadata (PID, epoch, state, partitions)
  3. State machine: Ongoing -> PrepareCommit -> CompleteCommit
  4. Failover = leader election + replay rebuilds map
  5. PrepareCommit logged before markers => crash-safe re-drive

basics

~20 s

__transaction_state is an internal, log-compacted, replicated topic keyed by transactional.id. It stores each producer's TransactionMetadata (PID, epoch, state, involved partitions). When a coordinator broker fails, a replica becomes the new partition leader, replays the compacted log to rebuild metadata in memory, and resumes as coordinator.

solid answer

~40 s

The coordinator keeps all durable transaction state in the internal topic __transaction_state, keyed by transactional.id. Each value is a TransactionMetadata record holding the producer ID, producer epoch, current transaction state (Empty, Ongoing, PrepareCommit, PrepareAbort, CompleteCommit, CompleteAbort), the set of partitions enrolled in the current transaction, and timestamps. The topic is log-compacted so the latest metadata per id is preserved and replicated (transaction.state.log.replication.factor, min.insync default for durability). Because state is in a replicated log, coordinator failover is just normal Kafka leader election: when the broker leading a __transaction_state partition dies, an in-sync replica is elected leader, loads (replays) that partition to reconstruct the in-memory transaction map, and becomes the coordinator. The coordinator also writes intent records (PrepareCommit/PrepareAbort) before writing commit/abort markers to data partitions, so an interrupted commit can be completed deterministically after failover.

go deeper

for a junior

Know there's an internal __transaction_state topic that stores transaction info so a new broker can take over.

for a middle

Explain it is compacted/replicated, keyed by transactional.id, holding PID/epoch/state/partitions, and that replay rebuilds the coordinator.

for a senior

Detail the Prepare/Complete intent records and how crash-mid-commit is recovered by re-driving idempotent markers.

for a principal

Reason about replication factor/min.isr for durability, immutable partition count, COORDINATOR_LOAD_IN_PROGRESS behavior, and the state machine's recovery guarantees.

## Why durable state is required A transaction can be mid-flight when a broker dies. To preserve exactly-once semantics across failover, the coordinator's knowledge — which producer, which epoch, which partitions, and how far the commit got — must be **durable and recoverable**, not just in memory. ## The __transaction_state topic Kafka stores this in an **internal, log-compacted, replicated** topic called **`__transaction_state`**: - **Key:** `transactional.id`. - **Value:** a serialized **`TransactionMetadata`** record. - **Compaction:** keeps the latest metadata per id; old versions are garbage-collected so the log stays bounded while always retaining current state. - **Replication:** governed by `transaction.state.log.replication.factor` (3 in production defaults) and `transaction.state.log.min.isr`, giving the same durability story as any replicated Kafka topic. ## What TransactionMetadata contains - **Producer ID (PID)** and **producer epoch** — the fencing identity. - **Transaction state**, a state machine: `Empty -> Ongoing -> PrepareCommit/PrepareAbort -> CompleteCommit/CompleteAbort` (with `Dead`/`PrepareEpochFence` internals). - **The set of topic-partitions** added to the current transaction (so the coordinator knows where to write markers). - **Timestamps** used for `transaction.timeout.ms` enforcement. ## The commit protocol and why intent is logged On `EndTxn(commit)` the coordinator: 1. Writes a **PrepareCommit** record to `__transaction_state` (durable intent). 2. Writes **commit markers** (control records) to every data partition in the transaction. 3. Writes **CompleteCommit** to `__transaction_state`. If the broker crashes between steps 1 and 3, the new coordinator sees `PrepareCommit` on recovery and **re-drives** the commit to completion — markers are idempotent, so re-writing them is safe. This is what makes the two-phase commit crash-tolerant. ## Failover walkthrough 1. Broker leading `__transaction_state` partition P fails. 2. Controller elects an in-sync replica as the new leader of P. 3. The new leader **loads/replays** partition P, rebuilding the in-memory `transactional.id -> TransactionMetadata` map. 4. It becomes the coordinator for every id hashing to P; producers re-issue `FindCoordinator` and reconnect. 5. Any transaction in a `Prepare*` state is completed deterministically. ## Edge cases - **Loading delay:** while a partition is being loaded, the coordinator returns `COORDINATOR_LOAD_IN_PROGRESS`; clients retry. - **Partition count is effectively immutable** — it defines the id->partition hash, so it must not change once in use. - **Replication factor too low** (e.g., 1 in a dev default) undermines durability; production must set 3. - **Compaction lag** doesn't hurt correctness because the latest record per key always wins on replay.

  • Why does the coordinator write a PrepareCommit record before writing commit markers to data partitions?
    So the commit intent is durable: if the broker crashes mid-commit, the new coordinator sees PrepareCommit on replay and re-drives the (idempotent) markers to completion, preserving atomicity.
  • What does a client receive while the new coordinator is still loading the __transaction_state partition?
    COORDINATOR_LOAD_IN_PROGRESS, a retriable error; the client backs off and retries until loading finishes.

saying these in an interview costs you the question

  • Claiming coordinator state lives only in memory or in ZooKeeper/KRaft metadata rather than __transaction_state
  • Saying the topic is regular (non-compacted) or single-replica by design
  • Thinking failover loses in-flight transactions instead of replaying them
  • Asserting commit markers go to __transaction_state (they go to the data partitions; intent goes to __transaction_state)

context