Exactly-Once in Kafka Streams: One Crash, Two Guarantees, and Where the Transaction Stops

Lab run with apache/kafka:4.3.1 (KRaft), kafka-streams 4.3.1 and Java 17; 3 integration tests

Revision note (2026-09-15). The earlier version of this article told readers to set processing.guarantee=exactly_once, a value the current Kafka Streams client rejects; said brokers 0.11 or newer suffice, where exactly_once_v2 needs 2.5 or newer; pinned kafka-streams 3.4.0 and "Java 8+" (4.x clients need Java 11); named the metrics transaction.commit-rate and transaction.abort-rate, which do not exist; and claimed that after a crash the application "reprocesses messages, guaranteeing no duplicates" without saying what that covers. This version runs one crash scenario under both guarantees and counts the records.

The guarantee, stated precisely

With processing.guarantee=exactly_once_v2, Kafka Streams writes the output records, the state store changelog records, and the consumer offsets of each task inside one Kafka transaction and commits them together every commit.interval.ms. If the application dies before the commit, the transaction is aborted: the records are still physically in the topics, but consumers with isolation.level=read_committed never see them, and the offsets go back to the previous commit, so the restart reprocesses the input from there.

That is the whole guarantee. It says nothing about work done outside Kafka topics during processing, and nothing about consumers that read with read_uncommitted, which is the consumer default. The lab makes both edges visible.

Tested versions: apache/kafka:4.3.1 broker (KRaft, one node, transaction.state.log.replication.factor=1), kafka-streams 4.3.1, Amazon Corretto 17.0.14, Gradle 8.8, Docker 27.4.0. Lab: examples/kafka-streams-eos.

The topology

The earlier article's count-per-key topology, with two hooks: a counter incremented in peek() that stands in for any external side effect (an HTTP call, an e-mail, a row in another database), and a switch that throws on the n-th record.

builder.<String, String>stream(input)
    .peek((key, value) -> {
        int n = sideEffects.incrementAndGet();          // "external" side effect, not transactional
        if (n == crashAt.get()) {
            throw new IllegalStateException("simulated crash after " + n + " records");
        }
    })
    .groupByKey()
    .count(Materialized.as("counts-store"))
    .toStream()
    .to(output, Produced.with(Serdes.String(), Serdes.Long()));

Configuration, identical for both runs except for the guarantee:

props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, guarantee);   // exactly_once_v2 or at_least_once
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 5000);           // the crash lands before the first commit
props.put(StreamsConfig.STATESTORE_CACHE_MAX_BYTES_CONFIG, 0);      // forward every count update immediately
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1);

Without the cache set to zero, count updates sit in the record cache until the commit, and the "what does read_uncommitted see" question has a less interesting answer. With EOS the default commit.interval.ms is 100 ms; the lab raises it so that the crash reliably happens inside the first, uncommitted transaction.

The scenario

  1. 100 input records, keys k0 to k9, 10 each. The correct final count for every key is 10.
  2. Start the application with the crash armed on the 60th record. The side effect runs, then the exception is thrown. Kafka Streams' default StreamsUncaughtExceptionHandler chooses SHUTDOWN_CLIENT; the test waits for state ERROR. 59 count updates were forwarded before the crash.
  3. While the application is down, read the output topic with read_uncommitted and with read_committed.
  4. Start a second instance with the same application.id and state.dir, no crash, and read the output again with both isolation levels, keeping the last value per key.

exactly_once_v2

[lab] exactly_once_v2: side effects after crash=60, total=160 (inputs=100)
[lab] exactly_once_v2: while down, output records visible: read_uncommitted=59 read_committed=0
[lab] exactly_once_v2: after restart, output records: read_committed=100 read_uncommitted=159
[lab] exactly_once_v2: final count per key (read_committed) = {k0=10, k1=10, k2=10, k3=10, k4=10, k5=10, k6=10, k7=10, k8=10, k9=10}

Read from the bottom up. The final counts are right and there are exactly 100 committed output records, one per input: the aborted run's 59 writes and the changelog entries that went with them were rolled back, the state store was rebuilt without them, and the input was reprocessed from offset 0. That is the guarantee working.

Now the two edges. A consumer with read_uncommitted saw 59 records while the application was down and 159 at the end. Those 59 are the aborted writes; they are real records in the log with real offsets. read_uncommitted is the Java consumer's default, so a downstream service that never set isolation.level reads the duplicates. And the side-effect counter reached 160: 60 in the crashed run (the 60th record ran its side effect before throwing) plus 100 in the restart. Anything your processor does that is not a write to a Kafka topic or a Kafka Streams state store happened 1.6 times per input here, and Kafka has no way to undo it.

at_least_once, same crash

[lab] at_least_once: side effects after crash=60, total=160 (inputs=100)
[lab] at_least_once: while down, output records visible: read_uncommitted=59 read_committed=59
[lab] at_least_once: after restart, output records: read_committed=159 read_uncommitted=159
[lab] at_least_once: final count per key (read_committed) = {k0=20, k1=20, k2=20, k3=14, k4=13, k5=13, k6=20, k7=13, k8=13, k9=13}

Nothing was committed before the crash, so the restart also reprocessed from offset 0. But the 59 output records and the changelog writes of the first run were ordinary, non-transactional writes: the state store was restored from that changelog, the counts continued from there, and every key ended above 10. All 159 records are visible to every consumer. Kafka did nothing wrong; the application counted twice, which is what at-least-once means.

The two runs differ only in processing.guarantee. The side-effect count is 160 in both.

What "exactly-once" covers, in one table

Effect of processingexactly_once_v2at_least_once
Records written to Kafka output topics, read with read_committedonce (100 records, counts of 10)duplicated (159 records, counts up to 20)
Records in the same topics, read with read_uncommittedduplicates visible (159)duplicates visible (159)
State store contents after restartrebuilt without the aborted workincludes the aborted work
Anything else your code does per recordran 160 times for 100 inputsran 160 times for 100 inputs

Consumers of your output must set isolation.level=read_committed to get the first row; Kafka Streams applications downstream do that automatically only if they themselves run with exactly_once_v2. External effects need their own idempotence: a natural key in the target table, an idempotency key on the HTTP call, or moving the effect out of the topology into a connector or consumer that reads with read_committed and is idempotent on its own.

Configuration facts that changed

The value exactly_once and the constant StreamsConfig.EXACTLY_ONCE are gone from the 4.x client, not deprecated:

[lab] processing.guarantee=exactly_once -> Invalid value exactly_once for configuration processing.guarantee: String must be one of: at_least_once, exactly_once_v2
[lab] StreamsConfig.EXACTLY_ONCE and EXACTLY_ONCE_BETA: no such field in kafka-streams 4.3.1

Other facts, from the 4.3 configuration reference and the running client:

  • exactly_once_v2 requires brokers 2.5 or newer, not 0.11.
  • The 4.x clients and Kafka Streams need Java 11 or newer; brokers need Java 17. kafka-streams 4.3.1 is the current release on Maven Central.
  • By default EOS requires transaction.state.log.replication.factor=3 and min.isr=2 on the brokers; the lab sets both to 1 because it has one broker, and the Streams documentation is explicit that a replication factor below 3 "effectively voids EOS". The lab shows the client-side semantics, not durability.
  • Under EOS the producer's transaction.timeout.ms defaults to 10000 ms; it is a producer setting, passed to Streams as producer.transaction.timeout.ms, not a Streams property.
  • The transaction metrics the client actually exposes are producer metrics: txn-init-time-ns-total, txn-begin-time-ns-total, txn-send-offsets-time-ns-total, txn-commit-time-ns-total, txn-abort-time-ns-total (plus txn-prepare-time-ns-total in 4.3.1), and the Streams thread metric commit-rate. The lab lists them from KafkaStreams.metrics(); there is no transaction.commit-rate.

Error handling inside a transactional topology

The earlier version suggested a try/catch in the processor with a dead-letter topic. Two things to know before doing that. Catching and swallowing means the record is committed as processed, inside the transaction, so it will not be retried; that is a choice, and it should be an explicit one. And an exception that escapes, as in the lab, does not go to a dead-letter topic: with the default processing.exception.handler (LogAndFailProcessingExceptionHandler) and the default uncaught exception handler, the client shuts down, which is what the ERROR state in the lab was. If you want records skipped or routed, configure processing.exception.handler (default LogAndFailProcessingExceptionHandler in the 4.3 reference) rather than catching inside the topology, and write the dead-letter record through the same Streams producer so it is part of the transaction.

What this does not cover

  • Durability. One broker, replication factor 1. Broker failure during a transaction and producer fencing between two live instances were not exercised.
  • transaction.timeout.ms expiry and what happens to a task whose commit takes longer than that.
  • The throughput and latency cost of exactly_once_v2. Not measured.
  • A real external sink. The side effect is an in-memory counter, which is enough to count executions but says nothing about how a particular sink behaves.

Reproduce it

cd examples/kafka-streams-eos
./start-broker.sh          # apache/kafka:4.3.1 on localhost:29092
gradle test --no-daemon    # 3 tests, about 2 minutes (waits for the aborted transaction and restores)
./stop-broker.sh

Sources