Kafka transactions enable stream processing applications to execute atomic read-process-write operations across multiple topic partitions and commit input consumer offsets within the same atomic boundary. A dedicated broker called the transaction coordinator tracks transaction states, while downstream consumers set isolation.level=read_committed to avoid reading uncommitted or aborted messages. To eliminate zombie producers caused by network partitions or long garbage collection pauses, Kafka assigns a persistent transactional.id and uses producer epoch increments to fence off superseded producer instances.
The atomic read-process-write loop
In a transactional streaming job (such as consuming from orders and writing to payments):
- The producer registers with the transaction coordinator using a static transactional.id.
- The producer begins a transaction with beginTransaction().
- It writes transformed output records to target topic partitions.
- It calls sendOffsetsToTransaction(), routing consumer offsets directly to the transaction coordinator rather than __consumer_offsets.
- Calling commitTransaction() instructs the coordinator to write a two-phase commit marker across all involved partitions.
Downstream consumers configured with isolation.level=read_committed buffer messages until commit markers appear, ensuring aborted batches are skipped.
Zombie fencing with producer epochs
If a worker node encounters a 45-second Java GC pause, the cluster coordinator may assume it crashed and spin up a replacement container. If the old worker wakes up, two processes would attempt to write under the same identity:
- When the new worker initializes with the same transactional.id, the transaction coordinator increments the producer epoch counter.
- The coordinator assigns the higher epoch to the new instance.
- When the old zombie instance attempts to produce or commit, the broker rejects its requests with a ProducerFencedException.
- The zombie instance is fenced off immediately, preventing split-brain log corruption.
Scope of exactly-once guarantees
Kafka transactions provide exactly-once semantics exclusively within the Kafka boundary (Kafka to Kafka). Once messages leave Kafka for an external database or data warehouse, two-phase commits cannot span the boundary. External sinks must rely on idempotent upserts or unique deduplication keys.