Kafka Stream Processing with Exactly-Once Semantics
Stream processing systems face a fundamental tension between correctness and performance. At-most-once processing loses data on failure. At-least-once processing creates duplicates. Exactly-once semantics — processing each record exactly one time with exactly one output — was considered impossible in distributed systems for decades. Apache Kafka's implementation of exactly-once processing, introduced in KIP-98 and refined through subsequent releases, solved this problem within the Kafka ecosystem by combining idempotent producers, transactional messaging, and atomic consumer offset commits into a single coordinated mechanism.
Understanding how exactly-once works requires understanding why it is difficult. In a read-process-write pipeline, a consumer reads from an input topic, processes the message, writes the result to an output topic, and commits the consumption offset. If the process crashes after writing the output but before committing the offset, the message will be reprocessed on restart, producing a duplicate output. If it crashes after committing the offset but before writing, the message is lost. Exactly-once semantics ties these operations together atomically — either all three happen or none do.
Idempotent Producers
The first building block is the idempotent producer, which prevents duplicate writes caused by producer retries. When a producer sends a message and the broker acknowledgment is lost due to a network partition, the producer retries. Without idempotency, this retry creates a duplicate message in the topic. The idempotent producer eliminates this by assigning each message a sequence number that the broker uses to deduplicate retries.
// Java — Idempotent producer configuration
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// Enable idempotent producer
props.put("enable.idempotence", "true");
// These are set automatically with idempotence, but shown for clarity:
// acks=all — wait for all in-sync replicas
// retries=Integer.MAX_VALUE — retry indefinitely
// max.in.flight.requests.per.connection=5 — allows batching with ordering
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// Sends are automatically deduplicated by the broker
producer.send(new ProducerRecord<>("orders", orderId, orderJson));
The broker maintains a map of (ProducerID, Partition) → SequenceNumber. Each new message from a producer must have a sequence number exactly one greater than the last committed sequence number for that partition. If the sequence number matches the last committed number, the broker returns success without writing — the message was already committed by a previous attempt. If the sequence number is greater than expected + 1, the broker rejects it as an out-of-order write, which indicates a bug in the producer implementation.
Transactional Messaging
Idempotent producers prevent duplicates within a single producer session, but they do not solve the read-process-write atomicity problem. Transactional messaging extends idempotency to span multiple partitions and tie producer writes to consumer offset commits within a single atomic operation.
// Transactional producer — read-process-write loop
Properties props = new Properties();
props.put("bootstrap.servers", "kafka-1:9092,kafka-2:9092,kafka-3:9092");
props.put("enable.idempotence", "true");
props.put("transactional.id", "order-processor-1"); // stable across restarts
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
consumer.subscribe(List.of("raw-orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
// Process the record
String enrichedOrder = enrichOrder(record.value());
// Write output within the transaction
producer.send(new ProducerRecord<>(
"enriched-orders", record.key(), enrichedOrder
));
}
// Atomically commit output records AND consumer offsets
Map<TopicPartition, OffsetAndMetadata> offsets = getOffsetsToCommit(records);
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (ProducerFencedException e) {
// Another instance with the same transactional.id is active
producer.close();
throw e;
} catch (Exception e) {
producer.abortTransaction();
}
}
The transactional.id is the key to cross-session exactly-once. When a producer starts with a transactional.id that was used by a previous instance (which may have crashed), the transaction coordinator fences the old producer and aborts any uncommitted transactions it left behind. This prevents zombie producers — crashed instances that come back after being replaced — from committing transactions that would create duplicates. The design philosophy is similar to how Apache Iceberg handles concurrent writes through optimistic concurrency control.
Kafka Streams EOS
Kafka Streams provides exactly-once semantics as a configuration option that handles all the transactional plumbing internally. Setting processing.guarantee to exactly_once_v2 enables transactional read-process-write for all stream processing operations including stateful transformations, aggregations, and joins.
// Kafka Streams exactly-once configuration
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-enrichment");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);
// Commit interval controls transaction frequency
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100);
// State store configuration for exactly-once
props.put(StreamsConfig.STATE_DIR_CONFIG, "/var/kafka-streams/state");
StreamsBuilder builder = new StreamsBuilder();
KTable<String, OrderAggregate> orderCounts = builder
.stream("raw-orders", Consumed.with(Serdes.String(), orderSerde))
.groupByKey()
.aggregate(
OrderAggregate::new,
(key, order, aggregate) -> aggregate.add(order),
Materialized.<String, OrderAggregate, KeyValueStore<Bytes, byte[]>>
as("order-counts-store")
.withKeySerde(Serdes.String())
.withValueSerde(aggregateSerde)
);
orderCounts.toStream().to("order-aggregates",
Produced.with(Serdes.String(), aggregateSerde));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
The exactly_once_v2 mode (introduced in Kafka 2.5) uses a single transaction coordinator per Streams task instead of per partition, which reduces the number of transactions and improves throughput. It also handles state store changes within the transaction — if a stateful aggregation is interrupted mid-transaction, the state store is restored to its pre-transaction state on recovery, preventing state corruption.
Consumer Isolation Levels
For downstream consumers to see only committed transactional output, they must use read_committed isolation level. Without this setting, consumers read uncommitted transaction data that may later be aborted, violating the exactly-once guarantee from the consumer's perspective.
// Consumer configured for transactional reads
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "kafka-1:9092");
consumerProps.put("group.id", "downstream-consumer");
consumerProps.put("isolation.level", "read_committed"); // critical
consumerProps.put("enable.auto.commit", "false");
// With read_committed, the consumer only sees messages from
// committed transactions. Aborted transaction records are filtered
// out transparently by the broker.
The read_committed consumer introduces a slight latency increase because it cannot read past the Last Stable Offset (LSO) — the offset of the earliest open transaction. If a producer has an open transaction with messages at offset 100 and committed messages exist at offsets 101-200, the read_committed consumer blocks at offset 100 until the transaction commits or aborts. This means long-running transactions stall downstream consumers, which is why transaction commit intervals should be kept short.
Performance Considerations
| Configuration | Latency Impact | Throughput Impact | When to Use |
|---|---|---|---|
| At-most-once | Lowest | Highest | Metrics, logs where loss is tolerable |
| At-least-once | Low | High | Idempotent consumers, key-value writes |
| Exactly-once (100ms commit) | +3-5% | -2-3% | Financial, order processing |
| Exactly-once (10ms commit) | +10-15% | -8-10% | Low-latency EOS (unusual) |
The commit interval is the primary tuning lever. At 100ms intervals (the default), the transaction overhead is amortized across hundreds or thousands of messages, making the per-message cost negligible. At 10ms intervals, each transaction contains fewer messages, so the fixed cost of the two-phase commit dominates. Most production deployments use 100-200ms intervals, which provides a good balance between latency and overhead. These tradeoffs are comparable to the consistency-performance decisions encountered in data quality framework design.
Failure Scenarios
Producer Crash During Transaction
When a producer crashes after beginning a transaction but before committing, the transaction coordinator detects the absence of heartbeats after transaction.timeout.ms (default 60 seconds) and aborts the transaction. All messages written within the aborted transaction are marked with abort markers, and read_committed consumers skip them. On restart, the new producer instance fences the old producer ID and starts fresh.
Broker Failure During Commit
If the transaction coordinator crashes during the two-phase commit, the new coordinator (elected from the ISR) replays the transaction log. If the commit decision was written before the crash, the coordinator completes the commit. If not, it aborts the transaction. The client retries the commit, which either succeeds (if the coordinator completed it) or receives an abort notification and retries the entire transaction.
Consumer Rebalance
When Kafka Streams tasks are rebalanced between instances, the new owner of a task restores its state store from the changelog topic and resumes processing from the last committed offset. Because state store updates and offset commits are within the same transaction, the restored state is always consistent with the committed offset — there is no window where state reflects uncommitted processing.
Exactly-once semantics in Kafka is not magic — it is a carefully constructed protocol that ties together multiple mechanisms to prevent duplicates and data loss within the Kafka ecosystem. Understanding these mechanisms helps you make informed decisions about when to use EOS (financial data, order processing, inventory management) and when at-least-once with idempotent consumers is simpler and sufficient (analytics, metrics, search indexing).