skip to content

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%

answer

  1. initTransactions once → fences + recovers
  2. begin → send → sendOffsets → commit
  3. offset = record.offset()+1
  4. pass consumer.groupMetadata(), not group.id
  5. abortTransaction in catch

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.

solid answer

~40 s

Setup (once): configure the producer with a transactional.id and call initTransactions(), which fences any prior producer with that id and recovers pending transactions. Set the consumer to enable.auto.commit=false. The per-loop cycle: poll() a batch; beginTransaction(); for each record, do the transform and producer.send() the outputs; build a map of TopicPartition -> OffsetAndMetadata for the next offset to read on each input partition; call producer.sendOffsetsToTransaction(offsetMap, consumer.groupMetadata()); then commitTransaction(). On any exception, abortTransaction() and let the consumer re-poll from the last committed position. sendOffsetsToTransaction forwards the offsets to the transaction coordinator so they're written to __consumer_offsets as part of the transaction's commit markers — the offset advance and the output records become durable in one atomic step. Passing consumer.groupMetadata() (not just the raw group.id) enables generation-based fencing of zombie consumers.

go deeper

for a junior

Recognize the begin/send/sendOffsets/commit shape and that the consumer must not auto-commit.

for a middle

Reproduce the loop correctly including offset+1 and abort-on-error.

for a senior

Explain initTransactions fencing/recovery and why groupMetadata() matters for zombie fencing.

for a principal

Reason about the coordinator's commit-marker protocol and design transactional.id assignment so fencing and recovery work across deployments.

## One-time setup ``` Producer config: transactional.id = <stable, unique per logical task> (enable.idempotence=true is implied) Consumer config: enable.auto.commit = false isolation.level = read_committed // for the *next* stage ``` Call **`producer.initTransactions()`** exactly once at startup. This: - Registers the `transactional.id` with the **transaction coordinator** (a broker role). - **Fences** any earlier producer instance that used the same `transactional.id` by bumping the **producer epoch**, so a zombie predecessor can no longer write. - Aborts or completes any in-flight transaction left over from a prior crash, giving you a clean slate. ## The per-batch loop ``` while (running) { records = consumer.poll(timeout) producer.beginTransaction() try { Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>() for (record : records) { // transform producer.send(new ProducerRecord<>(outTopic, ...)) offsets.put( new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)) // NEXT offset to read } producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata()) producer.commitTransaction() } catch (Exception e) { producer.abortTransaction() } } ``` ### Key details - **`record.offset() + 1`** — a committed offset is the position of the *next* record to consume, so you store last-processed + 1, exactly as a normal `commitSync` would. - **`sendOffsetsToTransaction(offsets, groupMetadata)`** — instead of the consumer committing offsets itself, the *producer* sends them to the transaction coordinator. They are buffered and only written to `__consumer_offsets` when the transaction commits. So the offset advance is **part of** the transaction. - **`consumer.groupMetadata()`** — supplies the consumer's group id plus its **member id** and **generation id**. The coordinator uses these to fence a stale (rebalanced-out) consumer: if this consumer's generation is no longer current, the offset commit is rejected, preventing a zombie from corrupting offsets. (The deprecated overload taking a bare `group.id` string lacks this fencing.) - **`commitTransaction()`** writes transaction *commit markers* to every involved partition (output topics and `__consumer_offsets`). It blocks until those are durable. **`abortTransaction()`** writes abort markers; `read_committed` consumers will skip the aborted records. ## Why this is atomic The output records and the offset records are tied to the same producer transaction. The coordinator's two-phase commit makes both sets of commit markers appear together. A crash before `commitTransaction()` leaves the whole transaction aborted, so neither the output nor the offset advance is visible, and the loop safely reprocesses. ## Common mistakes - Calling `consumer.commitSync()` anywhere — that commits offsets *outside* the transaction and breaks atomicity. - Storing `record.offset()` instead of `+ 1` — causes one record of reprocessing each restart. - Creating a new `transactional.id` per run — defeats fencing and recovery.

  • Why pass consumer.groupMetadata() instead of the plain group.id string?
    groupMetadata() carries member id and generation id, letting the coordinator fence a consumer that has been rebalanced out (a zombie). The bare-group.id overload is deprecated and lacks this fencing.
  • What offset value do you put in OffsetAndMetadata for a record at offset 100?
    101 — the committed offset is the next offset to read, i.e. last-processed offset + 1. Storing 100 would reprocess record 100 after a restart.

saying these in an interview costs you the question

  • Calling consumer.commitSync()/commitAsync() in the loop — that's an out-of-band commit that breaks the transaction.
  • Using record.offset() without +1.
  • Calling initTransactions() inside the loop instead of once at startup.
  • Passing group.id string to sendOffsetsToTransaction instead of groupMetadata(), losing generation fencing.

context