Topic 5.2
Kafka Transactions
In one line
A transactional producer (with a stable transactional.id) can write to several partitions and commit consumer offsets as one atomic unit: either all records and the offset commit become visible, or none do. Consumers with isolation.level=read_committed skip aborted records. A transaction coordinator on the brokers manages the protocol.
Think of it like this
A bank transfer that debits one account and credits another. Either both entries appear in the ledger or neither does; nobody reading the ledger ever sees a debit without its credit.
Key ideas
- 01
API:
initTransactions()once (registers the transactional ID and fences older instances with the same ID), then per unit of workbeginTransaction(),send(...)to any topics,sendOffsetsToTransaction(offsets, consumer.groupMetadata()), andcommitTransaction()orabortTransaction(). - 02
Coordinator: each
transactional.idmaps to a transaction coordinator broker, which records state in the internal__transaction_statetopic and writes commit or abort markers into every partition the transaction touched. - 03
Fencing: a restarted instance with the same
transactional.idgets a higher producer epoch; the old instance (a zombie that paused and woke up) is fenced and its writes rejected withProducerFencedException. That's what prevents duplicate outputs when an instance is replaced. - 04
Consumers: with
isolation.level=read_committed, a consumer reads only committed transactional records (and non-transactional ones) up to the last stable offset (LSO); withread_uncommitted(default) it sees everything, including records later aborted. - 05
Costs: a few extra round trips per transaction and latency until commit (consumers see records only after commit). Batch multiple records per transaction (for example per poll) for throughput. Transactions only cover Kafka writes and offsets, not databases or APIs.
Code & diagrams
producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "billing-processor-" + instanceIndex);
producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
producer.initTransactions();
while (running) {
ConsumerRecords<String, Order> records = consumer.poll(Duration.ofMillis(200));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (ConsumerRecord<String, Order> r : records) {
producer.send(new ProducerRecord<>("billing-events", r.key(), toBillingEvent(r.value())));
offsets.put(new TopicPartition(r.topic(), r.partition()), new OffsetAndMetadata(r.offset() + 1));
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction(); // outputs + input offsets, atomically
} catch (ProducerFencedException e) {
producer.close(); // a newer instance took over: stop
throw e;
} catch (KafkaException e) {
producer.abortTransaction(); // outputs discarded, offsets not committed
rewindToLastCommitted(consumer); // reprocess the batch
}
}When it breaks
Consumers of transactional output use read_uncommitted (the default)
What you see
They process records from aborted transactions, producing duplicates or phantom effects even though the producer used transactions.
Fix & prevent
Set isolation.level=read_committed on every downstream consumer.
A long-open transaction
What you see
The last stable offset can't advance past it; read_committed consumers of those partitions stall, and lag grows for everyone.
Fix & prevent
Keep transactions short; transaction.timeout.ms aborts stuck ones; monitor open transaction duration.
Explain it without notes
How does transactional.id fencing prevent zombie duplicates?
Practice
Run the transactional processor, kill it mid-batch, restart it, and verify billing-events contains each order's billing event exactly once (read with read_committed).
Trade-offs
- ↔
Transactions give atomic multi-partition writes and exactly-once Kafka-to-Kafka processing, at the cost of latency, throughput per transaction, and operational complexity.
Done when you can
I can implement a transactional read-process-write loop and explain fencing and read_committed.