Exactly‑Once Semantics in Kafka: Eliminating Duplicate Processing in Stream Pipelines
Kafka’s Exactly‑Once Semantics (EOS) guarantees that each record is processed once, even across failures. This post walks through the problem, the EOS mechanism, how to enable it, a hands‑on Java example, trade‑offs, and how to verify your setup.
13 Jul 2025, 09:16 UTC

Problem: Duplicate Records in Stream Processing
When a record is written to a Kafka topic, the producer may retry on transient failures. If the retry is not idempotent, the same message can end up twice in the topic. Downstream consumers that aggregate or update state will then process the record twice, leading to incorrect totals, duplicate inventory adjustments, or erroneous fraud alerts.
In financial services, inventory systems, or any domain where a duplicate operation can have serious consequences, the cost of a single duplicate record can outweigh the performance hit of a more conservative delivery model.
EOS Mechanism: How Kafka Guarantees “Exactly Once”
Kafka’s Exactly‑Once Semantics (EOS) is a combination of three features:
- Idempotent Producer – The producer assigns a monotonically increasing sequence number per transaction and can safely retry sends without duplicating records.
- Transactional Writes – The producer groups a set of sends into a transaction, committing or aborting the whole batch atomically.
- Read‑Committed Consumer Isolation – Consumers configured with
isolation.level=read_committedsee only committed records, never the intermediate state of an uncommitted transaction.
With these three layers, Kafka ensures that a consumer will never see a record more than once, even if the producer or broker restarts.
Enabling EOS: Configuration Checklist
EOS requires a minimum set of broker and producer settings. The following table lists the critical properties and the recommended values for a production‑ready deployment.
| Component | Property | Recommended Value |
|---|---|---|
| Broker | transaction.state.log.replication.factor | ≥ 1 (ideally 3 for HA) |
| Broker | transaction.state.log.min.isr | ≥ 1 (or the broker’s replication factor) |
| Broker | transaction.state.log.segment.bytes | default 1073741824 (1 GB) – tune for log compaction |
| Producer | enable.idempotence | true |
| Producer | transaction.id | unique string per logical producer instance |
| Producer | max.in.flight.requests.per.connection | 5 (default) – keep default for EOS |
| Consumer | isolation.level | read_committed |
All topics involved in a transaction must have a replication factor of at least 1 and be configured with min.insync.replicas matching the broker’s transaction.state.log.min.isr.
Broker‑Side Commands
# Verify broker version (needs 0.11+)
kafka-topics.sh --version
# Example: create a topic for EOS with replication factor 3
kafka-topics.sh \
--create \
--topic eos-demo \
--partitions 3 \
--replication-factor 3 \
--bootstrap-server localhost:9092
Producer‑Side Setup
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "txn-demo-1");
// ... other producer configs
Producer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
Worked Example: Java Producer & Consumer
Below is a minimal Java snippet that demonstrates a transactional send and a consumer that reads only committed records.
// Transactional Producer
producer.beginTransaction();
producer.send(new ProducerRecord<>("eos-demo", "key1", "value1"));
producer.send(new ProducerRecord<>("eos-demo", "key2", "value2"));
producer.commitTransaction();
// Consumer with read_committed isolation
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "eos-consumer");
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
Consumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(Collections.singletonList("eos-demo"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("%s: %s%n", record.key(), record.value());
}
}
Running this code will print each key/value pair once, even if the producer is restarted between sends. The consumer never sees the intermediate state of the transaction.
Trade‑offs and Limitations
- Latency – Each transaction requires the broker to log the commit offset and replicate the transaction state log. This adds a few milliseconds to each batch.
- Throughput – Transaction boundaries limit the maximum number of records that can be sent in a single batch. Larger batches reduce overhead but increase the time to abort on failure.
- Disk Usage – The transaction state log is a separate log that can grow quickly if many producers issue frequent commits.
- Compatibility – EOS requires Kafka 0.11+ and that all topics in the transaction have a replication factor ≥ 1. Mixing EOS and non‑EOS topics in the same transaction is not supported.
- Application‑Level Deduplication – EOS handles intra‑transaction duplicates. Out‑of‑order or delayed messages that arrive after a transaction commit still need application‑level deduplication logic.
Verification Checklist
- Broker version – Ensure
kafka-topics.sh --versionreports ≥ 0.11. - Broker configs – Confirm
transaction.state.log.replication.factor≥ 1 andtransaction.state.log.min.isr≥ 1. - Topic replication – All EOS topics must have replication factor ≥ 1.
- Producer init – Verify
enable.idempotence=true,transaction.idset, andinitTransactions()called. - Consumer isolation – Ensure
isolation.level=read_committedand that the consumer is not usingread_uncommitted. - Test run – Send a transaction with two messages, abort the transaction, then commit a second transaction. Consume and confirm that only the committed messages appear once.
- Monitoring – Watch broker logs for “transaction aborted” events and consumer lag for EOS topics.
Actionable Takeaway
If duplicate processing can corrupt state or violate business rules, enabling Kafka’s Exactly‑Once Semantics is the most robust solution. Start by adding the configuration changes to a single producer and consumer pair, run the verification checklist, and measure latency/throughput. Once the metrics are acceptable, roll out EOS to all critical topics in your architecture. Remember that EOS is a trade‑off: correctness for a measurable performance cost.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.