Home / Writing / Kafka

Why Kafka consumers process messages twice, and how to stop it

Duplicates after a deploy are not a Kafka bug. They are the default delivery guarantee meeting a side effect that isn't idempotent. Here is where they come from and the fixes that hold up.

Short answer

A Kafka consumer processes a message twice when it finishes the work but its offset commit doesn't happen before the partition moves to another consumer or the process restarts. The new owner resumes from the last committed offset and repeats the work. This is Kafka's at-least-once delivery working as designed. The durable fix is to make processing idempotent, usually with a deduplication key written in the same database transaction as the side effect. Committing on partition revocation, cooperative rebalancing and static membership reduce how often it happens, but they don't remove it.

The guarantee you actually have

A consumer group tracks progress with a committed offset per partition: the position of the next record to read. If a consumer processes records and then dies, or loses the partition, before committing, whoever takes over starts from the old committed offset. Those records get delivered again.

That is at-least-once delivery. You can trade it for at-most-once by committing before processing, but then a crash loses the records instead. For anything involving money, messages to users or counters, losing data is worse than repeating it, so most systems sit on at-least-once whether they chose it or not.

Timeline of a duplicate caused by a rebalance Consumer c-2 processes offset 88412, then is removed by a rebalance before committing. Consumer c-1 is assigned the partition, resumes from committed offset 88412 and processes it again. consumer c-2 consumer c-1 process 88412 rebalance commit never sent partition 3 reassigned · committed offset = 88412 process 88412 same record, side effect runs again commit 88413
The window between finishing the work and committing the offset is where duplicates come from.

The four ways duplicates happen

1. A rebalance lands between processing and commit

Rolling deploys, autoscaling and a consumer that misses its heartbeats all trigger rebalances. With the classic eager protocol, every member gives up its partitions. Anything processed since the last commit is replayed by the next owner. This is the most common source of "duplicates after every deploy".

2. Processing takes longer than max.poll.interval.ms

If the gap between two poll() calls exceeds max.poll.interval.ms (five minutes by default), the consumer is considered stuck and is removed from the group. Its partitions are reassigned, and the slow batch it was working on gets processed again by someone else. A large max.poll.records combined with a slow downstream call is a common way to hit this, and it tends to repeat: the same slow batch keeps getting kicked out.

3. A crash after the side effect, before the commit

The process charges the card, then the pod is OOM-killed. No rebalance protocol can help here. Only idempotent processing does.

4. Producer retries without idempotence

Sometimes the duplicate is already in the topic. A producer that times out waiting for an acknowledgement and retries can write the same record twice. Idempotent producers (enable.idempotence=true, the default in modern clients) prevent this for retries within a producer session, but an application that re-sends after a restart can still create a logical duplicate with a new offset. Your consumer has to tolerate that too.

Fixes, in order of importance

Make the side effect idempotent

This is the only fix that covers all four causes. Give every message a stable key (an event ID set by the producer, or a natural key such as an order ID), and record that you applied it in the same transaction as the effect itself.

-- one row per message you have applied
CREATE TABLE processed_events (
  consumer_group text        NOT NULL,
  event_id       text        NOT NULL,
  processed_at   timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (consumer_group, event_id)
);
// inside the handler, one database transaction
await db.transaction(async (tx) => {
  const inserted = await tx.query(
    `INSERT INTO processed_events (consumer_group, event_id)
     VALUES ($1, $2) ON CONFLICT DO NOTHING`,
    ["orders-svc", event.id]
  );
  if (inserted.rowCount === 0) return;   // already applied: skip
  await tx.query(
    "INSERT INTO charges (order_id, amount) VALUES ($1, $2)",
    [event.orderId, event.amount]
  );
});
// commit the Kafka offset only after the transaction commits

When the effect is a call to another service rather than a database write, pass the same key as an idempotency key to that service. Most payment providers support this. If the downstream service can't deduplicate, write the intent to an outbox table in your transaction and have a single sender deliver it with the key.

Keep deduplication state in durable storage. An in-memory set forgets everything on restart, which is exactly when duplicates arrive.

Commit before you lose the partition

Register a rebalance listener and synchronously commit the offsets of everything you have finished when partitions are revoked (onPartitionsRevoked in the Java client, or the equivalent hook in yours). This shrinks the replay window during graceful rebalances. It does nothing for crashes.

Rebalance less, and less disruptively

  • Cooperative rebalancing. With partition.assignment.strategy=CooperativeStickyAssignor, consumers only give up the partitions that actually move, instead of everything. Kafka 4.0 also made the new consumer group protocol (KIP-848) generally available, which moves assignment to the broker and avoids the group-wide stop. Both reduce disruption; neither changes the delivery guarantee.
  • Static membership. Set a stable group.instance.id per instance and a session.timeout.ms longer than a normal restart. A pod that restarts during a rolling deploy rejoins with its old identity and its old partitions, with no rebalance at all.
  • Bounded batches. Size max.poll.records so a full batch finishes well inside max.poll.interval.ms even when downstream is slow. Put timeouts on every downstream call.

Use transactions only where they apply

Kafka's exactly-once semantics cover read-process-write loops that stay inside Kafka: consume from one topic, produce to another, and commit the consumer offsets in the same producer transaction. They don't extend to your database, your email provider or a payment API. For those, you still need idempotence.

Try it: break it, then fix it

Process a few orders, then trigger a deploy or a crash between the charge and the offset commit and see what the next consumer does. Then turn on the idempotency key and try again.

Committed offset0
Orders received0
Duplicate charges0

Process a few orders, then deploy or crash between the charge and the commit.

What about auto-commit?

With enable.auto.commit=true, the client commits the offsets returned by earlier polls on a timer during later poll() calls. In a simple synchronous loop, that is still at-least-once. If you hand records to other threads and poll again before they finish, auto-commit can commit records that were never processed, and a crash loses them. If you process asynchronously, turn auto-commit off and commit explicitly after the work is done.

Checklist

  • Every message has a stable ID set by the producer, or a natural key.
  • The dedupe record and the side effect are written in one transaction.
  • External calls carry an idempotency key, or go through an outbox.
  • Offsets are committed after the work, and on partition revocation.
  • Static membership is set for deployments that restart in place.
  • A full batch finishes well within max.poll.interval.ms.
  • A test forces a restart between the effect and the commit, and checks the effect happened once.

Frequently asked questions

Is exactly-once delivery possible with Kafka?

Kafka's exactly-once semantics cover read-process-write loops that stay inside Kafka, where consumed offsets and produced records are committed in one transaction. They do not cover side effects in databases or external APIs. For those, you need idempotent processing.

Should I turn off enable.auto.commit?

If you process records synchronously in the poll loop, auto-commit still gives at-least-once delivery. If you process records on other threads or asynchronously, turn it off and commit explicitly after the work is done, or a crash can lose records.

Doesn't the deduplication table grow forever?

Delete rows older than the longest window in which a message could be redelivered, such as your topic retention or maximum replay period, plus a margin. If you might replay a topic from the beginning, keep keys as long as that data exists, or use natural keys with upserts.

Does cooperative rebalancing stop duplicates?

It reduces them, because consumers only give up the partitions that move. It does not stop them: crashes and slow batches still cause redelivery. Only idempotent processing removes the effect of duplicates.