skip to content

Consume-Transform-Produce EOS

Binding consumer offset commits into the producer transaction so a read-process-write loop is atomic. The canonical exactly-once pattern, and a favourite whiteboard question.

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

questions

5

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

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

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

You're scaling a hand-rolled consume-transform-produce service across many instances and partitions. How do you assign transactional.id values, and what are the trade-offs between the pre-KIP-447 per-partition model and the post-KIP-447 per-thread model?

level: principalimportance: should knowfreq 25%

basics

~20 s

On modern brokers (KIP-447), bind one transactional.id to each stable processing thread/instance and rely on consumer generation fencing for rebalance safety — far fewer producers. The legacy model required a deterministic transactional.id per input partition so fencing worked, which exploded the producer count at scale.

open as a page