Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Transactions and zombie fencing

Kafka · Internals

Transactions and zombie fencing

Hardkafka-47
transactionszombie-fencingexactly-onceread-committed

Question

How do Kafka transactions work, and what is zombie fencing?

Solution

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.

PreviousNext