Idempotent Consumer Pattern in Kafka Explained

Stéphane Derosiaux October 1, 2026 12 min read

An idempotent consumer is a consumer-side design pattern in which processing the same Kafka record more than once has the same effect as processing it once. It complements the idempotent Kafka producer: producer idempotence prevents duplicate writes to Kafka, while consumer idempotence protects databases, APIs, and other downstream systems.

enable.idempotence=true does not make a consumer idempotent. It only deduplicates producer retries between one producer session and the broker, per partition. A consumer that writes to a database and crashes before committing its offset still processes the record again.

Where duplicates come from

The diagram shows the two most common sources: a producer retry after a lost acknowledgement, and a consumer crash before the offset commit.

Four layers of

  • Producer resends outside idempotence's scope. With enable.idempotence off, a lost acknowledgement followed by a retry writes the record twice. With it on (the default since Kafka 3.0), an application-level resend after delivery.timeout.ms expires, or after a producer restart, still writes a second copy, because deduplication only applies within one producer session and partition. The log then contains two copies at different offsets.
  • Consumer crash before offset commit. The side effect happened, the offset did not move. The next instance replays from the last committed offset.
  • Rebalance after a crash or eviction. When a consumer dies or exceeds max.poll.interval.ms, its partitions move to another member of the consumer group, and everything processed since the last commit on those partitions is redelivered to the new owner. A graceful rebalance normally commits first, through auto-commit or a commit in onPartitionsRevoked.
  • Dead-letter replays. Records sent to a dead letter queue and later replayed arrive again, often hours later.
  • Cross-cluster failover. After a failover to a MirrorMaker 2 replica, consumers resume from the last translated checkpoint, which lags the primary, so a window of already-processed records is consumed again.

Example: a payment consumer writes to a ledger, then crashes before committing its Kafka offset. On restart, Kafka delivers the payment again because the committed offset still points to the earlier position. Kafka cannot see that the ledger write already happened.

Choosing the idempotency key

The dedup key must be identical on every redelivery of the same message, including producer retries, dead-letter replays, and cross-cluster moves.

KeyStable across producer retriesStable across replays and cluster movesNotes
Producer-supplied UUID in the record (header or payload)YesYesGenerated once at the origin, travels with the record
Business key (payment_id, order_id plus version)YesYesSame properties when the domain already has a unique id
(topic, partition, offset)No: a retried write gets a new offsetNo: offsets change across clusters and replaysIdentifies a log position, not a message
Use an id created with the business action and carried in the record. For example, the service that accepts an HTTP request can generate a UUID, add it to the Kafka record, and pass the same value through every downstream call.

Implementing an idempotent consumer: the inbox table

Recording "seen" and performing the side effect must succeed or fail together. If they are separate, a crash can produce a duplicate (the work completed but the key was not recorded) or lose a message (the key was recorded but the work did not complete).

In a relational database, insert the dedup row and the business row in one transaction. This is the inbox table:

Inbox pattern: dedup insert and business write in one database transaction, offset commit after

-- Postgres. The primary key is what makes ON CONFLICT detect a duplicate.
CREATE TABLE processed_messages (
  message_id     text        NOT NULL,
  consumer_group text        NOT NULL,
  processed_at   timestamptz NOT NULL DEFAULT NOW(),  -- used to expire old keys
  PRIMARY KEY (message_id, consumer_group)
);

-- Per record. One statement, one transaction: the claim and the business
-- write commit together, or not at all.
BEGIN;
WITH claimed AS (
  INSERT INTO processed_messages (message_id, consumer_group)
  VALUES ('pay-1002', 'ledger-writer')
  ON CONFLICT (message_id, consumer_group) DO NOTHING
  RETURNING message_id                       -- 0 rows on a duplicate
)
INSERT INTO ledger_entries (payment_id, account_id, amount_cents)
SELECT 'pay-1002', 'acc-7781', 1350 FROM claimed;  -- nothing on a duplicate
COMMIT;
-- Then, and only then: consumer.commitSync()

Commit the database transaction first, then the Kafka offset. A crash between the two replays the record, the INSERT conflicts, zero rows reach ledger_entries, and the offset advances. A crash before the database commit rolls back both rows, and the replay does the work for real.

The primary key also stops two consumers from processing the same record at once. This happens when a consumer evicted for exceeding max.poll.interval.ms is still processing a batch while the partition's new owner processes the same records. A separate SELECT followed by an INSERT is a race: both consumers can see no row and start the work. With INSERT ... ON CONFLICT, the second transaction blocks on the unique index until the first commits, then inserts nothing (or proceeds normally if the first rolled back). Under REPEATABLE READ or SERIALIZABLE isolation it fails with a serialization error instead and must be retried.

The deduplication window

An inbox table grows by one row per processed record. To bound it, delete rows older than a fixed window, or expire them with a TTL in Redis or a Kafka Streams windowed store. The window must outlast the latest duplicate:

  • Producer retries and crash replays arrive within seconds to minutes.
  • Rebalance replays cover what was processed since the last commit (with auto-commit, at most auto.commit.interval.ms, 5 seconds by default) and arrive once the group detects the failure: after session.timeout.ms (45 seconds by default) for a crash, or max.poll.interval.ms (5 minutes by default) for a stalled consumer.
  • DLQ re-drives and cross-cluster failovers can arrive hours or days later.

A duplicate that arrives after its key was deleted is processed as new. Size the window for the slowest redelivery path, usually the retention period of the dead-letter topic. If dead-letter records can be replayed for seven days, keep deduplication keys for at least seven days.

A deliberate replay is different: it is meant to rebuild the output. Before a full replay, clear the inbox for that consumer group together with the read model. Otherwise, the retained keys will make the replay skip some records.

Idempotent consumer vs Kafka exactly-once

For a consumer whose only output is another Kafka topic, Kafka transactions replace the inbox table. The consumed offsets and produced records commit together, and a downstream consumer with isolation.level=read_committed does not see aborted writes (the default is read_uncommitted). Kafka Streams provides this with processing.guarantee=exactly_once_v2.

Use transactions when the output is Kafka only and the job needs atomic writes across several partitions, or a Kafka Streams state store that stays consistent with its output topics. To avoid duplicates alone, an idempotent consumer is enough (see when exactly-once is worth the cost). Producer throughput barely drops: Confluent's original 2017 benchmark measured 3% against acks=all for 1 KB records in 100 ms transactions. The cost shows up downstream. A read_committed consumer only reads up to the last stable offset, so a producer that leaves a transaction open stalls every reader of that partition until the transaction commits, aborts, or reaches transaction.timeout.ms (60 seconds by default).

Kafka transactions cannot include an external database or API, and that holds inside Kafka Streams too. exactly_once_v2 covers input offsets, state stores, and output topics. An email sent from foreach() is sent again when a task fails, its transaction aborts, and the records are reprocessed. The fix is to write the intent ("send receipt for order 42") to a topic and let an idempotent consumer perform it (see where Kafka exactly-once stops).

Output of the consumerMechanism
Kafka topic only (stream processor, enrichment, routing)Kafka transactions, read_committed downstream
Relational databaseInbox table in the business transaction, offset commit after
Key-value store or cache with atomic conditional writesCompare-and-set on the key, or an upsert keyed by the idempotency key
Email, payment, HTTP call, anything that cannot roll backIdempotency key forwarded to the callee (next section)
A full-state record needs no inbox. When the record carries the whole new state of an entity, an upsert keyed by the entity is idempotent, provided the write is guarded by a version (ON CONFLICT ... DO UPDATE ... WHERE t.version < EXCLUDED.version). Without the guard, a late replay of an older record rolls the row back to stale state. It stops working when the record carries a delta ("add 10 to stock") or triggers anything outside the row.

Payments, emails, HTTP calls: pass the idempotency key downstream

A payment, email, or HTTP call cannot be rolled back. If a payment call times out, the consumer does not know whether the provider created the charge.

Pass the record's idempotency key to the external API, for example in an Idempotency-Key header. The provider can then return the result of the first request instead of repeating the action. The callee keeps keys for a limited time (Stripe prunes them after 24 hours), so a dead-letter re-drive older than that executes again: the consumer's own deduplication window still applies. If the API has no idempotency support, use a compensating workflow or accept and monitor the risk of duplicates.

Operating idempotent consumers

  • Keep the dedup record where the side effect happens. A Redis key written next to a Postgres insert is two writes again, with the same crash window as before. A separate dedup store is safe only when the side effect is itself keyed, such as an upsert or an Idempotency-Key call. Removing duplicates later in a warehouse does not protect a payment, email, or database write that already ran twice.
  • Test with duplicates. Inject duplicate records and concurrent deliveries. Conduktor Gateway's chaos interceptor can duplicate a percentage of records on produce or consume for this type of test.
  • Treat offset resets as replays. Decide whether the reset should skip past effects or rebuild them, then keep or clear the inbox accordingly.

For Java implementations, see building idempotent Kafka consumers, whose Redis SETNX variant is a separate store and falls under the first rule above; for the offset commit strategies underneath, see delivery semantics for Kafka consumers.

Summary

An idempotent consumer uses a stable message key and records that key in the same transaction as its side effect. It commits the Kafka offset afterwards and keeps keys long enough to cover the slowest replay path. Kafka transactions replace it only for Kafka-to-Kafka flows that need atomic multi-partition writes or consistent Streams state. Anything that leaves Kafka, including a side effect inside a Streams topology, still needs its own deduplication.

What is an idempotent consumer in Kafka?

A consumer built so that processing the same record twice has the effect of processing it once. It records a per-message identifier and performs the side effect in the same transaction, so redeliveries from crashes, rebalances, or replays are detected and skipped rather than executed again.

Does enable.idempotence=true make my consumer idempotent?

No. That setting deduplicates producer retries between one producer session and the broker, per partition. It does nothing for a consumer that crashes after writing to a database and before committing its offset, or for records replayed from a dead letter topic. Consumer idempotency is application code.

What should I use as the idempotency key?

A producer-supplied identifier that travels with the record: a UUID generated at the origin of the business action, or a business key like a payment id. Avoid topic, partition, and offset as the key, because a retried write gets a new offset and offsets change across clusters and dead letter replays.

Why must the dedup record and the side effect be in one transaction?

Because a crash between them either duplicates the work (side effect done, key not recorded) or loses the message (key recorded, side effect not done). One transaction makes both changes succeed or fail together.

Does Kafka exactly-once semantics replace the idempotent consumer pattern?

Only when the consumer's output is another Kafka topic, and the job needs atomic multi-partition writes or Kafka Streams state that stays consistent with its output. For a database, an email, or an HTTP call, Kafka transactions cannot roll back the external effect. That includes side effects inside a Streams topology: an email sent from foreach() under exactly_once_v2 is sent again after a task failure.

How long should the deduplication window be?

Longer than the slowest redelivery path in the system. Producer retries and crash replays arrive within minutes, but a dead letter re-drive or a cross-cluster failover can arrive days later. A duplicate that arrives after its key was evicted is processed as new, so the window is usually sized against the DLQ topic's retention.

Sources and References