Command Palette

Search for a command to run...

Hectal
PHASE 5Advanced ~7 min· topic 3 of 5

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.

0/5 · 0%

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

  1. 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.

  2. 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.

  3. 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.

  4. 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.

  5. 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

BillingStream.javajava
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

01

What does exactly_once_v2 cover in Kafka Streams, and what doesn't it cover?

Practice

01

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.