# Idempotent Consumer Pattern in Kafka Explained

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](https://www.conduktor.io/kafka/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 "delivered": client buffer, leader ack, ISR ack, consumer offset commit](https://www.conduktor.io/assets/images/glossary/idempotent-consumer-pattern-0.webp)

- **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](https://www.conduktor.io/glossary/kafka-consumer-groups-explained), 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](https://www.conduktor.io/glossary/dead-letter-queues-for-error-handling) and later replayed arrive again, often hours later.
- **Cross-cluster failover.** After a failover to a [MirrorMaker 2](https://www.conduktor.io/glossary/kafka-mirrormaker-2-for-cross-cluster-replication) 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.

| Key | Stable across producer retries | Stable across replays and cluster moves | Notes |
|---|---|---|---|
| Producer-supplied UUID in the record (header or payload) | Yes | Yes | Generated once at the origin, travels with the record |
| Business key (`payment_id`, `order_id` plus version) | Yes | Yes | Same properties when the domain already has a unique id |
| `(topic, partition, offset)` | No: a retried write gets a new offset | No: offsets change across clusters and replays | Identifies 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](https://www.conduktor.io/assets/images/glossary/idempotent-consumer-pattern-1.webp)

```sql
-- 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](https://www.conduktor.io/glossary/kafka-transactions-deep-dive) 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](https://www.conduktor.io/blog/kafka-exactly-once-performance-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](https://www.conduktor.io/blog/exactly-once-semantics-when-it-works)).

| Output of the consumer | Mechanism |
|---|---|
| Kafka topic only (stream processor, enrichment, routing) | Kafka transactions, `read_committed` downstream |
| Relational database | Inbox table in the business transaction, offset commit after |
| Key-value store or cache with atomic conditional writes | Compare-and-set on the key, or an upsert keyed by the idempotency key |
| Email, payment, HTTP call, anything that cannot roll back | Idempotency 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](https://www.conduktor.io/blog/building-idempotent-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](https://www.conduktor.io/kafka/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.

## Related Pages

- [Idempotent Kafka Producer](https://www.conduktor.io/kafka/idempotent-kafka-producer): the producer-side mechanism (PID and sequence numbers) that this pattern complements.
- [Delivery Semantics for Kafka Consumers](https://www.conduktor.io/kafka/delivery-semantics-for-kafka-consumers): at-most-once, at-least-once, and the offset commit strategies behind them.
- [Exactly-Once Semantics in Kafka](https://www.conduktor.io/glossary/exactly-once-semantics-in-kafka): what Kafka transactions guarantee and where the guarantee stops.
- [Kafka Transactions Deep Dive](https://www.conduktor.io/glossary/kafka-transactions-deep-dive): the consume-process-produce loop that replaces dedup for Kafka-to-Kafka flows.
- [When Kafka Exactly-Once Is Worth the Performance Cost](https://www.conduktor.io/blog/kafka-exactly-once-performance-cost): the decision between transactions and idempotent consumers.
- [Kafka Exactly-Once: When It Works and When It Doesn't](https://www.conduktor.io/blog/exactly-once-semantics-when-it-works): sink connectors, side effects, and other places the guarantee stops.
- [Outbox Pattern for Reliable Event Publishing](https://www.conduktor.io/glossary/outbox-pattern-for-reliable-event-publishing): the producer-side twin of the inbox table, solving the dual-write problem.
- [Dead Letter Queues for Error Handling](https://www.conduktor.io/glossary/dead-letter-queues-for-error-handling): the replay path that sets the lower bound on the dedup window.
- [Kafka Consumer Groups Explained](https://www.conduktor.io/glossary/kafka-consumer-groups-explained): why a rebalance redelivers uncommitted records.

## Sources and References

- [Apache Kafka Design: Message Delivery Semantics](https://kafka.apache.org/43/design/design/)
- [Confluent Docs: Message Delivery Guarantees](https://docs.confluent.io/kafka/design/delivery-semantics.html)
- [KIP-98: Exactly Once Delivery and Transactional Messaging](https://cwiki.apache.org/confluence/display/KAFKA/KIP-98+-+Exactly+Once+Delivery+and+Transactional+Messaging)
- [KIP-447: Producer Scalability for Exactly Once Semantics](https://cwiki.apache.org/confluence/display/KAFKA/KIP-447%3A+Producer+scalability+for+exactly+once+semantics)
- [Apache Kafka: Geo-Replication (MirrorMaker 2)](https://kafka.apache.org/43/operations/geo-replication-cross-cluster-data-mirroring/)
- [PostgreSQL Documentation: INSERT and ON CONFLICT](https://www.postgresql.org/docs/current/sql-insert.html)
- [Apache Kafka: Consumer Configs (`isolation.level`)](https://kafka.apache.org/43/configuration/consumer-configs/)
- [Apache Kafka: Producer Configs (`transaction.timeout.ms`)](https://kafka.apache.org/43/configuration/producer-configs/)
- [Kafka Streams Core Concepts: Processing Guarantees](https://kafka.apache.org/43/streams/core-concepts/)
- [Confluent: Exactly-Once Semantics Are Possible (2017 benchmark)](https://www.confluent.io/blog/exactly-once-semantics-are-possible-heres-how-apache-kafka-does-it/)
- [Enterprise Integration Patterns: Idempotent Receiver](https://www.enterpriseintegrationpatterns.com/patterns/messaging/IdempotentReceiver.html)
- [microservices.io: Pattern, Idempotent Consumer](https://microservices.io/patterns/communication-style/idempotent-consumer.html)
- [Stripe API Reference: Idempotent Requests](https://docs.stripe.com/api/idempotent_requests)
