skip to content

Transaction Coordinator and Zombie Fencing

The transaction coordinator, the __transaction_state log, and how producer epoch bumps fence zombie instances. This is the how-does-it-really-work follow-up after the transaction API question.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What is the Kafka transaction coordinator, and how does a producer find the one assigned to it?

level: juniorimportance: must knowfreq 55%

answer

  1. Coordinator = broker component, not a process
  2. hash(transactional.id) % 50 -> __transaction_state partition
  3. Leader of that partition = the coordinator
  4. FindCoordinator(type=TRANSACTION)
  5. Deterministic mapping enables fencing

basics

~20 s

The transaction coordinator is a broker-side component that manages a transactional producer's state. A producer locates it by hashing its transactional.id to a partition of the internal __transaction_state topic; the leader of that partition is the producer's coordinator.

solid answer

~40 s

Every transactional producer is served by exactly one transaction coordinator, which is a module running inside a broker (not a separate process). The coordinator owns the transaction lifecycle: assigning producer IDs and epochs, tracking which partitions a transaction touches, and writing commit/abort markers. A producer discovers its coordinator with a FindCoordinator request keyed by its transactional.id. The broker hashes that id modulo the number of partitions of the internal compacted topic __transaction_state (default 50 partitions); the leader broker of the resulting partition is the coordinator. Because the assignment is deterministic on transactional.id, the same logical producer always maps to the same coordinator partition, which is essential for fencing zombies that reuse the id.

go deeper

for a junior

Know it is a broker component found by hashing transactional.id to a __transaction_state partition whose leader is the coordinator.

for a middle

Explain FindCoordinator, the 50-partition compacted topic, and what state the coordinator tracks.

for a senior

Articulate why the deterministic id->coordinator mapping is the foundation for zombie fencing and how failover rebuilds state.

for a principal

Reason about partition-count immutability, replication factor for durability, and the coordinator's role in the two-phase commit protocol overall.

## What a transaction is here A Kafka transaction lets a producer atomically write to multiple topic-partitions (and commit consumer offsets) so that downstream `read_committed` consumers see either all of the writes or none. To coordinate this, Kafka needs a single authoritative place that knows the current state of each in-flight transaction. ## The transaction coordinator The **transaction coordinator** is not a standalone server — it is a component embedded in a normal Kafka broker. Each transactional producer is served by exactly **one** coordinator at a time. The coordinator is responsible for: - Allocating a **producer ID (PID)** and **producer epoch** (via `InitProducerId`). - Recording which topic-partitions are part of the current transaction (`AddPartitionsToTxn`). - Driving the two-phase commit: on `EndTxn`, it writes **commit or abort markers** (control records) into every involved partition. - Persisting all of this state durably so it survives broker restarts/failover. ## The __transaction_state log The coordinator persists state in an internal topic named **`__transaction_state`**. Key properties: - It is **log-compacted**, so the latest state per `transactional.id` is retained. - It defaults to **50 partitions** (`transaction.state.log.num.partitions`) and replication factor `transaction.state.log.replication.factor` (default 3 in production configs). - The key is the `transactional.id`; the value is a serialized `TransactionMetadata` (PID, epoch, state, partitions, timestamps). ## Coordinator discovery (the deterministic mapping) A producer finds its coordinator with a **`FindCoordinator`** request of type TRANSACTION, passing its `transactional.id`. The broker computes: ``` partition = abs(hash(transactional.id)) % numPartitions(__transaction_state) ``` The **leader of that `__transaction_state` partition is the producer's transaction coordinator**. Because the mapping is a pure function of `transactional.id`, the *same* logical producer (even after a restart, even a zombie copy) always resolves to the same coordinator partition. That determinism is what makes **zombie fencing** possible — all producers sharing an id contend at one coordinator, which can compare epochs and reject the stale one. ## Failover If the broker hosting a `__transaction_state` partition leader fails, leadership for that partition moves to a replica; the new leader replays the compacted log to rebuild in-memory `TransactionMetadata` and becomes the coordinator. Producers transparently re-issue `FindCoordinator` and reconnect. ## Edge cases - Changing `transaction.state.log.num.partitions` after data exists reshuffles the hash mapping and is unsafe — treat it as fixed. - A producer without a `transactional.id` (only `enable.idempotence=true`) has **no** transaction coordinator; it gets a PID from a partition leader instead and is idempotent but not transactional.

  • What is __transaction_state and why is it compacted?
    It is the internal topic where the coordinator durably stores per-transactional.id metadata (PID, epoch, state, partitions). Compaction keeps the latest state per id so the coordinator can rebuild memory after failover without unbounded log growth.
  • Does a purely idempotent producer (no transactional.id) have a transaction coordinator?
    No. Idempotence alone gives a PID/epoch from a partition leader for dedup, but transactions and the coordinator only apply when transactional.id is set.

saying these in an interview costs you the question

  • Saying the coordinator is a separate dedicated server/process rather than a broker component
  • Claiming the producer picks any broker or the controller as coordinator
  • Thinking discovery is random rather than a deterministic hash of transactional.id
  • Confusing __transaction_state with __consumer_offsets

context

open as a page

Walk through what InitProducerId does and how the producer epoch is used to fence a previous instance sharing the same transactional.id.

level: seniorimportance: must knowfreq 60%

basics

~20 s

InitProducerId asks the coordinator for a producer ID and epoch. When a producer reuses an existing transactional.id, the coordinator bumps the epoch and aborts any in-flight transaction from the old instance. Brokers then reject writes carrying the now-stale (lower) epoch, fencing the old producer.

open as a page

Which broker and producer configs govern the transaction coordinator and zombie fencing, and what do they control?

level: middleimportance: should knowfreq 35%

basics

~10 s

Producer side: transactional.id (enables fencing) and transaction.timeout.ms. Broker side: transaction.state.log.num.partitions (coordinator sharding), transaction.state.log.replication.factor and min.isr (durability), and transaction.max.timeout.ms (cap on producer timeout).

open as a page

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

level: seniorimportance: should knowfreq 40%

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.

open as a page

In Kafka Streams exactly-once, how did zombie fencing evolve from per-partition transactional.ids to a single producer per instance, and why?

level: principalimportance: nice to knowfreq 22%

basics

~20 s

Original EOS (exactly_once) used one transactional.id (and producer) per input partition, so fencing was per-task. EOS v2 (exactly_once_v2, KIP-447) uses one producer per stream thread/instance and fences via consumer group metadata, drastically cutting producers and connections while keeping correct fencing.

open as a page