“Does Kafka support exactly-once delivery?” is a senior-level interview question with a trap built in. A plain “yes” is wrong, and so is “exactly-once is impossible.” The precise answer is that Kafka provides exactly-once processing for pipelines that read from Kafka and write to Kafka, using idempotent producers and transactions, and that anything outside Kafka needs its own idempotence. Interviewers then ask how transactions work, what a zombie is, what read_committed changes and why a slow transaction delays consumers.
This guide walks through the transactional consume-transform-produce loop in the Java client, the coordinator protocol behind it, the consumer side, and the limits. It targets Kafka 3.8/4.0. The code is illustrative: it follows the documented KafkaProducer transactional API but was not run against a cluster.
Before you start
You should be comfortable with offset commits and the difference between at-least-once and at-most-once; see Kafka offset commits and delivery semantics. You should also know what the idempotent producer does with producer IDs and sequence numbers, covered in Kafka idempotent producer and batching.
The short answer
In a consume-transform-produce loop, duplicates come from a crash between writing output and committing input offsets. Kafka transactions remove that gap: a producer with a transactional.id begins a transaction, sends the output records, adds the consumer’s offsets with sendOffsetsToTransaction, and commits, so outputs and offsets become visible together or not at all. The transactional.id and its epoch fence older “zombie” instances. Downstream consumers must use isolation.level=read_committed to skip aborted records and stop at the last stable offset. The guarantee covers Kafka-to-Kafka processing only; external side effects need idempotent writes.
How it works
Consider an at-least-once processor that reads payments, converts currency and writes payments-eur. If it writes the output and crashes before committing the input offset, a restart processes the same input again and writes a second output. Neither the idempotent producer (new session, new producer ID) nor careful commit ordering can prevent this, because two separate systems, output partitions and the group’s offsets, must change together.
Transactions make them change together. The pieces:
transactional.id: a stable name you configure, for examplefx-converter-0. It enables idempotence automatically.- Transaction coordinator: a broker chosen by hashing the
transactional.idonto a partition of the internal__transaction_statetopic, where it logs transaction state. - Producer ID and epoch:
initTransactions()gets the producer ID for thistransactional.idand increments its epoch. Requests carrying an older epoch are rejected, so a frozen old instance cannot write after a new one starts. Any transaction left open by the previous epoch is aborted. - Markers: records are written to their partitions immediately during the transaction. On commit or abort, the coordinator writes a control marker into every partition the transaction touched, including the
__consumer_offsetspartition holding the group’s offsets.
The commit is a two-phase protocol run by the coordinator: it records PrepareCommit in __transaction_state, writes the commit markers, then records CompleteCommit. Once PrepareCommit is logged, the transaction will complete even if the producer dies.
On the read side, a read_committed consumer only reads up to the last stable offset (LSO): the offset before the earliest still-open transaction in that partition. Records beyond it wait, even non-transactional ones, until that transaction ends. Aborted records are filtered out using an index of aborted transactions. The default, read_uncommitted, returns everything up to the high watermark, aborted records included.
Since Kafka 2.5 (KIP-447), passing consumer.groupMetadata() to sendOffsetsToTransaction lets the group coordinator also fence instances from an older group generation. That is why one transactional producer per application instance is enough, instead of one per input partition as in earlier versions.
Step-by-step walkthrough
Step 1: Configure the producer and consumer
The consumer must not commit on its own, and it should read only committed input if upstream is transactional.
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
p.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "fx-converter-" + instanceOrdinal); // stable per instance
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
Properties c = new Properties();
c.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
c.put(ConsumerConfig.GROUP_ID_CONFIG, "fx-converter");
c.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
c.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
c.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
c.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);The transactional.id must be stable across restarts of the same instance, which is what lets the replacement fence its predecessor, and unique across concurrently running instances.
Step 2: Initialize once and run the transactional loop
producer.initTransactions();
consumer.subscribe(List.of("payments"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> r : records) {
producer.send(new ProducerRecord<>("payments-eur", r.key(), toEur(r.value())));
}
producer.sendOffsetsToTransaction(nextOffsets(records), consumer.groupMetadata());
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
producer.close(); // fatal: another instance owns this transactional.id
throw e;
} catch (KafkaException e) {
producer.abortTransaction(); // abortable: undo outputs and offsets
rewindToCommitted(consumer);
}
}nextOffsets builds a map of each partition to its last processed offset plus one, the same rule as ordinary commits. The fatal exceptions mean this producer can never succeed again and must be closed; any other KafkaException is handled by aborting and retrying the batch.
Step 3: Rewind after an abort
Aborting undoes the outputs and the offsets, but the consumer’s in-memory position has already moved past the batch. Without a rewind, the next poll would skip those records.
void rewindToCommitted(KafkaConsumer<String, String> consumer) {
Map<TopicPartition, OffsetAndMetadata> committed = consumer.committed(consumer.assignment());
for (TopicPartition tp : consumer.assignment()) {
OffsetAndMetadata om = committed.get(tp);
if (om != null) consumer.seek(tp, om.offset());
else consumer.seekToBeginning(List.of(tp)); // or follow your auto.offset.reset policy
}
}Step 4: Read the output correctly downstream
Every consumer of payments-eur that must not see aborted data needs read_committed. This is easy to forget, because the default works and shows data: it just shows the wrong data during failures.
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");Step 5: Keep transactions short
A read_committed consumer cannot read past an open transaction, so end-to-end latency includes the transaction’s duration. Commit every poll batch rather than every few minutes. transaction.timeout.ms (60 seconds by default, capped by the broker’s transaction.max.timeout.ms of 15 minutes) makes the coordinator abort a transaction whose producer disappeared. Kafka Streams wraps all of this behind processing.guarantee=exactly_once_v2 and commits every 100 ms by default in that mode.
Worked scenario
A team ran the currency converter with at-least-once processing. A daily report summed payments-eur and occasionally showed totals a few hundred euros too high. The duplicates lined up with deploys: during each rebalance, batches that had been written but not yet committed were processed again by the new owner.
They converted the processor to the transactional loop above, and the next deploy still produced inflated totals for a few minutes. The report job used the default read_uncommitted, so it counted records from transactions that had been aborted when the old instances were fenced. Setting isolation.level=read_committed on the report consumer fixed the totals.
A week later, alerts showed the report’s lag climbing for several minutes on one partition while the converter was healthy. A bug made one instance open a transaction and then block on a slow reference-data call for minutes before committing. Every read_committed consumer on that partition waited at the LSO for it. The fix moved the reference-data lookup out of the transaction and lowered transaction.timeout.ms to 30 seconds so a stuck transaction would be aborted sooner.
Common mistake
- “Exactly-once means each message is delivered once over the network.” Records can still be fetched more than once; what is exactly-once is the committed effect inside Kafka.
- “Turning on transactions makes my database write exactly-once.” The database is outside the transaction. Use idempotent writes or store offsets in the database.
- Forgetting
read_committeddownstream. The default consumer sees aborted records. - A random
transactional.idper start. The new instance then cannot fence the old one, so zombie writes are possible. - Committing offsets with
consumer.commitSync()inside the loop. Offsets must go throughsendOffsetsToTransaction, or they are not part of the transaction.
Verify the behavior
Run the converter, then read its output with both isolation levels while forcing an abort, for example by throwing an exception after the sends in a test build:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic payments-eur \
--from-beginning --isolation-level read_uncommitted # shows records from aborted transactions
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic payments-eur \
--from-beginning --isolation-level read_committed # shows only committed recordsInspect transaction state directly with the transactions tool (Kafka 3.0+):
kafka-transactions.sh --bootstrap-server localhost:9092 list
kafka-transactions.sh --bootstrap-server localhost:9092 describe --transactional-id fx-converter-0
kafka-transactions.sh --bootstrap-server localhost:9092 find-hanging --broker-id 1describe shows the producer ID, epoch, state and the partitions in the current transaction; a transaction that stays Ongoing for longer than expected is the cause of a stalled read_committed consumer.
Follow-up questions
What is a zombie, and how is it fenced? An old instance that is still running after it was presumed dead, for example after a long GC pause. The new instance’s initTransactions() bumps the epoch, and the old one’s next transactional request fails with ProducerFencedException.
Why does a long transaction delay consumers? read_committed consumers stop at the last stable offset, which cannot move past an open transaction.
What does exactly_once_v2 change in Kafka Streams? It uses one transactional producer per stream thread with group-metadata fencing (KIP-447) instead of one per task, which scales much better. It requires brokers on 2.5 or later.
How would you get exactly-once into PostgreSQL? Write the rows and the partition offset in one database transaction and seek to the stored offset on assignment, or upsert by a unique event ID.
Interview exercise
A transactional processor reads orders, writes invoices, and calls sendOffsetsToTransaction. It crashes after sending 50 invoices but before commitTransaction(). A replacement with the same transactional.id starts 10 seconds later. Describe what happens to the 50 invoices, the group’s offsets, a read_committed consumer of invoices and a read_uncommitted consumer of invoices.
Answer and reasoning
The 50 invoices are already in the invoices log, but their transaction never committed. When the replacement calls initTransactions(), the coordinator bumps the epoch and aborts the open transaction, writing abort markers to the touched partitions, including the offsets partition. The group’s offsets stay at the last committed transaction, so the replacement reprocesses those orders and writes 50 new invoices in a new transaction. A read_committed consumer waited at the last stable offset during those seconds, then skips the aborted invoices and sees each invoice exactly once. A read_uncommitted consumer saw the first 50 as they were written and then sees the 50 replacements, so it observes duplicates. Without the replacement, the coordinator would have aborted the transaction itself after transaction.timeout.ms.
Continue learning
Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka offset commits and delivery semantics, Kafka idempotent producer and batching and the outbox pattern. Primary sources: the KafkaProducer Javadoc for Kafka 4.0, KIP-98 on exactly-once delivery and transactions, KIP-447 on producer scalability for exactly-once and Confluent’s article Transactions in Apache Kafka.