In a read-process-write pipeline, how does Spring Kafka make the output sends and the consumer offset commit atomic (exactly-once)?
answer
- container begins tx per batch
- sendOffsetsToTransaction(offsets, groupMetadata)
- offsets committed BY producer, not consumer
- enable.auto.commit=false is mandatory
- EOSMode.V2 = one producer per group
basics
~20 sConfigure the listener container with a KafkaTransactionManager. For each batch it begins a transaction, runs your listener's sends, then commits the consumed offsets into that same transaction via the producer, and commits once. If anything fails it aborts, so neither the output nor the offset advance.
solid answer
~40 sRead-process-write means: consume a record, process it, produce results, and commit the consumer offset — and you want all of that to be atomic. Spring drives it through the listener container. When the container has a KafkaTransactionManager (set on the ConcurrentKafkaListenerContainerFactory), for each poll batch it: begins a Kafka transaction on a transactional producer, invokes your @KafkaListener where you send output via KafkaTemplate, then calls producer.sendOffsetsToTransaction(offsets, consumerGroupMetadata) to write the input offsets *into the same transaction*, and finally commits. Because the offset commit and the output records commit together, an abort rewinds both — the record will be re-read and reprocessed, and downstream read_committed consumers never saw the aborted output. Requirements: transactional producer (transactionIdPrefix), consumer enable.auto.commit=false, downstream isolation.level=read_committed. Note offsets are committed by the producer, not the consumer, so auto-commit must be off.
code
java · 30 lines@Configuration
class RpwConfig {
@Bean
ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
ConsumerFactory<String, String> cf,
KafkaTransactionManager<String, String> tm) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
factory.setConsumerFactory(cf); // cf must have enable.auto.commit=false
// Container starts a tx per batch and sends offsets to it:
factory.getContainerProperties().setKafkaAwareTransactionManager(tm);
return factory;
}
}
@Component
class Processor {
private final KafkaTemplate<String, String> template;
Processor(KafkaTemplate<String, String> template) { this.template = template; }
@KafkaListener(topics = "orders-in", groupId = "orders")
public void onMessage(String order) {
// Runs inside a container-started Kafka transaction.
template.send("orders-out", enrich(order));
// Offsets for orders-in are committed into THIS transaction by the container
// via producer.sendOffsetsToTransaction(...). Throw -> everything aborts.
}
private String enrich(String s) { return s.toUpperCase(); }
}go deeper
Know the container wraps produce + offset commit in one transaction that aborts together.
Name sendOffsetsToTransaction and the enable.auto.commit=false requirement.
Explain 'effectively once', EOSMode.V2, and non-Kafka side-effect idempotency.
Reason about fencing, scaling with V2, LSO stalls, and cross-resource consistency design.
## The pattern **Read-process-write**: a service consumes from an input topic, transforms, produces to output topic(s), and advances the consumer offset. Naively these are separate operations — a crash between 'produce output' and 'commit offset' causes duplicates (reprocess) or, if ordered the other way, loss. Kafka transactions make **produce + offset-commit one atomic unit**. ## How Spring wires it You set a `KafkaTransactionManager` on the `ConcurrentKafkaListenerContainerFactory` (Boot does this automatically when `transaction-id-prefix` is present and there's a `KafkaTransactionManager`). Then the **listener container** (specifically `TransactionalContainer` logic in `KafkaMessageListenerContainer`) orchestrates each batch: 1. **Begin**: obtain a transactional producer and `beginTransaction()`. 2. **Process**: invoke your `@KafkaListener`; any `KafkaTemplate.send(...)` in that thread enrolls in the transaction. 3. **Send offsets**: the container calls `producer.sendOffsetsToTransaction(offsetsToCommit, consumer.groupMetadata())`. This records the *input* offsets **inside the transaction** on the internal `__consumer_offsets` topic. Using `groupMetadata()` (vs. a bare group id) enables proper fencing (EOSMode.V2). 4. **Commit / abort**: on success `commitTransaction()` (atomically commits output records *and* the offsets); on any thrown exception `abortTransaction()` — both the output and the offset advance are rolled back, so the record is redelivered. ## Mandatory configuration - Producer: `transactionIdPrefix` set (transactional + idempotent). - Consumer: `enable.auto.commit=false` — **the producer commits offsets, not the consumer**. Auto-commit would double-commit and break atomicity. - Downstream consumers: `isolation.level=read_committed` so aborted output is never observed. ## Why it's 'exactly-once' (really 'effectively once') On abort/crash, the record is re-consumed and reprocessed, but because the previous output was aborted and downstream readers use `read_committed`, no consumer ever observed the duplicate output, and the offset never advanced past unprocessed work. The *processing* may run more than once, but the *observable output and offset progression* happen exactly once. This is why Kafka's EOS is precisely a **read-process-write within Kafka** guarantee. ## EOSMode and producers Spring's `EOSMode.V2` (the default; requires broker/client ≥ 2.5) uses a **single producer per consumer group** and `sendOffsetsToTransaction` with consumer group metadata for fencing. The older `V1` (per-partition producers) is removed in Spring Kafka 3.x. V2 scales far better because it doesn't need a producer per topic-partition. ## Gotchas - Forgetting `enable.auto.commit=false` — offsets get committed twice / non-atomically. - Side effects to non-Kafka systems (DB, external API) inside the listener are **not** covered by the Kafka transaction unless separately coordinated; reprocessing can repeat them. Make them idempotent. - A hung listener holds the transaction open, stalling downstream `read_committed` consumers at the LSO. - The guarantee is only end-to-end if downstream is `read_committed`.
- Why must enable.auto.commit be false for transactional read-process-write?Because the offset commit is done by the transactional producer via sendOffsetsToTransaction so it commits atomically with the output. If the consumer also auto-committed, offsets would advance outside the transaction, defeating atomicity and risking data loss on abort.
- Are database writes inside the listener covered by the Kafka transaction?No. The Kafka transaction only spans Kafka sends and offset commits. External side effects (DB rows, HTTP calls) are not rolled back on abort and may repeat on reprocessing, so they must be made idempotent or coordinated with a separate transaction manager.
- What does EOSMode.V2 change versus the older V1 mode?V2 (default, needs broker 2.5+) uses a single producer per consumer group and passes consumer group metadata to sendOffsetsToTransaction for fencing, instead of V1's producer-per-topic-partition. It scales much better and V1 was removed in Spring Kafka 3.x.
saying these in an interview costs you the question
- Saying the consumer commits the offsets in a transactional pipeline (the producer does).
- Leaving enable.auto.commit=true with transactions.
- Believing DB side effects are rolled back by the Kafka transaction.
- Thinking exactly-once means the listener body runs only once (it can re-run; output is effectively once).
- Omitting downstream read_committed and still claiming end-to-end EOS.