05 / 05

What does the transaction coordinator do and what failure scenarios must it tolerate?

Difficulty: 9/10
Transactional processing, Exactly-once, Idempotency

The Transaction Coordinator: Responsibilities, Failures, and Recovery

The transaction coordinator is a Kafka broker component responsible for managing the lifecycle of transactions. Each transactional.id is mapped to a specific coordinator partition in the __transaction_state internal topic, and that partition's leader broker acts as the coordinator for that transactional.id. The coordinator's job is to maintain the state of each transaction (ongoing, prepare-commit, prepare-abort, complete-commit, complete-abort), write transaction markers to the involved partitions, and handle timeouts and fencing. When a producer calls initTransactions(), it finds its coordinator, registers its transactional.id, and receives a producer ID (PID) and epoch. When the producer calls commitTransaction() or abortTransaction(), the coordinator orchestrates the two-phase commit: it writes a prepare marker to the transaction state log, then writes commit or abort markers to the data partitions, then writes a complete marker. If the producer fails mid-transaction, the coordinator eventually times out the transaction and aborts it.

The failure scenarios the coordinator must tolerate are significant. First, the producer can crash or become unreachable. The coordinator handles this with transaction.timeout.ms: if the producer does not commit or abort within the timeout, the coordinator aborts the transaction. This prevents an open transaction from blocking the Last Stable Offset indefinitely. Second, the coordinator itself can crash. Since the transaction state is stored in the __transaction_state topic with replication, a new coordinator can be elected and recover the state. On recovery, the new coordinator may need to complete or abort in-flight transactions; this is why the prepare markers exist. Third, the producer can be fenced. If a new producer instance starts with the same transactional.id, the coordinator bumps the epoch, and the old producer's writes are rejected with a ProducerFencedException. This prevents zombie producers from writing after a restart. Fourth, the coordinator can be slow or overloaded. Since transaction state is partitioned by transactional.id, a hot transactional.id can overload a single coordinator broker. This is why it is important to spread transactional.ids across many partitions and not use a single transactional.id for a high-throughput pipeline.

A common mistake is to assume that the transaction coordinator is a single point of failure. It is not, because the __transaction_state topic is replicated and partitioned, and coordinators can fail over. However, a coordinator failover during a transaction can cause delays while the new coordinator recovers state. Another mistake is to set transaction.timeout.ms too high, which means a stuck transaction can block the LSO for a long time and delay read_committed consumers. Setting it too low can cause legitimate long-running transactions to be aborted. The broker config transaction.max.timeout.ms caps what producers can request; if a producer asks for a higher timeout, the request is rejected. The trade-off is between giving transactions enough time to complete and not blocking consumers. Version note: KIP-447 (Kafka 2.5) improved transaction coordinator scalability by allowing a single producer to write to many partitions and by improving how the coordinator handles large transactions. In KRaft mode (Kafka 3.x), the transaction coordinator is still a broker component, but metadata management is handled by the KRaft controller, which changes some failure and recovery paths. Always check the version and mode when diagnosing coordinator issues.

javascript
  1. 1

    The transaction coordinator manages transaction state, markers, timeouts, and fencing.

  2. 2

    State is stored in the replicated __transaction_state topic, partitioned by transactional.id.

  3. 3

    It tolerates producer crashes via transaction.timeout.ms, which aborts stuck transactions.

  4. 4

    It tolerates its own crash via replica failover and recovery from prepare markers.

  5. 5

    It fences zombie producers by bumping the epoch when a new instance uses the same transactional.id.

  6. 6

    A hot transactional.id can overload a single coordinator; spread transactional.ids.

  7. 7

    KIP-447 (Kafka 2.5) improved scalability; KRaft mode changes metadata but not the coordinator role.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.