Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?
answer
- at-least-once => possible duplicates
- crash between process and offset commit
- key: business id OR (topic,partition,offset)
- UNIQUE constraint / Redis SETNX
- TTL to bound growth
basics
~20 sKafka 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 sKafka'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
Know that Kafka can deliver the same message twice and you keep a set of seen ids to skip repeats.
Explain at-least-once cause (crash between process and commit) and pick a key plus a UNIQUE/SETNX store with TTL.
Contrast business key vs (t,p,o), insist on atomic check-and-set, and reason about TTL sizing vs duplicate window.
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.