skip to content

Delivery Semantics and Transactions

At-most-once, at-least-once and exactly-once in Kafka, plus the transactional producer, read_committed consumers, and consume-transform-produce. Interviewers use this area to separate buzzword answers from real understanding.

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

explore

questions

page 1 of 2

What is the consume-transform-produce pattern in Kafka, and why does plain at-least-once delivery fall short of exactly-once for it?

level: juniorimportance: must knowfreq 70%

answer

  1. read-process-write loop
  2. two side effects: output + offset
  3. no safe ordering across a crash
  4. atomic produce + offset commit
  5. disable auto-commit

basics

~20 s

It's the read-process-write loop: a consumer reads from an input topic, the app transforms records, and a producer writes results to an output topic. At-least-once can produce duplicates or commit offsets without the output, so the output and the offset commit aren't atomic.

solid answer

~40 s

Consume-transform-produce (also called read-process-write) is the core stream-processing loop: poll records from an input topic, apply business logic, and produce derived records to one or more output topics. For end-to-end exactly-once you need two things to happen atomically: the produced output records and the advance of the consumer's input offsets. Plain at-least-once decouples these. If you commit the offset before the produce is durable, a crash loses the output (effectively at-most-once). If you produce first then crash before committing the offset, the input is reprocessed and you get duplicate output. Kafka solves this by wrapping the produce and the offset commit in a single producer transaction via sendOffsetsToTransaction, so either both commit or neither does.

go deeper

for a junior

Know the three steps and that output + offset must move together or you get duplicates / lost data.

for a middle

Explain why neither ordering of produce-then-commit is safe and name the transactional fix.

for a senior

Tie the atomicity to sendOffsetsToTransaction and read_committed consumers, and call out auto-commit must be off.

for a principal

Frame the scope boundary: EOS is Kafka-to-Kafka; reason about external side effects and when Kafka Streams' exactly_once_v2 replaces hand-rolling this loop.

## The pattern **Consume-transform-produce** (CTP), a.k.a. **read-process-write**, is the canonical stream-processing loop in Kafka: 1. **Consume** — a `KafkaConsumer` polls records from one or more *input* topic-partitions. 2. **Transform** — application code maps, filters, joins, or aggregates those records. 3. **Produce** — a `KafkaProducer` writes the derived records to one or more *output* topics. Progress through the input is tracked by **consumer offsets**: a per-(group, topic, partition) integer stored in the internal `__consumer_offsets` topic, marking the next record to read. Committing an offset is how the consumer says "I'm done with everything before here." ## Why at-least-once isn't exactly-once The loop has *two* externally visible side effects: (a) the records written to the output topic, and (b) the new committed input offset. End-to-end **exactly-once** means: for every input record, its effect on the output appears once and only once. That requires (a) and (b) to be **atomic** — all-or-nothing. Without transactions you must order them, and either order is broken: - **Commit offset, then produce** — if the app crashes after committing but before the produce is durable, the input record is never reprocessed (offset already advanced) and its output is lost. This is *at-most-once*. - **Produce, then commit offset** — if the app crashes after the produce but before committing, on restart the consumer re-reads the same input and produces the output *again*. This is *at-least-once* with **duplicates**. There's no safe ordering because a crash can land in the gap between the two operations. ## The fix Kafka transactions make the produce and the offset commit part of **one atomic transaction**. The producer is configured with a `transactional.id`; the app calls `beginTransaction()`, produces output, then calls `producer.sendOffsetsToTransaction(offsets, groupMetadata)` to fold the consumer's offset commit *into* the same transaction, then `commitTransaction()`. The offsets are written to `__consumer_offsets` by the transaction coordinator as part of the commit, so they become visible exactly when the output records do. Downstream consumers reading with `isolation.level=read_committed` only see committed output. ## Edge cases - You must **disable consumer auto-commit** (`enable.auto.commit=false`) — otherwise the consumer commits offsets out-of-band, defeating atomicity. - Exactly-once here is about *Kafka-to-Kafka* effects. Side effects to external non-transactional systems (e.g. an HTTP call) are not covered. - Kafka Streams implements this loop for you when `processing.guarantee=exactly_once_v2`.

  • Which Kafka API call makes the offset commit part of the producer transaction?
    producer.sendOffsetsToTransaction(offsetMap, consumerGroupMetadata), called between beginTransaction and commitTransaction.
  • Is end-to-end exactly-once guaranteed for a side effect like calling an external REST API in the transform step?
    No. Kafka transactions only make Kafka writes (output records + offsets) atomic; non-Kafka side effects are not part of the transaction and can repeat on reprocessing.

saying these in an interview costs you the question

  • Claiming at-least-once already gives exactly-once if you 'just commit carefully' — no ordering survives a crash in the gap.
  • Thinking exactly-once covers arbitrary external side effects, not just Kafka writes.
  • Forgetting that the offset commit, not just the produce, must be transactional.

context

open as a page

What is the scope boundary of Kafka's exactly-once semantics (EOS), and why doesn't it automatically extend to external systems?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Kafka EOS only guarantees exactly-once within a single Kafka cluster: a read-process-write where the input topic, output topic, and consumer offsets all live in that same cluster. External databases, APIs, or other clusters are outside the transaction, so it can't cover them.

open as a page

What are the three delivery semantics in Kafka (at-most-once, at-least-once, exactly-once), and what does each guarantee about message delivery?

level: juniorimportance: must knowfreq 80%

basics

~20 s

At-most-once: each message is delivered zero or one time (loss possible, no duplicates). At-least-once: delivered one or more times (duplicates possible, no loss) — Kafka's default. Exactly-once: delivered once and only once (no loss, no duplicates).

open as a page

What does it mean to make a Kafka consumer idempotent, and why does that let you live safely with at-least-once delivery?

level: juniorimportance: must knowfreq 75%

basics

~20 s

An idempotent consumer can process the same message more than once with no extra effect. Kafka's at-least-once delivery can redeliver a record after a crash; if processing is idempotent, those duplicates are harmless, so you get effectively exactly-once results.

open as a page

What does the consumer setting isolation.level do in Kafka, and what are its two possible values?

level: juniorimportance: must knowfreq 65%

basics

~10 s

isolation.level controls whether a consumer sees records from in-progress or aborted transactions. read_uncommitted (default) returns all records; read_committed only returns records from committed transactions, hiding aborted and not-yet-committed ones.

open as a page

What is Kafka's transactional producer and what problem does it solve compared to a plain producer?

level: juniorimportance: must knowfreq 70%

basics

~10 s

A transactional producer lets you write to multiple topic-partitions atomically: either all writes commit and become visible together, or all are aborted. A plain producer writes each record independently with no all-or-nothing guarantee.

open as a page

What is the Kafka transaction coordinator, and how does a producer find the one assigned to it?

level: juniorimportance: must knowfreq 55%

basics

~20 s

The transaction coordinator is a broker-side component that manages a transactional producer's state. A producer locates it by hashing its transactional.id to a partition of the internal __transaction_state topic; the leader of that partition is the producer's coordinator.

open as a page

What is a transaction marker (control batch) in Kafka, and who writes it to a partition?

level: juniorimportance: must knowfreq 55%

basics

~20 s

A transaction marker is a special COMMIT or ABORT record the transaction coordinator writes at the end of a transaction into every partition the transaction touched. It tells consumers whether the preceding transactional records are valid.

open as a page

Why must enable.auto.commit be false in a transactional consume-transform-produce loop, and what breaks if it's left on?

level: middleimportance: must knowfreq 55%

basics

~20 s

Auto-commit makes the consumer commit offsets on its own timer, outside the producer transaction. That decouples the offset advance from the output produce, so a crash can advance offsets without the matching output (lost data) — defeating exactly-once. Offsets must be committed only via sendOffsetsToTransaction.

open as a page

Walk through the exact producer/consumer call sequence for a transactional read-process-write loop, and explain what sendOffsetsToTransaction does.

level: middleimportance: must knowfreq 65%

basics

~10 s

Once, call producer.initTransactions(). Per batch: poll records, beginTransaction(), produce outputs, sendOffsetsToTransaction(offsets, consumer.groupMetadata()), then commitTransaction() (or abortTransaction() on error). sendOffsetsToTransaction writes the consumer offsets into the same transaction so they commit atomically with the output.

open as a page

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%

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.

open as a page

On the consumer side, how do the order of offset-commit vs record-processing produce at-most-once versus at-least-once, and what crash window causes loss or duplicates in each?

level: middleimportance: must knowfreq 70%

basics

~10 s

Commit-before-process = at-most-once: if you crash after committing but before processing, the record is skipped (lost). Process-before-commit = at-least-once: if you crash after processing but before committing, you reprocess it (duplicate).

open as a page

What is the Last Stable Offset (LSO), and how does it gate reads for a read_committed consumer?

level: middleimportance: must knowfreq 55%

basics

~20 s

The LSO is the offset of the first record belonging to a still-open transaction. A read_committed consumer is never allowed to read past the LSO, so it stops there until that transaction commits or aborts.

open as a page

Walk through the transactional producer API lifecycle: transactional.id, initTransactions, beginTransaction, commitTransaction, abortTransaction.

level: middleimportance: must knowfreq 65%

basics

~10 s

Configure transactional.id, call initTransactions() once at startup. Then per unit of work: beginTransaction(), send records, and either commitTransaction() on success or abortTransaction() on failure. Repeat begin/commit for each transaction.

open as a page

Explain the transactional outbox pattern and why teams choose it over distributed 2PC to integrate a database with Kafka.

level: seniorimportance: must knowfreq 60%

basics

~20 s

The outbox pattern writes the business row and an 'outbox' event row in one local database transaction, then a separate relay (often Debezium CDC) reads the outbox and publishes to Kafka. It avoids two-phase commit (2PC) across the DB and Kafka by relying only on the local DB transaction plus at-least-once publishing with idempotency.

open as a page

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%

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.

open as a page

When would you choose idempotent at-least-once processing over Kafka's exactly-once semantics (EOS) for a consumer?

level: seniorimportance: must knowfreq 65%

basics

~20 s

Choose idempotent consumers when the side effect lands outside Kafka — a database, HTTP call, cache, or other system. Kafka EOS only makes consume-transform-produce atomic within Kafka, so for external sinks idempotence is simpler and the only thing that actually works.

open as a page

What is ProducerFencedException, when does Kafka throw it, and how should an application handle it?

level: seniorimportance: must knowfreq 50%

basics

~20 s

ProducerFencedException means another producer with the same transactional.id registered with a newer epoch, so this older instance is fenced out. It is fatal: you cannot continue or abort — close the producer and let only the newer instance proceed.

open as a page

Walk through what InitProducerId does and how the producer epoch is used to fence a previous instance sharing the same transactional.id.

level: seniorimportance: must knowfreq 60%

basics

~20 s

InitProducerId asks the coordinator for a producer ID and epoch. When a producer reuses an existing transactional.id, the coordinator bumps the epoch and aborts any in-flight transaction from the old instance. Brokers then reject writes carrying the now-stale (lower) epoch, fencing the old producer.

open as a page

Walk through the two-phase atomic commit Kafka uses to commit a transaction across multiple partitions.

level: seniorimportance: must knowfreq 50%

basics

~20 s

On commit, the coordinator first durably writes a PrepareCommit record to __transaction_state (phase 1), then writes COMMIT markers to every touched partition and finally a CompleteCommit record (phase 2). The prepared state lets it retry markers after a crash so the commit is atomic.

open as a page

On the producer side, how do acks and retries determine whether you get at-most-once or at-least-once, and what failure produces a duplicate write?

level: middleimportance: should knowfreq 60%

basics

~20 s

Retries cause duplicates when a write succeeds on the broker but the ack is lost — the producer resends, so the record is stored twice (at-least-once). Disabling retries (or acks=0) avoids duplicates but risks losing records that weren't acked (at-most-once).

open as a page

Some operations like sending an email or charging a card can't be made naturally idempotent. How do you make a consumer effectively exactly-once for those, conceptually?

level: middleimportance: should knowfreq 55%

basics

~20 s

If the side effect can't be a simple overwrite, you make it idempotent by giving each unit of work a stable unique id and recording 'already done' so a redelivered record is recognized and skipped. The action plus the 'done' record must commit together.

open as a page

A teammate says 'we enabled enable.idempotence on the producer, so our consumers are exactly-once now.' What's wrong with that statement?

level: middleimportance: should knowfreq 60%

basics

~20 s

Producer idempotence only stops the producer's own retries from writing duplicate records into a partition. It does nothing for consumer-side processing. A consumer can still reprocess a record after a crash, so consumer-side idempotence is a completely separate concern.

open as a page

Which broker and producer configs govern the transaction coordinator and zombie fencing, and what do they control?

level: middleimportance: should knowfreq 35%

basics

~10 s

Producer side: transactional.id (enables fencing) and transaction.timeout.ms. Broker side: transaction.state.log.num.partitions (coordinator sharding), transaction.state.log.replication.factor and min.isr (durability), and transaction.max.timeout.ms (cap on producer timeout).

open as a page

How do transaction markers, the Last Stable Offset (LSO), and read_committed consumers interact?

level: middleimportance: should knowfreq 40%

basics

~20 s

A read_committed consumer can only read up to the Last Stable Offset (LSO) — the offset before the earliest still-open transaction. Markers let the LSO advance: once a COMMIT/ABORT marker lands, the broker knows the outcome and exposes those records (delivering committed, skipping aborted).

open as a page

Explain how consumer.groupMetadata() in sendOffsetsToTransaction prevents zombie consumers from corrupting offsets, and how this fencing relates to producer transactional.id fencing.

level: seniorimportance: should knowfreq 40%

basics

~20 s

groupMetadata() carries the consumer's group id, member id, and generation id. The transaction coordinator rejects an offset commit whose generation is stale, so a consumer that was rebalanced out (a zombie) can't commit. This is generation-based fencing, complementing the producer's epoch-based transactional.id fencing.

open as a page

Why does cross-cluster replication with MirrorMaker 2 not provide exactly-once semantics, and what are the consequences for consumers on the target cluster?

level: seniorimportance: should knowfreq 45%

basics

~20 s

MirrorMaker 2 is a Kafka Connect consume-from-source, produce-to-target pipeline across two clusters. Because the source consume and target produce span different clusters with no shared transaction, it's at-least-once: failures cause duplicate records on the target. Offsets and producer state don't carry over, so target consumers must tolerate duplicates and remapped offsets.

open as a page

Walk through exactly where the duplicate-creating window is in a consumer that does an idempotent write, and how offset-commit choices interact with it.

level: seniorimportance: should knowfreq 50%

basics

~20 s

The window is between performing the side effect and committing the offset: a crash there causes the record to be reprocessed after rebalance/restart. Idempotent writes make that reprocessing harmless. Offset-commit choice (auto vs manual, before vs after) only shifts how often duplicates happen, not whether idempotence is needed.

open as a page

How does a read_committed consumer actually avoid delivering records from aborted transactions?

level: seniorimportance: should knowfreq 35%

basics

~20 s

The broker keeps an aborted-transaction index per segment and sends the relevant aborted (producerId, firstOffset) entries with each fetch. The consumer reads records below the LSO and drops any whose producerId+offset fall inside an aborted transaction's range.

open as a page

What latency and operational trade-offs does isolation.level=read_committed introduce, and how would you mitigate them?

level: seniorimportance: should knowfreq 40%

basics

~10 s

read_committed adds end-to-end latency because consumers can't read past the LSO until a transaction commits, and one long/stuck transaction stalls the whole partition. Mitigate with short transactions, a sane transaction.timeout.ms, and LSO-lag monitoring.

open as a page

showing 1–30 of 40