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.
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
- Watch
- Try it
- Predict
- Capture
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.
Highlighted lines are the ones running in the diagram right now.
def send(record):record.pid = self.pid # assigned by InitProducerIdrecord.seq = self.next_seq[record.partition]self.next_seq[record.partition] += 1record.epoch = self.epoch # checked by brokerresp = leader.produce(record)if resp == DuplicateSequenceException:return # retry already landedif resp == ProducerFencedException:raise # zombie — give up
def initProducerId(transactional_id):state = self.txn_log[transactional_id]if state.has_open_txn:abort_transaction(state) # roll back zombie's writesstate.epoch += 1 # monotonic, the fencestate.pid = state.pid or new_pid()self.txn_log[transactional_id] = statereturn (state.pid, state.epoch)
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-447group_meta.member_id,)
def handleProduce(record):state = self.producer_state[record.pid]if record.epoch < state.epoch:return ProducerFencedException # cross-restart fenceexpected = state.last_seq[record.partition] + 1if record.seq < expected:return DuplicateSequenceException # within-session dedupeif record.seq > expected:return OutOfOrderSequenceExceptionlog.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.
Designs that use this
- Twitter / X TimelinePush or pull? Both. The canonical fanout problem.
- Uber / Lyft — Match Drivers and RidersMatch a rider to the closest acceptable driver in under 3 s. Geohash, S2, surge.
- Slack / DiscordChannels and history. Push or pull — and how a hot-channel fanout doesn't melt the gateway.
- WhatsApp / MessengerHundreds of millions of long-lived sockets, sub-second 1:1 + group delivery, E2E-encrypted, multi-device, multi-region active-active.