Kafka Exactly-Once Semantics

In distributed systems, messages can be lost or duplicated whenever a network call fails ambiguously: did the broker receive it or not? Exactly-once semantics (EOS) means each record's effect happens once, even when producers retry, brokers fail over, or consumers crash mid-batch. Kafka provides EOS through two building blocks: idempotent producers, which prevent duplicate writes to a partition, and transactions, which write to multiple partitions and commit consumer offsets atomically.

The crucial caveat: Kafka's exactly-once guarantees apply to reading from and writing to Kafka. Once a side effect leaves Kafka (a database write, an email, an HTTP call), you need idempotent design on that side to get the same effect.

TL;DR

Quick Example

A consume-transform-produce loop with transactions (Java):

Either the enriched records and the input offsets commit together, or neither does. There are no duplicates in orders.enriched and no skipped inputs.

Core Concepts

Delivery Semantics

Idempotent Producers

Each producer gets a producer ID (PID) and numbers every batch per partition. If a retry re-sends a batch the broker already wrote, the broker recognizes the sequence number and discards the duplicate. It also rejects out-of-order batches, preserving order. This is enabled by default in modern clients with acks=all. Its scope is one producer session and one partition. See Kafka producers.

Transactions

A transactional producer can:

  1. beginTransaction().
  2. Send records to any number of partitions and topics.
  3. Add consumed offsets with sendOffsetsToTransaction().
  4. commitTransaction() or abortTransaction().

A transaction coordinator on the broker tracks the transaction in the internal __transaction_state topic, and on commit writes commit markers into every involved partition. Records from aborted transactions stay in the log but are marked aborted.

transactional.id and Zombie Fencing

The transactional.id is a stable identity across restarts. When a new producer instance calls initTransactions() with the same ID, the coordinator bumps an epoch and fences the old instance. If a "zombie" (a paused or partitioned old process) wakes up and tries to commit, it gets ProducerFencedException. That prevents two instances from both committing results for the same input. With Kafka 2.5+ (KIP-447), fencing also uses consumer group metadata, so you no longer need one transactional ID per input partition.

read_committed Consumers

Consumers default to isolation.level=read_uncommitted, which also returns records from open or aborted transactions. Downstream consumers of transactional topics must set read_committed to see only committed data. They then read up to the last stable offset (the first offset of any still-open transaction), so long transactions increase end-to-end latency.

Kafka Streams and Other Frameworks

In Kafka Streams, exactly-once is a single setting:

Streams manages transactions, state store changelogs, and offset commits together, so stateful aggregations and joins stay consistent through failures. Apache Flink achieves end-to-end exactly-once with Kafka using its checkpointing plus Kafka transactional sinks (two-phase commit). See stream processing.

Exactly-Once With External Systems

Kafka transactions can't include a Postgres write or an API call. Options:

Best Practices

Use EOS Where Duplicates Are Costly

Financial ledgers, inventory counts, and aggregations where duplicates corrupt results benefit from transactions. For idempotent workloads (upserts, cache updates), at-least-once with idempotent consumers is simpler and faster.

Keep Transactions Short

Transactions hold back read_committed consumers until they commit. Commit per poll batch, keep transaction.timeout.ms reasonable, and avoid slow external calls inside the transaction.

Handle Fencing and Abort Paths Correctly

ProducerFencedException means stop, because another instance owns this work. Other errors mean abort, rewind the consumer to the last committed offsets, and retry. Test these paths deliberately with broker restarts and killed instances.

Make Every Downstream Consumer read_committed

One consumer left on read_uncommitted sees records from aborted transactions, silently breaking the guarantee for that pipeline.

Common Mistakes

Believing Idempotence Alone Gives End-to-End Exactly-Once

Idempotent producers prevent duplicate writes from retries. If your application crashes after processing and re-sends on restart, that's a new send with a new sequence, so it's duplicated. You need transactions or idempotent consumers.

Committing Offsets Outside the Transaction

Use sendOffsetsToTransaction so outputs and input offsets commit atomically.

Random transactional.id per Start

Generating a new random transactional.id on every restart defeats zombie fencing: old instances are never fenced. Use a stable ID per logical producer instance.

FAQ

Is exactly-once delivery actually possible?

Exactly-once delivery over an unreliable network is impossible in general, but exactly-once processing (each record's effect applied once) is achievable within a system that controls both sides. That's what Kafka provides for Kafka-to-Kafka pipelines, using idempotence, transactions, and atomic offset commits.

Does exactly-once hurt performance?

Somewhat. Transactions add coordinator round trips and commit markers, and read_committed consumers wait for commits. With reasonable batch sizes the throughput cost is modest, often single-digit percentages, and end-to-end latency grows with commit interval.

Do I need transactions if I only produce, not consume?

Only if you need atomic writes across multiple partitions or topics ("all these events or none"). For single-record publishing, the idempotent producer already prevents retry duplicates.

How does Kafka Streams achieve exactly-once?

It wraps each processing cycle in a transaction: output records, state store changelog updates, and input offset commits are committed atomically. After a failure, uncommitted work is aborted and reprocessed from the last committed offsets, with state restored from the changelog.

Related Topics

References