skip to content

Idempotent Consumers and Deduplication

Actually building a dedup store, such as a unique constraint or a Redis key on a business idempotency key, and ordering the write correctly. Interviewers want the implementation, not just the phrase 'make it idempotent'.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?

level: juniorimportance: must knowfreq 70%

answer

  1. at-least-once => possible duplicates
  2. crash between process and offset commit
  3. key: business id OR (topic,partition,offset)
  4. UNIQUE constraint / Redis SETNX
  5. TTL to bound growth

basics

~20 s

Kafka redelivers messages after rebalances or restarts (at-least-once), so the same record can arrive more than once. You record each processed message's id in a store (a DB unique constraint or Redis key) and skip any id you've already seen.

solid answer

~50 s

Kafka's default delivery is at-least-once: if a consumer crashes or rebalances after processing but before its offset commit lands, those records are redelivered, so your handler can run twice for the same message. If processing has side effects (charge a card, send an email, insert a row), duplicates cause double work. An app-level dedup store fixes this: pick a key that uniquely identifies the message — either a business idempotency key carried in the payload/header, or the (topic, partition, offset) triple — and persist it. Before processing, check whether the key already exists; if it does, skip. The simplest implementations are a database table with a UNIQUE constraint on the key (insert fails on a duplicate) or Redis SETNX (set-if-not-exists) on the key. Add a TTL or cleanup job so the store doesn't grow unbounded.

go deeper

for a junior

Know that Kafka can deliver the same message twice and you keep a set of seen ids to skip repeats.

for a middle

Explain at-least-once cause (crash between process and commit) and pick a key plus a UNIQUE/SETNX store with TTL.

for a senior

Contrast business key vs (t,p,o), insist on atomic check-and-set, and reason about TTL sizing vs duplicate window.

for a principal

Frame dedup as one option in a spectrum (idempotent ops, transactional outbox, EOS) and weigh store cost/consistency trade-offs at scale.

## The problem: at-least-once delivery Kafka consumers track progress with **offsets** — a monotonically increasing position per (topic, partition). After processing records, a consumer commits the offset so that on restart it resumes from there. The catch: **processing and committing are two separate steps**. If the consumer processes a record, performs a side effect, but then crashes (or a **rebalance** reassigns the partition) before the offset commit is durable, the next consumer reads from the last committed offset and **re-processes** the already-handled record. This is the **at-least-once** guarantee: every record is delivered *one or more* times, never lost, but possibly duplicated. Duplicates are harmless for **idempotent** operations (e.g. `SET balance = 100`) but dangerous for non-idempotent ones (incrementing a counter, charging a card, sending an email, appending a row). ## The fix: a deduplication store Maintain a record of which messages you have already processed, keyed by a **stable identifier**: - **Business idempotency key**: a value the producer puts in the message (e.g. `orderId`, `paymentRequestId`, a UUID in a header). Best when the same logical event might appear at different offsets (e.g. re-published, or arriving on multiple topics). - **(topic, partition, offset)**: the physical coordinate of the record. Unique within a topic and trivially available from `ConsumerRecord`. Good when there is no natural business key, but it does NOT dedup the same logical event published twice (those have different offsets). ### Two common implementations 1. **Relational DB UNIQUE constraint**: a table `processed_messages(key PRIMARY KEY, processed_at)`. Attempt to `INSERT` the key; a duplicate raises a unique-violation, which you catch and treat as "already processed, skip." 2. **Redis SETNX**: `SET key 1 NX` returns success only if the key did not exist. If it already exists, you skip. Redis is fast and supports per-key TTL natively. ### Housekeeping The store grows forever unless you bound it. Use a **TTL** (Redis `EX`, or a scheduled DELETE of rows older than N days in SQL). The TTL must be longer than the maximum window in which a duplicate could plausibly arrive (covering retries, consumer lag, and DLQ replays). ## Edge cases - The dedup check and the side effect must be ordered carefully (see related questions) or you can still double-process under crashes. - A too-short TTL re-admits old duplicates; a too-long TTL bloats the store. - Under high concurrency, the check-then-act must be **atomic** (the UNIQUE constraint or SETNX gives you that for free; a separate SELECT-then-INSERT does not).

  • Why is (topic, partition, offset) sometimes a worse dedup key than a business key?
    Because the same logical event re-published to Kafka gets a new offset, so (t,p,o) sees it as new and won't dedup it. A business idempotency key in the payload survives re-publishing and even cross-topic delivery.
  • What makes the UNIQUE constraint or SETNX better than SELECT-then-INSERT?
    They are atomic check-and-set operations, so two concurrent consumers can't both pass the check and both process. A separate SELECT then INSERT has a race window.

saying these in an interview costs you the question

  • Claiming Kafka itself guarantees no duplicates to consumers by default (it's at-least-once).
  • Using SELECT-then-INSERT and assuming it's race-free.
  • Forgetting the store needs a TTL/cleanup and will grow unbounded.
  • Thinking (topic,partition,offset) dedups re-published logical events.

context

open as a page

Walk through the read-process-write ordering that makes a dedup store actually safe under consumer crashes. Where exactly do you record the key, side effect, and offset?

level: seniorimportance: must knowfreq 50%

basics

~20 s

Make the side effect and the dedup-key write part of the same atomic transaction, so either both happen or neither does. Commit the Kafka offset only after that transaction succeeds. On redelivery, the existing key tells you to skip.

open as a page

Why must the dedup check-and-record be a single atomic operation, and what bug appears if you use SELECT-then-INSERT or GET-then-SET?

level: middleimportance: should knowfreq 45%

basics

~20 s

If you check 'have I seen this key?' and then separately record it, two concurrent consumers can both pass the check before either records, so both process the message. You need one atomic step: a UNIQUE-constraint INSERT or Redis SET NX that checks and records together.

open as a page

When should you key your dedup store on a business idempotency key versus (topic, partition, offset)?

level: middleimportance: should knowfreq 55%

basics

~20 s

Use a business key (like orderId) when the same logical event can reappear at a different offset or on another topic — it survives re-publishing. Use (topic,partition,offset) only when there's no natural business id and duplicates come purely from Kafka redelivery.

open as a page

How would you implement a Redis-based dedup store with SETNX, and what TTL/cleanup and durability concerns must you handle?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Use SET key value NX EX <ttl>: it sets the key only if absent (NX) and auto-expires after the TTL (EX). If the command returns 'not set', you've seen the message — skip it. TTL bounds memory; pick it longer than any plausible duplicate window.

open as a page