skip to content

Why can't ordering of commit and process alone achieve exactly-once, and how does Kafka's transactional read-process-write make the consume-transform-produce cycle effectively exactly-once?

level: seniorimportance: must knowfreq 55%

answer

  1. two writes can't be sequenced safe → need atomicity
  2. sendOffsetsToTransaction binds offsets into the txn
  3. transactional.id + epoch fencing of zombies
  4. read_committed + commit markers + LSO
  5. exactly_once_v2 in Streams; only within Kafka

basics

~20 s

Ordering always leaves a crash window that either loses or duplicates a record. Exactly-once closes it by making the produced output AND the consumed input offset commit part of one atomic transaction — if it aborts, neither is visible, so no duplicate and no loss.

solid answer

~50 s

With only commit/process ordering there's an irreducible crash window: commit-first risks loss, process-first risks duplicates — you can shrink it but never eliminate it. Exactly-once semantics (EOS) for the **consume-transform-produce** pattern works by making the *output records* and the *input offsets* commit **atomically** in a single Kafka transaction. The app uses a `transactional.id`, calls `initTransactions()`, then per cycle: `beginTransaction()`, produce the transformed records, call `sendOffsetsToTransaction(offsets, consumerGroupMetadata)` to write the consumed input offsets *inside the same transaction*, then `commitTransaction()`. The offsets land via the transaction, not via the consumer's normal commit. Downstream consumers set `isolation.level=read_committed` so they never see records from aborted or in-flight transactions. If the producer crashes mid-cycle the transaction aborts (or is fenced), so neither the output nor the advanced offsets become visible — on restart the input is reprocessed from scratch with no duplicate output. Kafka Streams enables all this with `processing.guarantee=exactly_once_v2`.

go deeper

for a junior

Know exactly-once exists and needs special config, not just careful ordering.

for a middle

Know it combines idempotent producer + transactions and that consumers need read_committed.

for a senior

Walk the transactional read-process-write API, offset-binding, markers, and fencing; state the external-system limit.

for a principal

Decide when EOS is worth its overhead vs idempotent at-least-once, and design for effects outside Kafka (outbox).

## Why ordering alone fails A consumer's two steps — **process** (produce output / write state) and **commit input offset** — cannot both be made crash-safe by sequencing: - **Commit-then-process**: a crash after commit skips the record → loss. - **Process-then-commit**: a crash after process, before commit, replays the record → the output is produced **twice** (duplicate). There is always a moment between the two steps where a crash leaves them inconsistent. No matter how tightly you sequence them, the window has nonzero width because they are two separate, independently-failing writes. Exactly-once requires them to **succeed or fail together** — i.e. atomicity. ## The consume-transform-produce pattern Many apps read from topic A, transform, and write to topic B (stream processing). The unit of work is: *consume input, produce output, advance the input offset*. EOS makes that whole unit atomic. ## Transactional API (mechanics) 1. Producer config: set a stable **`transactional.id`** (uniquely identifies this logical producer across restarts; required so the broker can **fence** zombie instances). Idempotence is implied/required. 2. `producer.initTransactions()` — registers the txn id with the **transaction coordinator** and fences any older producer with the same id (bumps the **epoch**). 3. Per cycle: - `beginTransaction()` - `producer.send(...)` the transformed output records to topic B. - `producer.sendOffsetsToTransaction(currentOffsets, consumer.groupMetadata())` — this writes the **consumed input offsets** to `__consumer_offsets` **as part of the same transaction**, not via the consumer's own `commitSync`. - `commitTransaction()` (or `abortTransaction()` on error). 4. The coordinator writes **transaction markers** (commit/abort control records) into every affected partition. Only on a commit marker do the output records and the offsets become visible. ## read_committed isolation Downstream consumers set **`isolation.level=read_committed`**. They buffer records until they see the transaction's commit/abort marker, and **never deliver records from aborted or still-open transactions**. The default `read_uncommitted` would expose duplicates from aborted attempts, breaking EOS for readers. The **Last Stable Offset (LSO)** bounds how far a read_committed consumer can advance. ## Crash handling - Producer crashes mid-transaction → the coordinator eventually **aborts** it; markers say 'abort'; neither outputs nor offsets are visible. On restart, `initTransactions()` fences the dead instance and the input is reprocessed cleanly — same output produced once, committed once. - A zombie (slow, presumed-dead) producer that wakes up is **fenced** by epoch, so it can't commit a stale transaction → no duplicate. ## Scope and limits - EOS is **within Kafka**: input topic → processing → output topic + offsets. A side effect to an **external** system (non-transactional DB, HTTP call) is *not* covered — you still need idempotency or a transactional outbox for those. - Kafka Streams wraps all of this: **`processing.guarantee=exactly_once_v2`** (the modern, single-producer-per-instance protocol; the older `exactly_once` is deprecated). ## One-line model EOS = idempotent producer (dedup retries) + transactions (atomic output + offset commit) + read_committed (hide aborted/in-flight) → the read-process-write cycle affects state once and only once.

  • Why must the input offsets be committed via sendOffsetsToTransaction instead of the consumer's commitSync?
    So the offset advance and the output records are in the SAME transaction and become visible atomically. A separate commitSync could succeed while the transaction aborts (or vice versa), reopening the duplicate/loss window.
  • Does Kafka EOS guarantee exactly-once when the processing step also writes to an external Postgres DB?
    No. Kafka transactions cover only Kafka topics + offsets. An external DB write isn't enrolled, so you need idempotent upserts or a transactional outbox to keep the external effect exactly-once.
  • What does isolation.level=read_committed change for a downstream consumer?
    It withholds records belonging to aborted or still-open transactions, delivering only committed records up to the Last Stable Offset — preventing it from seeing duplicates from aborted attempts.

saying these in an interview costs you the question

  • Claiming exactly-once means each record is transmitted only once (it means each record affects state once; retries still happen and are deduped).
  • Saying EOS extends to external systems automatically.
  • Forgetting read_committed on downstream consumers (without it, readers see aborted-transaction records).
  • Using the consumer's own commit instead of sendOffsetsToTransaction inside the transaction.
  • Confusing exactly_once with exactly_once_v2 / not knowing v2 is the current Streams setting.

context