Why this matters
- VaultCommerce payment-to-ledger pipeline uses
read_committedconsumers and transactional producer. - transactional.id uniquely identifies a producer instance across restarts — fencing zombies.
- EOS does not cover external DB writes — pair with idempotent sink or outbox.
- Higher latency than non-transactional — use only where duplicates are unacceptable.
Transaction flow
beginTransaction → produce to output topics → sendOffsetsToTransaction → commitTransaction. Abort on failure. Consumer sets isolation.level=read_committed to skip open transactions. VaultCommerce transactional.id=payment-ledger-pod-${POD_NAME}.
Key points
- transactional.id — fences stale producers; required for transactions
- read_committed — consumer isolation skipping uncommitted messages
- sendOffsetsToTransaction — atomic offset commit with output records
- Transaction coordinator — broker managing transaction state
- exactly_once_v2 — Kafka Streams EOS with improved protocol
Fencing and timeouts
transaction.max.timeout.ms bounds open transaction duration. New producer with same transactional.id fences old zombie. VaultCommerce keeps transactions under 30 seconds — longer work splits into smaller units.
VaultCommerce rollout checklist
Before promoting changes that touch the VaultCommerce order and payment event backbone, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. Transactions bind produces and offset commits. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.
// VaultCommerce payment ledger — transactional read-process-write
producer.initTransactions();
while (running) {
ConsumerRecords<String, PaymentEvent> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) continue;
try {
producer.beginTransaction();
for (var r : records) {
LedgerEntry entry = ledgerService.transform(r.value());
producer.send(new ProducerRecord<>("vaultcommerce.ledger.entries.v1",
entry.id(), entry));
}
Map<TopicPartition, OffsetAndMetadata> offsets = consumerPositions(records);
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
throw e;
}
}
Quick recall
Everything you need if you only revisit this box.
- Transactions bind produces and offset commits.
- Consumers need isolation.level=read_committed.
- External stores still need idempotency or outbox.
Test yourself
Answer these before moving on — recall is what makes it stick.