Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Transactional Producer

Transactions enable atomic writes to multiple partitions and exactly-once semantics when combined with transactional consumers.


ProducerTransactionBrokersProducerProducerTransactionCoordinatorTransactionCoordinatorBrokersBrokersinitTransactions()PID, epochbeginTransaction()send(topic-A)send(topic-B)sendOffsetsToTransaction()commitTransaction()write COMMIT markerscomplete

All writes within a transaction succeed or fail atomically.


# Required for transactions
transactional.id=my-app-instance-1
enable.idempotence=true # Implied
# Transaction timeout
transaction.timeout.ms=60000
# Read only committed transactions
isolation.level=read_committed
enable.auto.commit=false

Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-processor-1");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// Initialize once on startup
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("orders", key, value1));
producer.send(new ProducerRecord<>("audit", key, value2));
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException e) {
// Fatal - must close
producer.close();
} catch (KafkaException e) {
// Recoverable - abort and retry
producer.abortTransaction();
}
producer.initTransactions();
consumer.subscribe(Collections.singletonList("input"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
String output = process(record.value());
producer.send(new ProducerRecord<>("output", record.key(), output));
}
// Commit offsets atomically with output
producer.sendOffsetsToTransaction(
getOffsets(records),
consumer.groupMetadata()
);
producer.commitTransaction();
} catch (KafkaException e) {
producer.abortTransaction();
}
}

StateDescription
EmptyNo active transaction
OngoingTransaction in progress
PrepareCommitCommitting
PrepareAbortAborting
CompleteFinished

ExceptionSeverityAction
ProducerFencedExceptionFatalClose producer, another instance active
OutOfOrderSequenceExceptionFatalClose producer, state corruption
KafkaExceptionRecoverableAbort transaction, retry

When a producer with the same transactional.id starts, the old producer is fenced:

try {
producer.commitTransaction();
} catch (ProducerFencedException e) {
// Another instance took over
log.error("Fenced by new producer instance");
producer.close();
System.exit(1);
}

AspectImpact
Latency+10-50ms per transaction commit
Throughput30-50% lower than non-transactional
ComplexityHigher error handling requirements

Batch multiple records per transaction to amortize overhead.