skip to content

How does the broker actually detect and drop duplicate batches from an idempotent producer?

level: middleimportance: must knowfreq 70%

answer

  1. PID + epoch + per-partition seq
  2. leader tracks last seq
  3. seq==last+1 append, <=last drop, >last+1 reject
  4. OutOfOrderSequenceException on gap
  5. cache of last 5 -> in-flight<=5

basics

~20 s

Each producer gets a producer ID (PID). Each batch carries a per-partition sequence number that increases by one. The broker remembers the last sequence written per (PID, partition) and discards any batch it has already seen.

solid answer

~50 s

On its first request the producer obtains a unique producer ID (PID) from the broker via an InitProducerId request (the partition leader/coordinator assigns it). For every partition the producer writes to, it stamps each record batch with a monotonically increasing sequence number starting at 0. The leader tracks, per (PID, producer epoch, partition), the last sequence number it has appended. When a batch arrives: if its base sequence is exactly last+1, it's appended; if it equals an already-written sequence, it's a retry and the broker returns success without re-appending (a duplicate, dropped); if it's greater than last+1, there's a gap and the broker rejects with OutOfOrderSequenceException. This is how exactly-once-per-partition is achieved for retries. The broker keeps a small cache (the last 5 batches' metadata, hence in-flight <= 5) so it can recognize duplicates of recent batches.

go deeper

for a junior

Know there's a PID and an increasing number that lets the broker spot duplicates.

for a middle

Explain the per-partition sequence and the append/drop/reject decision.

for a senior

Tie the 5-batch cache to the in-flight cap and explain OutOfOrderSequenceException and leader-change recovery.

for a principal

Reason about PID expiration, replicated producer state, and the durability interplay with acks=all.

## The two identifiers Idempotent dedup rests on two pieces of metadata attached to every produced **record batch**: 1. **Producer ID (PID)** — a unique 64-bit ID the broker assigns to the producer. The producer requests it via an **InitProducerId** API call before its first send. For a plain idempotent (non-transactional) producer the PID is allocated fresh each session; for a transactional producer it's tied to the `transactional.id` and survives restarts. 2. **Producer epoch** — a short integer bumped when a producer is fenced/restarted (mainly relevant for transactions); together (PID, epoch) identify a producer instance. 3. **Sequence number** — a per-**partition** integer the producer assigns to each batch, starting at **0** and increasing by **1** per batch (it's actually the base sequence of the batch; records inside increment from there). ## What the broker remembers The partition **leader** maintains, for each `(PID, epoch, partition)`, the metadata of the **last 5** appended batches (base sequence, offset, timestamp, etc.). Five is not arbitrary — it's why `max.in.flight.requests.per.connection` must be `<= 5` with idempotence: the broker can only recognize duplicates among the most recent few in-flight batches. ## The decision the leader makes per incoming batch Let `lastSeq` be the last sequence written for that (PID, epoch, partition), and `baseSeq` the incoming batch's base sequence: - **baseSeq == lastSeq + 1** → in order → **append**, advance lastSeq. - **baseSeq <= lastSeq** (within the cached window) → this batch was already written (a retry whose ack was lost) → **drop it, return the original offset/success**. No duplicate is written. This is the core dedup. - **baseSeq > lastSeq + 1** → a gap; an earlier batch was lost or skipped → reject with **OutOfOrderSequenceException**. The producer treats this as fatal for that session (it can't safely fill the gap) and surfaces an error. - **Duplicate older than the cached window** → broker can't tell; it returns **DUPLICATE_SEQUENCE_NUMBER** / handles conservatively. This is the reason ordering and bounded in-flight matter. ## Why this yields exactly-once *per partition* Retries are the only way duplicates would otherwise occur, and the sequence check makes every retry recognizable and droppable. Combined with `acks=all` (so an acked write is durable on all in-sync replicas and won't be lost), you get **exactly-once writes to that partition** for the producer's session. ## Edge cases - **Log retention / segment deletion**: if the producer is idle and its state expires (`producer.id.expiration.ms`, default 7 days), the PID state is dropped; a later batch may be treated as new. This is why very-low-throughput producers can theoretically see UNKNOWN_PRODUCER_ID, handled by the client. - **Leader change**: producer state is part of the replicated log/snapshot, so a new leader rebuilds the sequence state and dedup continues to work.

  • Why is max.in.flight.requests.per.connection capped at 5 for idempotence?
    The broker only caches metadata for the last 5 in-flight batches per producer-partition, so it can only recognize and dedup retries within that window. More than 5 outstanding could let a duplicate slip past the cache.
  • What does OutOfOrderSequenceException mean?
    The broker received a batch with a sequence number higher than expected (a gap), meaning an earlier batch was lost. The producer can't safely fill the gap, so it's surfaced as a (typically fatal) error for that session.

saying these in an interview costs you the question

  • Saying the sequence number is global rather than per-partition.
  • Claiming the broker stores every sequence forever — it caches only recent batches (last 5).
  • Confusing producer epoch (fencing) with the sequence number (dedup ordering).
  • Saying duplicates are detected on the producer side — detection happens on the broker (leader).

context