Home › Data Engineering › Kafka Exactly-Once Broken by Producer Retries — Fix
Advanced 5 min · September 23, 2026

Kafka Exactly-Once Broken by Producer Retries — Fix

Duplicates despite idempotence plus transactions? Learn fencing, read_committed gaps, and retry storms that break EOS in production..

N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Drawn from code that ran under real load.

Follow
✓ Production
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
Before you start⏱ 25 min
  • ✓A transactional producer you can configure
  • ✓Consumer groups and offset commit basics
  • ✓Access to producer logs and downstream counts
 ● Production Incident 🔎 Debug Guide
⚡Quick Answer
  • EOS needs three parts: idempotent producer, fenced transactions, and read_committed consumers — miss one and duplicates flow
  • enable.idempotence dedups retried sends per partition; keep max.in.flight at 5 or below or ordering breaks
  • Unique transactional.id per instance lets the broker fence zombies; shared IDs make writers kill each other
  • Unbounded retries replay past broker dedup windows after blips — bound them and keep handlers idempotent
✦ Definition~90s read
What is Kafka Exactly-Once Semantics Broken by Producer Retries?

Exactly-once semantics (EOS) means each record takes effect once despite retries, crashes, and failovers. Kafka builds it from three mechanisms. Idempotent producers tag every batch with sequence numbers so the broker drops retried duplicates. Transactions group multi-partition writes plus consumer offset commits into atomic units with commit/abort markers.

★
Imagine a bank where tellers number deposit slips in order and the vault rejects numbers already filed — that is idempotence.

The read_committed isolation level makes consumers skip aborted and pending records.

Fencing is the bodyguard: each transactional.id carries an epoch the broker bumps on takeover, so zombie writers with stale epochs get rejected instead of duplicating. Retry bounds (delivery.timeout.ms, backoff) keep replays inside the broker's dedup window.

None of these work alone — idempotence without transactions duplicates across restarts, transactions without read_committed leak ghosts, and unbounded retries outrun every dedup structure.

True end-to-end EOS adds one more layer outside Kafka: idempotent downstream handlers keyed on business IDs. The broker guarantees each record is delivered once to your logic; your logic must guarantee applying it twice is harmless. Teams that skip this last layer discover that 'exactly-once' ended at their consumer while duplicates thrived in their database.

Plain-English First

Imagine a bank where tellers number deposit slips in order and the vault rejects numbers already filed — that is idempotence. Two-account transfers go in sealed envelopes filed whole or rejected whole — those are transactions. Junior clerks see only filed envelopes, never drafts — that is read_committed. But shared ID badges, out-of-order numbering in a rush, or clerks reading the trash move money twice.

Your ledger shows 4,800 duplicate payouts. The producer config says enable.idempotence=true. Transactions are on. The dashboard claimed exactly-once. And yet finance is holding a spreadsheet where 312 merchants got paid twice — $214,000 duplicated — because exactly-once is a contract with four signatories, and one of yours never signed.

Kafka's EOS story is genuinely strong but narrow: idempotent producers dedup retried sends, transactions make multi-partition writes atomic, and read_committed consumers hide aborted records. Break any one joint — in-flight requests too high, transactional.id shared across replicas, a downstream reader on read_uncommitted, a retry storm that outruns dedup windows — and duplicates flow while every individual setting looks correct.

This advanced guide traces all four joints with production failure modes. You'll learn how sequence numbers and epochs actually work, why fencing exceptions are protection (not errors to suppress), how retry storms defeat broker-side dedup, and how to verify EOS end to end with chaos tests instead of config reviews. By the end you'll treat exactly-once as a tested property, not a flag.

What Exactly-Once Promises (and What It Doesn't)

Exactly-once semantics promises each record takes effect once, even with retries, crashes, and rebalances. Kafka implements it as three cooperating mechanisms, and the promise holds only when all three are configured. Miss one and you silently downgrade to at-least-once while the config still says EOS — the most expensive kind of wrong.

The honest scope matters as much as the machinery. Kafka's EOS covers the produce path (broker dedups retried sends) and the consume-transform-produce loop (transactions make writes plus offset commits atomic). It does not cover your downstream database: if your handler applies a record twice because it replays, Kafka can't retract the second write. True end-to-end EOS always pairs broker guarantees with idempotent handlers — dedup keys, upserts, conditional writes.

Treat the flag as the start of verification, not the end. enable.idempotence=true with transactions enabled passes every config audit and still duplicates under shared transactional.ids, read_uncommitted readers, or retry storms past dedup windows. The sections below dissect each joint with its production failure mode, so you can check the machinery instead of trusting the label. Bookmark this scope before any incident: when duplicates appear, you will know exactly which of the four joints to interrogate first.

📊 Production Insight
A team passed three config audits with EOS 'enabled' while duplicating 4,800 payouts. Rule: EOS is a tested property (chaos-test it), not a flag to set.
🎯 Key Takeaway
Broker EOS covers sends and atomic loops — your database still needs idempotent handlers, and flags alone prove nothing.

Idempotent Producers: enable.idempotence Deep Dive

The idempotent producer attaches a sequence number to every batch per partition (PID plus epoch plus sequence). The broker tracks the highest sequence seen and drops any resend with a sequence it already applied — turning at-least-once retries into exactly-once writes for that partition. This survives transient network errors transparently: the client retries, the broker dedups, your app never knows.

Two boundaries limit the magic. First, ordering: dedup assumes batches arrive in sequence order, which holds only when max.in.flight.requests.per.connection stays at 5 or below (the protocol's safe window). Raise it for throughput and retried batches interleave out of order — sequence gaps the broker can't reconcile, duplicates it can't catch. Second, session scope: PID state resets on producer restart, so pre-restart retries dedup but post-restart replays don't — that's the gap transactions and epochs close.

Keep the baseline boring: enable.idempotence=true, acks=all, retries at max, in-flight at 5, delivery timeout bounded at 120s. Verify the live config with grep, not the wiki — one 'performance tuning' commit raising in-flight to 10 silently voids the guarantee. Idempotence is the foundation; the next sections build fencing and atomicity on top of it.

producer.propertiesPROPERTIES
1
2
3
4
5
6
7
8
9
10
# producer.properties — the safe EOS baseline
enable.idempotence=true
acks=all
retries=2147483647
max.in.flight.requests.per.connection=5
delivery.timeout.ms=120000
retry.backoff.ms=100

# verify what is actually live (not what the wiki says)
grep -E 'enable.idempotence|max.in.flight|retries|acks|delivery.timeout' /etc/app/producer.properties
⚠ In-Flight Over 5 Quietly Disables Your Dedup
max.in.flight.requests.per.connection above 5 with idempotence enabled breaks the ordering the broker's sequence dedup depends on. Retries then interleave out of order and duplicates slip through — with every EOS flag proudly set.
📊 Production Insight
One tuning commit raised in-flight to 10 and voided dedup for months. Rule: grep the live producer config; wikis lie, properties files don't.
🎯 Key Takeaway
Sequence-per-partition dedups retried sends within a session — keep in-flight at 5 or below and verify live config, not docs.

Transactions: transactional.id Fencing Explained

Transactions extend dedup across restarts and across partitions. Each transactional.id gets an epoch from the broker, bumped whenever a new producer instance claims the ID. The broker rejects writes carrying stale epochs with ProducerFencedException — which means a zombie instance (network-partitioned old primary, resurrected standby) can't silently duplicate the new primary's writes. Fencing is split-brain protection, and fencing exceptions are it working.

The identity discipline is absolute: one live writer per transactional.id, and the ID must be stable across restarts of the same logical writer (statefulset ordinal, stable pod name) but unique across distinct writers. Share one ID between primary and hot standby and they fence each other forever — abort rates spike past 30% while both instances believe they're the victim. Generate IDs from hostname or ordinal and the problem class vanishes.

Structure every loop as begin, produce, sendOffsetsToTransaction, commit — with abort on any exception. The offset commit inside the transaction is what makes consume-transform-produce atomic: either the outputs and the progress marker land together, or neither does. Keep transaction.timeout.ms above your slowest batch so long-but-healthy transactions aren't reaped mid-work, and alert on ProducerFencedException rate — a trickle means failovers, a flood means shared IDs.

txn_producer.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
import socket
from kafka import KafkaProducer

# Unique stable ID per instance: fencing needs 1:1 identity
TXN_ID = f"payouts-{socket.gethostname()}"  # or statefulset ordinal
producer = KafkaProducer(
    bootstrap_servers="kafka:9092",
    transactional_id=TXN_ID,
    enable_idempotence=True,
    retries=2147483647,
    acks="all",
)
producer.init_transactions()

try:
    producer.begin_transaction()
    producer.send("ledger", key=k, value=v)
    producer.send_offsets_to_transaction(offsets, group_id="payouts-group")
    producer.commit_transaction()  # atomic: writes + offsets or nothing
 except Exception:
    producer.abort_transaction()
    raise
📊 Production Insight
Shared transactional.id across primary/standby drove 31% abort rates. Rule: fencing exceptions are protection — a flood means your IDs are shared.
🎯 Key Takeaway
One stable unique ID per writer turns fencing into split-brain protection — share IDs and writers destroy each other.

read_committed: the Consumer Half Everyone Forgets

Transactions write abort markers for rolled-back work, and consumers choose whether to respect them. read_uncommitted (the default) delivers everything including aborted records — fast, but it exposes ghosts: records visible now, gone on re-read, reconciled as real by anyone downstream. read_committed filters to committed data only, holding back uncommitted reads until the transaction resolves.

The failure mode is organizational, not technical. Your team sets read_committed on the ledger consumer; the analytics team plugs the default console consumer into the same topic for a dashboard; finance reconciles ghosts against the ledger and pages you for 'lost' records that were correctly aborted. One default-config reader anywhere downstream re-exposes every abort the pipeline paid to hide.

Audit readers the way you audit producers: grep isolation.level across every properties file, Connect worker config, and stream topology touching transactional topics. Encode read_committed in shared consumer templates so new services inherit it, and add it to onboarding checklists for topic access — 'which isolation level and why' should gate every new subscription the way schema review gates new producers.

consumer.propertiesPROPERTIES
1
2
3
4
5
6
7
8
# Every downstream reader of transactional topics needs this
isolation.level=read_committed

# Audit ALL readers (apps, Connect sinks, stream processors)
grep -rnE 'isolation.level' /etc/app/*.properties /etc/kafka/connect-*.properties

# Expected: every line says read_committed.
# Any line missing or read_uncommitted = ghost reads live there.
📊 Production Insight
An analytics dashboard on defaults reconciled 9,200 ghosts as real money. Rule: isolation.level is a pipeline-wide property, not a per-team preference.
🎯 Key Takeaway
One read_uncommitted reader anywhere downstream resurrects every aborted ghost — audit all readers and template the default.

How Retry Storms Break EOS (and How to Bound Them)

Retry storms break EOS at the buffer, not the protocol. With unbounded retries, a 2-minute broker blip queues tens of minutes of unsent records in producer memory. When connectivity returns, the flood replays — exceeding the broker's sequence tracking window, arriving out of order, and duplicating everything the pre-blip instance already sent through a fenced predecessor. The broker dedups what fits its window; the storm doesn't fit.

Bound the blast radius with delivery.timeout.ms (120s is sane for most pipelines) so doomed sends fail into your error handling instead of replaying forever. Add retry.backoff.ms=100 so retries don't hammer a recovering broker into a second outage. Then accept the residual risk structurally: make downstream handlers idempotent on the business key (payout_id upserts, ON CONFLICT DO NOTHING), so any duplicate that escapes the broker becomes a no-op at the sink.

This layered posture — broker dedup first, bounded retries second, idempotent sinks third — is what 'exactly-once' means in production. No single layer holds under all failures; the three together hold under every failure you've tested. And testing is the operative word, which the final section covers.

retry_bounds.propertiesPROPERTIES
1
2
3
4
5
6
7
8
9
10
11
# Bounded retries: fail fast instead of replaying forever
delivery.timeout.ms=120000
retry.backoff.ms=100
request.timeout.ms=30000

# Correlate blips with later duplicate bursts
# grep -E 'NotEnoughReplicas|NetworkException|Disconnect' /var/log/app/producer.log
# duplicates starting 10-40 min after a 2-min blip = unbounded buffer replay

# Downstream last defense: idempotent apply on business key
# INSERT ... ON CONFLICT (payout_id) DO NOTHING  -- dupes become no-ops
📊 Production Insight
A 2-minute blip with infinite retries replayed 40 minutes of payouts past dedup windows. Rule: delivery.timeout.ms=120s turns storms into handled errors.
🎯 Key Takeaway
Bound delivery time plus backoff, then make sinks idempotent — broker dedup alone can't absorb unbounded replays.

Verifying EOS End to End With Chaos Tests

Config reviews can't prove EOS — only adversarial tests can. Build a nightly gate that kills producers mid-transaction (kill -9, no flush), partitions networks during commits, and restarts brokers mid-batch, then counts duplicates downstream with a DISTINCT check on the business key. Zero means the guarantee holds; anything else names the leaking joint with a number attached.

Cover the full matrix over time: kill -9 during begin, during send, during commit; broker bounce during commit; standby promotion while primary is partitioned (fencing check); consumer rebalance mid-transaction (offset atomicity check). Each scenario maps to one joint, so failures diagnose themselves. Run the suite after every producer config change — 'performance tuning' commits are the leading cause of silently voided EOS.

Track three rates as leading indicators between chaos runs: transaction abort rate (spikes mean fencing fights or timeouts), ProducerFencedException rate (trickle means failovers, flood means shared IDs), and downstream duplicate rate (any nonzero means a joint is already leaking). Dashboards on these three turn EOS from a launch-day claim into a continuously verified property — which is the only kind worth putting in front of finance.

chaos_eos.shBASH
1
2
3
4
5
6
7
8
9
10
# Nightly EOS gate: kill -9 mid-batch, count dupes downstream
#!/bin/bash
# chaos_eos.sh — fails the build on any duplicate
TXN_PRODUCER_PID=$(pgrep -f txn_producer.py)
kill -9 $TXN_PRODUCER_PID   # no flush, no close — worst case
sleep 60                     # let failover + fencing settle
DUPES=$(psql -tAc "SELECT count(*) - count(DISTINCT payout_id) FROM ledger")
echo "duplicate_payouts=$DUPES"
[ "$DUPES" -eq 0 ] || { echo "EOS GATE FAILED"; exit 1; }
# Also assert: abort rate sane, fencing rate sane, ghost reads zero
📊 Production Insight
Nightly chaos with a zero-duplicate gate caught a shared-ID regression in 1 day instead of at month-end close. Rule: EOS without a chaos gate is a hope.
🎯 Key Takeaway
Nightly kill -9 plus DISTINCT-count gates prove EOS continuously — track aborts, fencing, and duplicate rates between runs.
● Production incidentPOST-MORTEMseverity: high

The $214,000 Duplicate Payout Behind enable.idempotence=true

Symptom
Monday reconciliation found 4,800 duplicate payout events over 7 days — 312 merchants paid twice, $214,000 duplicated. Producer logs showed ProducerFencedException bursts (1,400/hour) interleaved with 31% transaction abort rates. The analytics topic (read_uncommitted) displayed 9,200 ghost records that vanished on re-read, and one 2-minute broker blip on Thursday preceded a 41-minute duplicate burst at 3,100 duplicates/min.
Assumption
The team believed enable.idempotence=true plus transactions equals exactly-once, full stop. They shared one transactional.id across primary and hot standby 'for seamless failover,' left the analytics consumer on default read_uncommitted, and set retries to infinite 'so nothing is ever lost.' Each choice felt safe alone; together they built a duplicate machine that config reviews approved three times.
Root cause
Three joints failed at once. Primary and standby shared transactional.id=payouts-1, so each failover probe fenced the live writer — abort rates hit 31% and both instances rewrote the same payouts. Retries were unbounded (delivery.timeout.ms infinite), so a 2-minute broker blip queued 40 minutes of buffered sends that replayed past the broker's sequence dedup window. The analytics consumer ran read_uncommitted, surfacing every aborted record as a ghost the ledger team reconciled as real.
Fix
They gave each instance a unique transactional.id (pod ordinal), set delivery.timeout.ms=120000 with retry.backoff.ms=100, flipped every downstream consumer to read_committed, and made the ledger writer idempotent on payout-id. A kill -9 chaos suite now runs nightly with a zero-duplicate gate. Duplicates fell from 4,800/week to zero over 11 days, and the $214,000 was recovered via merchant credits within the month.
Key lesson
  • A shared transactional.id turns failover into mutual fencing — unique stable IDs are what make fencing protect instead of destroy.
  • One read_uncommitted consumer anywhere downstream re-exposes every abort you paid transactions to hide; audit all readers, not just yours.
  • Infinite retries plus finite dedup windows guarantee duplicate storms after blips — bound delivery time and keep handlers idempotent regardless.
Production debug guideFive checks — idempotence, fencing, isolation, retry bounds, atomic commits — that find the leaking joint.5 entries
Symptom · 01
Duplicates appear after transient broker errors
→
Fix
Verify idempotence is real, not just declared: grep -E 'enable.idempotence|max.in.flight.requests.per.connection|retries|acks' /etc/app/producer.properties. You need enable.idempotence=true with max.in.flight ≤ 5 and retries enabled. If in-flight exceeds 5, ordering breaks under retries and the broker's per-partition sequence dedup can't save you — lower it first, then re-run the duplicate repro.
Symptom · 02
Fencing exceptions plus duplicate bursts together
→
Fix
Hunt zombie writers: grep -cE 'ProducerFencedException|InvalidProducerEpoch' /var/log/app/producer.log over the last hour, then grep -E 'transactional.id' /etc/app/producer.properties across every replica. Shared IDs across live plus standby instances mean mutual fencing — abort rates spike and both sides duplicate. Assign unique stable IDs (statefulset ordinal or pod name) so only one writer ever owns an ID.
Symptom · 03
Records visible to consumers that later vanish
→
Fix
Audit every downstream reader: grep -rnE 'isolation.level' /etc/app/.properties /etc/kafka/connect-.properties. Any consumer on transactional topics showing read_uncommitted (the default) exposes aborted records as ghosts — visible, then gone. Set isolation.level=read_committed everywhere downstream, including Kafka Connect sinks and stream processors, then confirm ghost reads stop.
Symptom · 04
Duplicate storms minutes after a short broker blip
→
Fix
Bound the blast radius: check delivery.timeout.ms (infinite by default in older clients) and retry.backoff.ms in producer configs, then correlate duplicate bursts with broker blips via grep -E 'NotEnoughReplicas|NetworkException|Disconnect' /var/log/app/producer.log. Blip followed 10-40 minutes later by duplicates means unbounded buffers replayed past dedup windows — set delivery.timeout.ms=120000 with retry.backoff.ms=100 and re-test with a chaos blip.
Symptom · 05
Need proof EOS actually holds end to end
→
Fix
Prove atomicity of the whole loop: confirm the consumer commits via sendOffsetsToTransaction (grep -rn 'sendOffsetsToTransaction' src/) and that transaction.timeout.ms exceeds your slowest batch. Then chaos-test: kill -9 a producer mid-batch and count duplicates downstream — zero means EOS holds, anything above zero names the leaking joint. Repeat after every config change; EOS is a tested property, not a reviewed one.
Broken Exactly-Once: Diagnose at a Glance
Root CauseHow to ConfirmFixPrevention
Idempotence disabled or misconfiguredDuplicates after retries; producer log shows no sequence tracking; in-flight > 5enable.idempotence=true, in-flight ≤ 5, retries onConfluent/lint check on producer configs in CI
transactional.id fencing misfireProducerFencedException storms; two instances share one ID; abort rate spikesUnique stable transactional.id per instance; one live writerStatefulSet ordinals; alert on fencing exceptions
Downstream reading uncommitted dataGhost records visible then gone; consumer isolation.level is read_uncommittedSet read_committed on every downstream readerAudit all consumers/connectors for isolation level
Retry storm beyond dedup windowDuplicate bursts after broker blips; delivery.timeout.ms infinite; huge buffersBounded retries with backoff; idempotent downstream handlersChaos-test broker blips; measure duplicate rate
⚙ Quick Reference
5 commands from this guide
FileCommand / CodePurpose
producer.propertiesenable.idempotence=trueIdempotent Producers
txn_producer.pyfrom kafka import KafkaProducerTransactions
consumer.propertiesisolation.level=read_committedread_committed
retry_bounds.propertiesdelivery.timeout.ms=120000How Retry Storms Break EOS (and How to Bound Them)
chaos_eos.shTXN_PRODUCER_PID=$(pgrep -f txn_producer.py)Verifying EOS End to End With Chaos Tests

Key takeaways

1
EOS needs all three
idempotent producer, transactions with fenced IDs, read_committed consumers — any gap duplicates.
2
Keep max.in.flight at 5 or below with idempotence, and give every instance a unique stable transactional.id.
3
Fencing exceptions protect against split-brain
alert on them, never swallow them.
4
Bound retries with delivery.timeout.ms plus backoff; infinite retries outrun broker dedup windows.
5
Commit offsets inside the transaction, and keep downstream handlers idempotent as the last defense.
6
Verify with kill -9 chaos tests measuring duplicate rate
config reviews can't prove EOS.

Common mistakes to avoid

5 patterns
×

Raising max.in.flight.requests past 5 with idempotence on

Symptom
Ordering breaks under retries, sequence gaps appear, and the broker's dedup window can't protect you — duplicates with idempotence proudly enabled.
Fix
Set enable.idempotence=true and leave max.in.flight.requests.per.connection at 5 or below with retries on. The defaults are safe together — custom tuning is what breaks them.
×

Sharing one transactional.id across standby replicas

Symptom
Live and standby fence each other in a loop, abort rates hit 30%, and both instances spend the day killing each other's transactions.
Fix
Give every producer instance a unique transactional.id (statefulset ordinal, pod name) so fencing works. Never share one ID across replicas or restarts.
×

Forgetting read_committed on downstream consumers

Symptom
Dashboards show records that later vanish after aborts. Finance reconciles against ghosts, and the EOS pipeline gets blamed for 'losing' data it correctly rolled back.
Fix
Set isolation.level=read_committed on every downstream consumer and audit third-party connectors. One read_uncommitted reader re-exposes the ghosts you paid transactions to bury.
×

Retrying forever with no delivery timeout

Symptom
The broker recovers in 2 minutes but the producer replays 40 minutes of buffered records — downstream dedup drowns and lag spikes to millions.
Fix
Bound retries with delivery.timeout.ms (e.g. 120s) and retry.backoff.ms=100 so doomed sends fail fast. Infinite retries turn a 2-minute broker blip into a 40-minute duplicate storm.
×

Committing consumer offsets outside the transaction

Symptom
Rebalances replay already-published records because the offset commit didn't share the transaction's fate — EOS on produce, at-least-once on consume.
Fix
Commit offsets inside the transaction (sendOffsetsToTransaction) and keep the consume-transform-produce loop atomic. The offset commit is part of the transaction, not an afterthought.
INTERVIEW PREP · PRACTICE MODE

Interview Questions on This Topic

Q01JUNIOR
What are the three pieces of Kafka exactly-once?
Q02SENIOR
How does the idempotent producer actually dedup?
Q03SENIOR
How does transactional.id fencing stop zombie writers?
Q04SENIOR
Why do retry storms defeat idempotence?
Q05SENIOR
Design EOS consume-transform-produce that survives kill -9 mid-batch.
Q01 of 05JUNIOR

What are the three pieces of Kafka exactly-once?

ANSWER
Exactly-once means each record takes effect once despite retries and failures. Kafka builds it from idempotent producers (dedup retried sends per partition) plus transactions (atomic multi-partition writes with offset commits) plus read_committed consumers (hide aborted data). All three or it isn't EOS.
FAQ · 6 QUESTIONS

Frequently Asked Questions

01
Is enable.idempotence=true enough for exactly-once?
02
What is producer fencing in one paragraph?
03
read_committed vs read_uncommitted — what's the real difference?
04
Can retries break exactly-once even with transactions?
05
What happens when a transactional producer crashes mid-commit?
06
Do transactions hurt throughput badly?
N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Drawn from code that ran under real load.

Follow
✓ Verified
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
🔥

That's Kafka. Mark it forged?

5 min read · try the examples if you haven't

←
Previous
Kafka RecordTooLargeException
4 / 4 · Kafka
Next
Airflow Task Stuck in Queued Forever
→