Exactly-once — three monotonic counters

Exactly-once is not one mechanism but three cooperating fences: producer sequence numbers dedupe retries within a session, producer epoch fences zombies across restarts, and consumer-group generation in sendOffsetsToTransaction fences zombie tasks after rebalance.

Previously

Retries can duplicate. Restarts can zombie. Rebalances can hand the same partition to a stale consumer. Exactly-once isn't one mechanism — it's three monotonic counters, each fencing a different failure mode.

Scene 08

Exactly-once — three monotonic counters

  1. Watch
  2. Try it
  3. Predict
  4. Capture
TOPIC A · INPUT0123offset 0Streams taskread → process → writeConsumerreads Topic A · inputProducerwrites Topic B · outputTxn CoordinatorstateIDLEFENCESPIDPID + seq#Epochtransactional.idGroup-Gengroup-gen (KIP-447)TOPIC B · OUTPUTLSO=0Downstreamisolation = read_committed
What to watch for

Watch one transaction at a time: the consumer reads a cell from Topic A, the producer writes a transformed cell to Topic B (uncommitted, dashed), the coordinator walks open → prepareCommit → committed, and the cell flips solid as the LSO snaps forward. The read_committed downstream below sees the batch appear atomically — that's what the LSO line is gating.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Producer.send
(pid, seq) tags every record — dedupe within a session
def send(record):
record.pid = self.pid # assigned by InitProducerId
record.seq = self.next_seq[record.partition]
self.next_seq[record.partition] += 1
record.epoch = self.epoch # checked by broker
resp = leader.produce(record)
if resp == DuplicateSequenceException:
return # retry already landed
if resp == ProducerFencedException:
raise # zombie — give up
Coordinator.initProducerId
epoch bump — every restart fences the previous incarnation
def initProducerId(transactional_id):
state = self.txn_log[transactional_id]
if state.has_open_txn:
abort_transaction(state) # roll back zombie's writes
state.epoch += 1 # monotonic, the fence
state.pid = state.pid or new_pid()
self.txn_log[transactional_id] = state
return (state.pid, state.epoch)
Producer.sendOffsetsToTransaction
the offset-commit path that carries the group generation
def sendOffsetsToTransaction(offsets, group_meta):
# group_meta = (group_id, generation, member_id)
coordinator.addOffsetsToTxn(
self.transactional_id, self.epoch,
group_meta.group_id,
)
coordinator.txnOffsetCommit(
offsets,
group_meta.generation, # KIP-447
group_meta.member_id,
)
Broker.handleProduce
three independent checks — drop any one, a zombie leaks
def handleProduce(record):
state = self.producer_state[record.pid]
if record.epoch < state.epoch:
return ProducerFencedException # cross-restart fence
expected = state.last_seq[record.partition] + 1
if record.seq < expected:
return DuplicateSequenceException # within-session dedupe
if record.seq > expected:
return OutOfOrderSequenceException
log.append(record)
state.last_seq[record.partition] = record.seq

Where this sits in Build Kafka

Scene 08 of 13, in the Guarantees act — Three monotonic counters; one design canvas.. PID, epoch, group-generation — three independent fences against zombies.

Up next. Design canvas — you've seen every knob. Time to assemble them for a real workload, name the trade you're making, and watch the warning chips light up when it's contradictory.

All 13 scenes in Build Kafka · Every curriculum

Built with Arqly
Every scene in Build Kafka builds on the one before it.All 13 Build Kafka scenes