skip to content

When a Kafka consumer or sink connector writes to an external system, why is at-least-once the practical floor, and how do you make the sink effectively exactly-once?

level: middleimportance: must knowfreq 65%

answer

  1. two systems, no shared txn
  2. apply-then-commit = duplicates on crash
  3. upsert by deterministic key
  4. store offset IN the sink store
  5. dedupe = delivery≠effect

basics

~20 s

Because the Kafka offset commit and the external write happen in two separate systems, a crash between them causes redelivery and a repeat write — so it's at-least-once. You reach effective exactly-once by making the sink idempotent: upsert by a deterministic key or dedupe on a stored offset/ID so reprocessing produces no duplicate effect.

solid answer

~50 s

A sink does two non-atomic things: write to the external system and commit the Kafka offset. If the process dies after writing but before committing, the record is redelivered and written again — hence at-least-once is the floor. Kafka transactions can't help because the external write isn't in the transaction. The fix is to push the dedupe responsibility into the sink so repeats are harmless: (1) idempotent upsert keyed by a deterministic business key or by (topic, partition, offset); (2) store the last processed offset transactionally in the same external store, so on recovery you skip already-applied records; or (3) a unique constraint / conditional write that rejects duplicates. Kafka Connect sinks commonly rely on idempotent upserts (e.g., JDBC sink in upsert mode with a primary key). The goal is effectively-once behavior even though delivery is at-least-once.

go deeper

for a junior

Know that sink writes can duplicate on crash and that idempotent writes (upserts) fix it.

for a middle

Explain the two-non-atomic-steps problem and at least two idempotency patterns (upsert by key, dedupe table).

for a senior

Distinguish idempotent vs non-idempotent operations and design offset-in-sink-store or dedupe-ledger solutions that survive rebalances.

for a principal

Set org-wide conventions: every external effect carries an idempotency key; mandate effective-once via store-local atomicity rather than trusting Kafka EOS.

## The two-step problem A **sink** consumes from Kafka and writes to an external system (database, search index, object store, API). It performs two distinct actions: 1. **Apply the effect** in the external system (INSERT/UPDATE/call). 2. **Commit the Kafka consumer offset** so it won't re-read that record. These live in **two different systems** with **no shared transaction**. Whatever order you do them in, a crash in between breaks exactness: - Commit offset **first**, then apply → crash loses the effect (**at-most-once**, data loss). - Apply effect **first**, then commit → crash redelivers and re-applies (**at-least-once**, duplicates). At-least-once is the sane default (never lose data), so **duplicates are the hazard** to neutralize. ## Why Kafka EOS can't rescue this Kafka transactions atomically bind Kafka writes + offset commits *within the cluster*. The external write is invisible to the transaction coordinator, so it can't be rolled back or committed atomically with the offset. ## Making the sink effectively exactly-once (idempotency patterns) **1. Idempotent upsert by deterministic key.** If each record maps to a stable key, an UPSERT (INSERT ... ON CONFLICT DO UPDATE) makes re-applying the same record a no-op or a harmless overwrite. The JDBC sink connector's `insert.mode=upsert` with `pk.mode` relies on this. Works when the operation is naturally idempotent (state replace), less so for increments. **2. Store the processed offset in the sink store transactionally.** Write the data and the `(topic, partition, offset)` high-water mark in the **same external transaction**. On restart, read that stored offset and skip anything at or below it. This effectively moves the offset commit *into* the external system, making the two steps atomic there. (This is how some exactly-once sink implementations work.) **3. Dedupe table / unique constraint.** Insert a row keyed by a unique message id (e.g., a producer-assigned UUID or `(partition, offset)`); a duplicate hits the unique constraint and is discarded. **4. Conditional / compare-and-set writes.** For stores supporting conditional puts (e.g., DynamoDB condition expressions, S3 with a content hash key), reject or no-op on the second write. ## Edge cases - **Non-idempotent operations** (e.g., "add $10") cannot be made idempotent by upsert alone — you need a dedupe ledger keyed by message id. - **Side effects with no key** (sending an email, charging a card) require an idempotency key passed to the downstream API. - **Rebalances** can replay uncommitted records; idempotency must survive partition reassignment, so key on stable identifiers, not in-memory counters. ## Bottom line Delivery stays at-least-once; **idempotency converts at-least-once delivery into exactly-once *effect*.** That distinction (delivery vs. effect) is the whole game for sinks.

  • An upsert makes 'set balance = X' idempotent, but the operation is 'add $10'. How do you dedupe that?
    Increments aren't idempotent under retry, so upsert won't help. Use a dedupe ledger: record each message's unique id (e.g., partition+offset or a producer UUID) in the same transaction as the increment, and skip the increment if that id already exists.
  • How does storing the Kafka offset inside the external store achieve effectively-once?
    By writing the data and the (topic, partition, offset) in one external transaction, the 'effect applied' and 'position advanced' facts become atomic in that store. On recovery you read the stored offset and skip already-applied records, eliminating the gap that caused duplicates.

saying these in an interview costs you the question

  • Saying enabling Kafka transactions makes a database sink exactly-once
  • Assuming every operation can be made idempotent by an upsert (increments/side effects can't)
  • Ignoring rebalances/replays when designing dedupe keys (using in-memory state)
  • Confusing exactly-once delivery with exactly-once effect

context