Implementing Transactions

1. Understanding Transactional Concept

Concept Description Detail
Transaction Atomic writes across partitions All or nothing
Consume-Transform-Produce Read, process, write atomically EOS pattern
Coordinator Manages transaction state __transaction_state

Example: Transactional pipeline shape

producer.beginTransaction();
producer.send(out1);
producer.send(out2);
producer.sendOffsetsToTransaction(offsets, consumerGroupMeta);
producer.commitTransaction();

2. Enabling Transactional Producer

Config Description Detail
transactional.id Unique stable ID Enables tx + fencing
idempotence Auto-enabled Required
acks Forced to all Implicit

Example: Configure transactional id

props.put("transactional.id", "payment-processor-1");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);

3. Initializing Transactions

Method Purpose Detail
initTransactions() Register producer + fence old Once at startup
Recover Aborts incomplete prior tx Clean state
Epoch Bump Invalidates zombie producers Fencing

Example: Init transactions

producer.initTransactions(); // call once before any send

4. Beginning Transaction

Method Purpose Detail
beginTransaction() Start a new tx Before sends
State Sends buffered as tx Not visible yet
One Active No nested transactions Commit/abort first

Example: Begin a transaction

producer.beginTransaction();
producer.send(new ProducerRecord<>("ledger", key, debit));

5. Sending Transactional Messages

Aspect Description Detail
send() Part of current tx Buffered
Multi-partition Atomic across topics One commit
Visibility Only after commit (read_committed) LSO gated

Example: Multi-topic atomic write

producer.send(new ProducerRecord<>("accounts", from, debit));
producer.send(new ProducerRecord<>("accounts", to, credit));
producer.send(new ProducerRecord<>("audit", txId, record));

6. Sending Offsets in Transaction

Method Purpose Detail
sendOffsetsTo
Transaction
Commit input offsets atomically EOS read-process-write
Group Metadata consumer.groupMetadata() Required arg
No Manual Commit Tx commits the offset Don't commitSync

Example: Offsets within transaction

Map<TopicPartition, OffsetAndMetadata> offsets = currentOffsets(records);
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();

7. Committing Transaction

Method Effect Detail
commitTransaction() Make all writes visible Atomic
Blocks Until coordinator confirms Durable
Markers Writes commit markers read_committed sees

Example: Commit

producer.commitTransaction(); // all-or-nothing visible now

8. Aborting Transaction

Method Effect Detail
abortTransaction() Discard all sends Rollback
On Error Abort in catch block Cleanup
read_committed Never sees aborted data Filtered

Example: Abort on failure

try {
    producer.beginTransaction();
    process();
    producer.commitTransaction();
} catch (KafkaException e) {
    producer.abortTransaction();
}

9. Configuring Transaction Timeout

Config Description Default
transaction.timeout.ms Max tx duration 60000
transaction.max.
timeout.ms
Broker upper bound 900000
Exceeded Coordinator aborts tx Auto-rollback

Example: Set transaction timeout

props.put("transaction.timeout.ms", 30000);

10. Setting Read Committed

Aspect Description Detail
isolation.level read_committed Consumer side
LSO Reads stop at last stable offset No open tx data
Required for EOS Pairs with transactional producer End-to-end

Example: Consumer read committed

props.put("isolation.level", "read_committed");
props.put("enable.auto.commit", "false");

11. Understanding Exactly-Once Semantics

Component Role Detail
Idempotence Per-partition dedup Foundation
Transactions Atomic multi-partition + offsets Read-process-write
Streams EOS v2 Built-in exactly-once v2

Example: Streams EOS

props.put("processing.guarantee", "exactly_once_v2");
props.put("commit.interval.ms", 100);