Command Palette

Search for a command to run...

Hectal
PHASE 5Advanced ~8 min· topic 2 of 5

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.

0/5 · 0%

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

  1. 01

    API: initTransactions() once (registers the transactional ID and fences older instances with the same ID), then per unit of work beginTransaction(), send(...) to any topics, sendOffsetsToTransaction(offsets, consumer.groupMetadata()), and commitTransaction() or abortTransaction().

  2. 02

    Coordinator: each transactional.id maps to a transaction coordinator broker, which records state in the internal __transaction_state topic and writes commit or abort markers into every partition the transaction touched.

  3. 03

    Fencing: a restarted instance with the same transactional.id gets a higher producer epoch; the old instance (a zombie that paused and woke up) is fenced and its writes rejected with ProducerFencedException. That's what prevents duplicate outputs when an instance is replaced.

  4. 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); with read_uncommitted (default) it sees everything, including records later aborted.

  5. 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

TransactionalProcessor.javajava
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
    }
}
txn.mermaiddiagram
Rendering diagram…

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

01

How does transactional.id fencing prevent zombie duplicates?

Practice

01

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.