Every streaming system eventually crashes in the middle of something. A worker has read a record, updated a total, written half its output, and then the process dies, the network partitions or a rebalance moves the partition elsewhere. When the work resumes, the system has to decide which records to process again. Re-process too little and records are lost; re-process too much and totals double. Exactly-once semantics is the name for designs where, after any number of such failures, the observable result is the same as if each input record had been processed one time.
This article is about the semantics and the Kafka-native machinery that implements them: the idempotent producer, transactions, the last stable offset, zombie fencing, and the patterns that carry the guarantee into a database or an external API. Flink's checkpoint barriers and two-phase-commit sink are covered separately in exactly-once stream processing architecture; the two pieces fit together, and this page explains the layer underneath.
What exactly-once actually promises
No network can deliver a message exactly once. A sender that gets no acknowledgement cannot tell whether the message was lost or the acknowledgement was, so it either retries (and risks a duplicate) or gives up (and risks a loss). That gives the two basic delivery guarantees: at-most-once, which never retries, and at-least-once, which retries until acknowledged.
Exactly-once systems do not beat that impossibility. They accept at-least-once delivery and make duplicates harmless, so each record's effect is applied once; hence the alternative name effectively-once.
The guarantee is end to end or it is nothing. It needs three things together: input that can be replayed from a known position, state and position saved atomically, and output that is either idempotent or committed atomically with the position. One HTTP call fired from inside the processor turns that effect back into at-least-once.
| Ingredient | Kafka mechanism |
|---|---|
| Replayable input | Retained, offset-addressed log partitions |
| Atomic state and position | Offsets committed inside the producer transaction; changelog topics for state |
| Idempotent or transactional output | Idempotent producer, transactions, read_committed readers |
Layer one: the idempotent producer
The first source of duplicates is the producer's own retry: the broker appends a batch, the acknowledgement is lost, and the producer sends the batch again.
The idempotent producer removes this case. The broker assigns it a producer id and an epoch, and the producer numbers its batches per partition with an increasing sequence number. The partition leader remembers the sequence numbers of each producer's most recent batches (the last five) and acknowledges a repeat without appending it. A sequence that skips ahead indicates lost data and raises an out-of-order sequence error.
Because only a window is remembered, max.in.flight.requests.per.connection must be at most 5, and acks=all is required so an acknowledged batch survives leader failover (see ISR replication). Since Kafka 3.0 both are defaults, but conflicting explicit settings such as acks=1 disable idempotence, so check the effective configuration.
Idempotence covers one producer session. A restarted producer gets a new id, and consumed input is not covered at all. For those you need transactions.
Layer two: transactions and the last stable offset
A transactional producer is configured with a transactional.id, a name that stays the same across restarts of the same logical worker. The id hashes to a partition of the internal __transaction_state topic, and the broker leading that partition acts as the transaction coordinator for the id.
The producer registers each partition it is about to write with the coordinator, then writes records as usual; they are durable and replicated but belong to an open transaction. To commit, the coordinator first logs prepare-commit, the point of no return, then writes a commit marker, a control record, into every partition the transaction touched. Abort writes abort markers instead.
The consumer setting isolation.level decides what readers see. The default, read_uncommitted, returns everything, including aborted records. With read_committed, a consumer reads only up to the last stable offset (LSO), the offset before the first still-open transaction, and skips aborted records. A downstream reader left on the default is the most common way exactly-once pipelines leak duplicates.
The cost: an open transaction holds back every read_committed reader of that partition, so their latency is at least the commit interval.
The consume-transform-produce loop
The pattern that makes a Kafka-to-Kafka stage exactly-once puts the consumed offsets into the producer's transaction. If the transaction commits, the output and the advanced input position become visible together; if it aborts, neither does, and the batch is read again. The loop below uses the Java client API; error handling follows the pattern in the KafkaProducer documentation.
producer.initTransactions(); // registers transactional.id, fences older producers
consumer.subscribe(List.of("payments")); // consumer: enable.auto.commit=false, read_committed
while (running) {
ConsumerRecords<String, Payment> batch = consumer.poll(Duration.ofMillis(200));
if (batch.isEmpty()) continue;
producer.beginTransaction();
try {
for (ConsumerRecord<String, Payment> r : batch) {
producer.send(new ProducerRecord<>("merchant-events", r.key(), enrich(r.value())));
}
producer.sendOffsetsToTransaction(nextOffsets(batch), consumer.groupMetadata());
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
producer.close(); // fatal: another instance owns this id, or config is wrong
throw e;
} catch (KafkaException e) {
producer.abortTransaction(); // nothing from this attempt becomes visible
rewindToCommitted(consumer); // the consumer's position moved; move it back
}
}
static Map<TopicPartition, OffsetAndMetadata> nextOffsets(ConsumerRecords<String, Payment> b) {
Map<TopicPartition, OffsetAndMetadata> out = new HashMap<>();
for (TopicPartition tp : b.partitions()) {
List<ConsumerRecord<String, Payment>> rs = b.records(tp);
out.put(tp, new OffsetAndMetadata(rs.get(rs.size() - 1).offset() + 1)); // next to read
}
return out;
}
static void rewindToCommitted(KafkaConsumer<String, Payment> consumer) {
Map<TopicPartition, OffsetAndMetadata> committed = consumer.committed(consumer.assignment());
for (TopicPartition tp : consumer.assignment()) {
OffsetAndMetadata om = committed.get(tp);
if (om != null) consumer.seek(tp, om.offset());
// with no committed offset, fall back to your auto.offset.reset policy
}
}Three details decide whether this loop is correct. Auto-commit must be off, or the consumer commits offsets outside the transaction. The offset committed is the next offset to read, the last processed offset plus one. And after an abort the consumer's in-memory position has already moved past the batch, so the loop must seek back; skipping rewindToCommitted silently drops the aborted batch, turning exactly-once into at-most-once.
The enrich step must also be deterministic. If it reads the wall clock, a random number or an external service, a retried batch can produce different output from the aborted attempt. The aborted records are hidden so there is no duplicate, but the committed result depends on which attempt succeeded. For stateful processing, Kafka Streams with processing.guarantee=exactly_once_v2 applies the same idea to state: store updates go to changelog topics inside the same transaction.
Fencing zombies
The hardest failure is a pause, not a crash. A worker stalls in garbage collection, the group decides it is dead and moves its partitions, and then the old worker wakes up holding a batch and tries to commit. Two writers now believe they own the same input.
Kafka fences the old worker with the epoch. When the replacement calls initTransactions() with the same transactional.id, the coordinator bumps the epoch and aborts any transaction the old instance left open. Requests with the old epoch are rejected with ProducerFencedException, which is fatal: close the producer and stop.
Passing consumer.groupMetadata() to sendOffsetsToTransaction adds a second fence on the consumer group's generation (KIP-447): a worker rebalanced out of the group cannot commit offsets for partitions it no longer owns. This lets each instance use one stable transactional.id rather than one per input partition. Rebalancing itself is covered in consumer rebalancing.
Since Kafka 4.0, brokers enable KIP-890, transactions server-side defense, by default; with 4.0 clients the epoch is bumped on every transaction, so a delayed write from a finished transaction cannot land in the next one and leave it hanging.
Exactly-once into a database
When the output is a database rather than a Kafka topic, the Kafka transaction cannot cover it. The standard answer is to stop treating Kafka's committed offsets as the source of truth and store the input position in the same database transaction as the output. Either the row changes and the offset advance commit together, or neither does.
def apply_batch(conn, worker_epoch, records):
"""One DB transaction: update totals and advance positions, or change nothing."""
with conn: # psycopg2: commit on success, rollback on error
with conn.cursor() as cur:
for r in records:
cur.execute(
"""INSERT INTO merchant_totals (merchant_id, total_cents)
VALUES (%s, %s)
ON CONFLICT (merchant_id)
DO UPDATE SET total_cents = merchant_totals.total_cents + EXCLUDED.total_cents""",
(r.key, r.amount_cents))
for (topic, part), (first, nxt) in offset_ranges(records).items():
# compare-and-set: only the owner at the expected position may advance it
cur.execute(
"""UPDATE consumer_positions
SET next_offset = %s, owner_epoch = %s
WHERE topic = %s AND partition = %s
AND next_offset = %s AND owner_epoch <= %s""",
(nxt, worker_epoch, topic, part, first, worker_epoch))
if cur.rowcount != 1:
raise StaleOwner(topic, part) # rolls back the whole batch
def on_partitions_assigned(consumer, conn, partitions):
for tp in partitions: # resume from the DB, not from Kafka
consumer.seek(tp, load_next_offset(conn, tp))The increment is not idempotent on its own; it becomes exactly-once because it commits atomically with the position. The compare-and-set on next_offset and owner epoch is database-side fencing: a zombie finds the position already advanced, matches zero rows, and its whole transaction rolls back.
Setting state is idempotent; adding to state is not. A keyed upsert of the latest balance can simply be replayed, while increments need the atomic position or a deduplication table keyed on an event id.
Effects outside your control
Emails, payment captures and partner API calls cannot join any transaction, and they fire on every attempt. Use idempotency keys and an outbox. Derive a stable key from the input, such as topic, partition and offset or a business event id, and send it with the request so a receiver that honours keys ignores retries. With an outbox, the processor writes an intent record transactionally and a relay performs the call at least once, keyed by the intent id.
Worked example: a crash mid-transaction
A payments aggregator reads payments partition 7, committed position 1,000, and writes enriched events to merchant-events. It polls offsets 1,000 to 1,199, begins a transaction, sends 200 output records, and adds offset 1,200 for partition 7 to the transaction. Before it calls commit, the host loses power.
The 200 records sit in an open transaction, and read_committed readers stop at the LSO in front of them. The coordinator would abort it after transaction.timeout.ms (60 seconds by default), but the replacement worker calls initTransactions() first: the epoch is bumped, abort markers are written, and the LSO moves on. The replacement resumes from offset 1,000, reprocesses the same 200 payments and commits. Committed readers see each event once; a read_uncommitted reader would have seen 400.
Failure modes
- A reader on read_uncommitted. It sees aborted records. Audit every consumer, including connectors and ad hoc tools.
- No rewind after abort. The aborted batch is silently skipped.
- Stuck transactions. One open transaction stalls the LSO for every committed reader of the partition.
- Timeout shorter than the work. Batches slower than
transaction.timeout.msare aborted. In Flink the sink timeout must exceed checkpoint interval plus recovery, within the broker'stransaction.max.timeout.ms(15 minutes by default). - Non-deterministic processing. Clock reads and random ids make retries differ. Derive ids from input coordinates and timestamps from the record.
- Unstable transactional ids. A random id per start disables fencing: zombie and replacement both commit.
- Retention shorter than an outage. Deleted input cannot be replayed; compacted topics need care too, see log compaction.
Operating it
Watch commit latency and abort rate, the gap between high watermark and LSO on output partitions, consumer lag on committed offsets, and fenced-producer errors, which should be near zero outside deployments. Each commit costs a coordinator round trip plus a marker per partition, so commit intervals of a few hundred milliseconds to a few seconds are typical; larger transactions trade reader latency for throughput.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| At-least-once plus idempotent upserts | Simplest, lowest latency, no transactions | Only works where every effect is a keyed overwrite |
| Kafka transactions end to end | True exactly-once for Kafka-to-Kafka stages | Commit latency for readers, coordinator load, every reader must be read_committed |
| Offsets stored in the sink database | Exactly-once into the database with no distributed commit | You own position management, seeking and fencing |
| Flink checkpoints plus two-phase-commit sinks | Exactly-once for large stateful jobs | Output latency tied to checkpoint interval, timeout tuning |
What to do next
- For each pipeline, write down the input, the state and every output, and mark which outputs are transactional, idempotent or neither.
- Confirm the producer's effective configuration shows idempotence on and acks=all, and give each worker a stable transactional.id.
- Audit every consumer of transactional topics for isolation.level=read_committed.
- Add the rewind-after-abort step and a test that kills the worker between send and commit, then checks output counts.
- For database sinks, move positions into the database with a compare-and-set update.
- Put idempotency keys on every external call, and alert on LSO lag, abort rate and fenced producers.