Topic 5.3
Exactly-Once Stream Processing (Read-Process-Write)
In one line
For pipelines that read from Kafka and write back to Kafka, exactly-once is a solved problem: wrap output writes and input offset commits in one transaction, or use Kafka Streams with processing.guarantee=exactly_once_v2, which also covers state store changelogs.
Think of it like this
A translator converting letters from one mailbox into another. If they put translated letters in the outbox and mark originals as done in one sealed step, a power cut can never leave a letter translated twice or marked done without a translation.
Key ideas
- 01
Read-process-write steps: consume input, compute output, produce output, commit input offsets. Without transactions, a crash after producing but before committing duplicates outputs on restart.
- 02
With transactions, output records, state changelog records and input offsets commit atomically; restarts resume from committed offsets and aborted outputs are invisible to read_committed consumers.
- 03
Kafka Streams:
processing.guarantee=exactly_once_v2(Kafka 2.5+ brokers) enables this automatically, using one transactional producer per stream thread; state stores are restored from their changelogs to match committed offsets. - 04
Scope: exactly-once covers inputs from Kafka, outputs to Kafka and Streams state. If the processor also calls a REST API or writes a database, those calls can repeat on retry and need their own idempotency.
- 05
Cost: commit interval (
commit.interval.ms, 100 ms under EOS by default) adds end-to-end latency, and throughput drops somewhat versus at-least-once.
Code & diagrams
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "billing-processor");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
StreamsBuilder b = new StreamsBuilder();
b.stream("commerce.orders", Consumed.with(Serdes.String(), orderSerde))
.filter((orderId, o) -> o.status() == Status.PAID)
.mapValues(BillingEvent::from)
.to("billing-events", Produced.with(Serdes.String(), billingSerde));
new KafkaStreams(b.build(), props).start();Interview problem
The problem
Exactly-once from orders to billing-events
Input topic orders, output topic billing-events. Design a processor that never produces duplicate billing records under retries, crashes or rebalances.
Explain it without notes
What does exactly_once_v2 cover in Kafka Streams, and what doesn't it cover?
Practice
Build the BillingStream, kill it repeatedly under load, and verify one billing event per paid order with a read_committed consumer.
Trade-offs
- ↔
EOS adds latency and some throughput cost; it's usually worth it for financial or count-sensitive pipelines and unnecessary for idempotent sinks.
Done when you can
I can build exactly-once Kafka-to-Kafka processing and state its boundary.