Skip to content

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

Exactly-Once Semantics

Exactly-once semantics (EOS) ensure Kafka transactional read-process-write pipelines process committed records exactly once. Kafka achieves this through idempotent producers, transactions, and transactional consumers.


Exactly-once guarantee from message send to unique processingExactly-once guarantee from message send to unique processingExactly-Once GuaranteeMessages Sent: 100Messages Delivered: 100Unique Processed: 100Achieved through:- Idempotent producer- Transactions- read_committed isolationall deliveredeach exactly once
PropertyGuarantee
Delivery countExactly 1
Message lossAvoided for committed transactions when replication is healthy
DuplicatesAvoided within Kafka transactions
AtomicityYes (transactions)

Components of exactly-once semanticsComponents of exactly-once semanticsExactly-Once ComponentsIdempotent ProducerTransactionsTransactional ConsumerProducer ID (PID)Sequence NumbersEpochTransactional IDTransaction CoordinatorAtomic Commitsread_committedConsumer Group Offsetenablescompletes
ComponentResponsibility
Idempotent ProducerPrevent duplicate writes from retries
Transaction CoordinatorManage transaction state
Transactional ConsumerRead only committed data
Consumer Group CoordinatorStores offsets committed atomically via the transaction coordinator

The idempotent producer assigns each producer instance a unique Producer ID (PID) and tracks sequence numbers per partition.

Idempotent producer deduplicating a retried produce requestProducerBrokerProducer(PID: 1000)Producer(PID: 1000)Broker(Partition Leader)Broker(Partition Leader)Idempotent producer enabledInitProducerIdPID=1000, epoch=0Produce(partition=0, seq=0, data=A)Store: PID=1000, seq=0ackProduce(partition=0, seq=1, data=B)Network timeoutRetry Produce(partition=0, seq=1, data=B)seq=1 exists for PID=1000ack (DuplicateSequenceException logged, no duplicate written)Produce(partition=0, seq=2, data=C)ack
StateAction
seq == expectedAccept record, increment expected
seq < expectedDuplicate; ignore
seq > expectedOut of order; reject
# Enable idempotent producer
enable.idempotence=true
# Automatically enforced:
# acks=all
# retries=Integer.MAX_VALUE
# max.in.flight.requests.per.connection <= 5
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
Producer<String, String> producer = new KafkaProducer<>(props);
// Retries are deduplicated automatically
for (int i = 0; i < 1000; i++) {
producer.send(new ProducerRecord<>("events", "key-" + i, "value-" + i));
}
producer.flush();
producer.close();
Idempotent producer coverage within a single producer sessionIdempotent producer coverage within a single producer sessionIdempotent Producer CoverageSingle SessionNot CoveredRetries deduplicatedNetwork failures handledProducer restartMultiple producersApplication crashes✅ Within producer lifecyclePID + sequence provides dedup❌ New PID on restartNeed transactions for cross-session

Transaction lifecycle from InitProducerId to commitProducerTransactionPartitionConsumerProducerProducerTransactionCoordinatorTransactionCoordinatorPartitionLeadersPartitionLeadersConsumerCoordinatorConsumerCoordinatorInitializationInitProducerId(transactional.id)PID=1000, epoch=0TransactionAddPartitionsToTxn(topic-0, topic-1)OKProduce(topic-0, records)Produce(topic-1, records)ackAddOffsetsToTxn(group-id)OKTxnOffsetCommit(offsets)OKCommitEndTxn(COMMIT)WriteTxnMarkers(COMMIT)WriteTxnMarkers(COMMIT)COMMIT complete
Transaction coordinator state transitionsTransaction coordinator state transitionsEmptyOngoingPrepareCommitPrepareAbortCompleteCommitCompleteAbortDeadInitProducerIdbeginTransaction()send(), addPartition()commitTransaction()abortTransaction()error/timeoutmarkers writtenmarkers writtentransaction completetransaction completetransactional.id expires
# Producer configuration
transactional.id=my-app-instance-1
enable.idempotence=true # Required for transactions
# Consumer configuration
isolation.level=read_committed # Only see committed transactions
enable.auto.commit=false # Manual offset management
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-processor-1");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// Initialize transactions (call once on startup)
producer.initTransactions();
try {
producer.beginTransaction();
// Send to multiple partitions atomically
producer.send(new ProducerRecord<>("orders", "order-1", "data-1"));
producer.send(new ProducerRecord<>("audit", "order-1", "audit-1"));
producer.send(new ProducerRecord<>("notifications", "user-1", "notify-1"));
// All writes commit together
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException e) {
// Fatal errors - cannot recover
producer.close();
throw e;
} catch (KafkaException e) {
// Abort and retry
producer.abortTransaction();
}

The read-process-write pattern consumes from input topics, processes, and produces to output topics atomically.

Read-process-write transaction covering output records and offset commitRead-process-write transaction covering output records and offset commitRead-Process-WriteInput TopicProcessing(Application)Output TopicConsumerOffsetsTransaction includes:1. Output records2. Consumer offset commitAtomic: all succeed or nonereadwritecommit
// Configure producer for transactions
Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "stream-processor-1");
producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// Configure consumer for read_committed
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "stream-processors");
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
producer.initTransactions();
consumer.subscribe(Collections.singletonList("input-topic"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
// Process and produce
for (ConsumerRecord<String, String> record : records) {
String output = process(record.value());
producer.send(new ProducerRecord<>("output-topic", record.key(), output));
}
// Commit offsets as part of transaction
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (TopicPartition partition : records.partitions()) {
List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);
long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
offsets.put(partition, new OffsetAndMetadata(lastOffset + 1));
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException e) {
// Fatal - close and restart
throw e;
} catch (KafkaException e) {
// Recoverable - abort and retry
producer.abortTransaction();
// Consumer will re-read uncommitted records
}
}
Read-process-write outcomes: success, abort, and zombie fencingRead-process-write outcomes: success, abort, and zombie fencingFailure HandlingSuccessFailure - AbortZombie Fencing1. Read records2. Process3. Produce output4. Commit offsets5. Commit transaction1. Read records2. Process3. Produce output4. Error occurs5. Abort transaction6. Re-read recordsOld producer continuesNew producer startsOld producer fenced

Isolation LevelBehavior
read_uncommittedSee all records including aborted transactions
read_committedSee only committed records; aborted filtered
Records visible to read_uncommitted and read_committed consumersRecords visible to read_uncommitted and read_committed consumersTopic: ordersread_uncommittedConsumerread_committedConsumerRecord A (committed)Record B (txn in progress)Record C (aborted)Record D (committed)Sees: A, B, C, DSees: A, DB filtered (in progress)C filtered (aborted)
# Transactional consumer
isolation.level=read_committed
enable.auto.commit=false
auto.offset.reset=earliest

The LSO is the offset up to which all transactions are complete.

Last stable offset in a partition log with an ongoing transactionLast stable offset in a partition log with an ongoing transactionPartition Log0: A (committed)1: B (committed)2: C (txn-1 ongoing)3: D (committed)4: E (txn-1 ongoing)5: F (committed)LSO = 2Consumer sees 0, 1HW = 6After txn-1 commits:LSO = 6Consumer sees 0, 1, 2, 3, 4, 5(C and E now visible)

Properties props = new Properties();
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-processor");
// Enable exactly-once v2 (Kafka 2.5+)
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);
KafkaStreams streams = new KafkaStreams(topology, props);
streams.start();
VersionKafka VersionDescription
exactly_once0.11.0+Original EOS; one producer per task
exactly_once_v22.5.0+Optimized; one producer per thread
Kafka Streams task committing output, changelog, and offsets atomicallyKafka Streams task committing output, changelog, and offsets atomicallyKafka Streams EOSStream TaskTransactionSourceProcessorSinkState StoreOutput recordsState changelogConsumer offsetsAll three committed atomically:- Output to downstream topics- State store changelog- Input consumer offsetsatomic commit

Latency timeline of at-least-once compared with exactly-once sendsLatency timeline of at-least-once compared with exactly-once sendsAt-Least-OnceSendAckDoneExactly-OnceSendTxnInitProduceCommitDone05101520
MetricAt-Least-OnceExactly-OnceOverhead
Latency (p50)~5ms (workload-dependent)~20ms (workload-dependent)+15ms (workload-dependent)
Latency (p99)~20ms (workload-dependent)~100ms (workload-dependent)+80ms (workload-dependent)
ThroughputHigh (workload-dependent)Moderate (workload-dependent)-30-50% (workload-dependent)
# Batch transactions for better throughput
# Process multiple records per transaction
max.poll.records=1000
# Longer commit intervals
transaction.timeout.ms=60000
# Increase producer batch size
batch.size=65536
linger.ms=10
ScenarioRecommendation
Financial calculationsUse EOS; correctness critical
Billing/meteringUse EOS; duplicates costly
Stream aggregationsUse EOS; state must be consistent
High-throughput loggingUse at-least-once; EOS overhead too high

When a producer crashes and restarts (or a new instance starts with the same transactional.id), the old producer is “fenced.”

Producer fencing after restart with the same transactional IDProducer ACoordinatorProducer A.Producer A(epoch=0)Producer A(epoch=0)CoordinatorCoordinatorProducer A'(epoch=1)Producer A'(epoch=1)InitProducerId(txn.id=app-1)PID=100, epoch=0Processing...Crash/RestartInitProducerId(txn.id=app-1)Increment epochPID=100, epoch=1A' is now the valid producerProduce (epoch=0)ProducerFencedExceptionOld producer fenced
try {
producer.commitTransaction();
} catch (ProducerFencedException e) {
// Another instance with same transactional.id is active
// This instance must shut down
log.error("Producer fenced - another instance is active", e);
producer.close();
System.exit(1);
}

Exactly-once coverage for Kafka topics compared with external systemsExactly-once coverage for Kafka topics compared with external systemsEOS CoverageKafka InternalExternal SystemsTopic A → Topic BConsumer offsetsState store changelogDatabaseCacheAPI✅ Full exactly-once guarantee⚠️ Requires additional patterns:- Idempotent writes- Two-phase commit- Outbox pattern
PatternDescriptionUse Case
Idempotent sinkSink handles duplicatesDatabase with unique constraints
Outbox patternWrite to Kafka via outbox tableDatabase + Kafka consistency
Saga patternCompensating transactionsDistributed workflow
// Use idempotency key for external writes
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
// Database write with idempotency
String idempotencyKey = record.topic() + "-" +
record.partition() + "-" + record.offset();
database.upsert(
"INSERT INTO events (idempotency_key, data) VALUES (?, ?) " +
"ON CONFLICT (idempotency_key) DO NOTHING",
idempotencyKey, record.value()
);
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}

FeatureMinimum Kafka Version
Idempotent producer0.11.0
Transactions0.11.0
Exactly-once (original)0.11.0
Exactly-once v22.5.0

ConfigurationRequiredDefaultDescription
enable.idempotenceYesfalseEnable idempotent producer
transactional.idFor txn-Unique transaction identifier
transaction.timeout.msNo60000Transaction timeout
max.in.flight.requests.per.connectionNo5Must be ≤ 5 for idempotence
ConfigurationRequiredDefaultDescription
isolation.levelYesread_uncommittedSet to read_committed for EOS
enable.auto.commitYestrueSet to false for EOS
ConfigurationDefaultDescription
transaction.state.log.replication.factor3Transaction log replication
transaction.state.log.min.isr2Minimum ISR for transaction log
transactional.id.expiration.ms604800000Transaction ID expiration