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