Build Idempotent Kafka Consumers: Patterns That Actually Work

Stéphane Derosiaux April 26, 2024 4 min read
Build Idempotent Kafka Consumers: Patterns That Actually Work

Your consumer will see the same message twice. This isn't a bug. It's how Kafka works.

A common production incident starts with a team assuming at-least-once means exactly-once. The producer retries after a timeout. The consumer crashes before committing. A rebalance triggers reprocessing. Each scenario creates duplicates, and each team learns the hard way.

The solution isn't to fight Kafka's delivery semantics. It's to make your consumer idempotent: process the same message twice, get the same result once.

Why Duplicates Happen

Three scenarios cause duplicates:

Producer retries: A producer sends a message, the broker writes it, but the acknowledgment gets lost. The producer retries, creating a duplicate.

Consumer crashes: Your consumer processes a message, calls an external API, then crashes before committing the offset. After restart, Kafka redelivers.

Rebalancing: During consumer group rebalances, partitions move between consumers. Messages processed but not committed get redelivered.

Kafka's idempotent producer (enable.idempotence=true) only prevents duplicates at the broker level. Consumer-side duplicates are your problem.

The Pattern

An idempotent consumer tracks which messages it has processed:

public void consume(ConsumerRecord<String, OrderEvent> record) {
    String key = extractIdempotencyKey(record);

    if (deduplicationStore.hasProcessed(key)) {
        log.debug("Skipping duplicate: {}", key);
        return;
    }

    processOrder(record.value());
    deduplicationStore.markProcessed(key);
}

The hard part is making this atomic. If you check, process, and mark in three steps, a crash between any of them creates inconsistency.

Designing Idempotency Keys

The key must be unique per logical operation, stable across retries, and ideally generated by the producer.

Option 1: Client-generated UUID. The safest approach. Generate a UUID at the API layer that flows through the system.

Option 2: Composite business key. Derive from business attributes: customerId:orderId:timestamp. Works when the combination is truly unique. Use the business event time, not the record's produce timestamp, which changes when the producer resends.

Option 3: Kafka coordinates. Use topic-partition-offset as a key. Simple but breaks if you replay from a different topic, if the producer resends (the copy gets a new offset), or if the data moves to another cluster.

Pattern 1: Database Constraint

The simplest pattern. Use a unique constraint to reject duplicates.

CREATE TABLE processed_messages (
    idempotency_key VARCHAR(255) PRIMARY KEY,
    processed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

In your consumer:

@Transactional
public void processOrder(ConsumerRecord<String, OrderEvent> record) {
    String key = extractIdempotencyKey(record);

    int inserted = jdbcTemplate.update(
        "INSERT INTO processed_messages (idempotency_key) VALUES (?) ON CONFLICT DO NOTHING",
        key);
    if (inserted == 0) {
        log.info("Duplicate detected: {}", key);
        return;
    }
    orderService.createOrder(record.value());  // same DataSource, same transaction
}

The @Transactional wrapper ensures atomicity. If the order creation fails, the deduplication record also rolls back. ON CONFLICT DO NOTHING returns 0 rows on a duplicate instead of throwing. Catching a unique-violation exception inside the transaction doesn't work: in PostgreSQL the failed statement aborts the whole transaction, and with JPA the insert may only flush at commit, outside your try.

Concurrency note: The primary key does the work. When two consumers insert the same key, the second blocks on the unique index until the first commits, then inserts nothing. Don't add a SELECT before the INSERT: that check-then-insert is the race. Under REPEATABLE READ or SERIALIZABLE, the second transaction fails with a serialization error instead and must be retried.

Tradeoff: The deduplication table becomes a write hotspot under high load. For high-throughput services, consider Redis, with the caveat below.

Pattern 2: Redis SETNX

For higher throughput, use Redis with SETNX:

public boolean tryProcess(String idempotencyKey) {
    Boolean wasSet = redis.opsForValue()
        .setIfAbsent("dedup:" + idempotencyKey, "1", Duration.ofDays(7));
    return Boolean.TRUE.equals(wasSet);
}

The TTL handles cleanup automatically. Set it longer than your maximum expected replay window.

Atomicity: Redis and your database are two separate writes. Set the key before the business write and a crash in between skips the message for good; set it after and a crash in between processes it twice. Use this pattern only when the side effect is itself idempotent (an upsert, or an API call that carries the idempotency key), with Redis as a fast first filter. When the side effect is a plain insert, keep the key in the database transaction (Pattern 1).

Redis persistence: Configure Redis with AOF (appendonly yes, appendfsync everysec) to survive restarts. With RDB-only snapshots, a Redis crash between snapshots loses recent deduplication state, causing duplicates. For critical workloads, use Redis Cluster with replication. Redis replication is asynchronous, so a failover can still lose the most recent keys.

Tradeoff: Redis adds a network hop and infrastructure dependency. If Redis is unavailable, your consumer stops. Consider whether this availability tradeoff works for your use case.

Pattern 3: Natural Idempotency

When your business logic supports upserts, you might not need a separate deduplication store:

INSERT INTO inventory (product_id, quantity, event_version)
VALUES (?, ?, ?)
ON CONFLICT (product_id)
DO UPDATE SET quantity = EXCLUDED.quantity, event_version = EXCLUDED.event_version
WHERE inventory.event_version < EXCLUDED.event_version;

The WHERE clause prevents out-of-order events from overwriting newer data.

This works for state replacement events. It doesn't work for delta operations ("add 10 to quantity") or side effects (sending emails).

Choosing a Pattern

PatternThroughputBest For
Database constraintLow-MediumCRUD services, existing DB
Redis SETNXHighSide effects that are themselves idempotent
Natural upsertMedium-HighState replacement events
For most services that read from Kafka and write to a database, the database constraint pattern is sufficient and simpler than full exactly-once transactions.

The Deduplication Window

Your window must be longer than your maximum expected delay. Consider:

  • How long might your consumer be down?
  • Will you ever replay from hours or days ago?
  • What's your producer's retry behavior?

Seven days is common. Balance storage costs against your replay requirements.

Testing Idempotency

Always verify your consumer handles duplicates:

@Test
void shouldHandleDuplicateMessages() {
    String key = UUID.randomUUID().toString();
    ConsumerRecord<String, OrderEvent> record = createRecord(key);

    processor.process(record);
    processor.process(record);

    assertThat(orderRepository.findAll()).hasSize(1);
}

Test concurrent duplicates too. Race conditions hide until production.

At-Least-Once + Idempotent = Effectively Exactly-Once

The combination of at-least-once delivery and idempotent consumers gives you effectively exactly-once processing:

At-least-once + Idempotent consumer = Effectively exactly-once

This is simpler than Kafka's transactional exactly-once, which requires transactional producers, read_committed consumers, and two-phase commit coordination.

For most services, idempotent consumers are sufficient. Build your consumers to handle duplicates, and you'll stop worrying about Kafka's delivery semantics.

Book a demo to see how Conduktor Console helps you trace message flow and identify deduplication issues across your Kafka clusters.