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.

Advertisement

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.

IngredientKafka mechanism
Replayable inputRetained, offset-addressed log partitions
Atomic state and positionOffsets committed inside the producer transaction; changelog topics for state
Idempotent or transactional outputIdempotent 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.

Advertisement

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.

Consume-transform-produce with Kafka transactions: output and input position commit togetherInput topicpayments, 12 partitionsProcessorconsumer + txn producerOutput topicmerchant-totalspollsend (in txn)Txn coordinatora broker, by txn id hashbegin / commit__transaction_statetxn log: ongoing, prepare__consumer_offsetsgroup offsets (in txn)sendOffsetsToTransactionmarkersDownstream consumerread_committedup to LSO1. Processor polls a batch and begins a transaction.2. Output records are written, durable but invisible to read_committed.3. Input offsets are added to the same transaction.4. Commit: coordinator logs PREPARE, then writes markers to every partition.5. Abort or crash: markers say ABORT; readers skip the records; input is re-read.Either both the output and the new input position become visible, or neither does.
A transaction spans output partitions and the consumer-group offsets partition. Markers make the output and the new input position visible together.

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.ms are aborted. In Flink the sink timeout must exceed checkpoint interval plus recovery, within the broker's transaction.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

ChoiceGainCost
At-least-once plus idempotent upsertsSimplest, lowest latency, no transactionsOnly works where every effect is a keyed overwrite
Kafka transactions end to endTrue exactly-once for Kafka-to-Kafka stagesCommit latency for readers, coordinator load, every reader must be read_committed
Offsets stored in the sink databaseExactly-once into the database with no distributed commitYou own position management, seeking and fencing
Flink checkpoints plus two-phase-commit sinksExactly-once for large stateful jobsOutput latency tied to checkpoint interval, timeout tuning

What to do next

  1. For each pipeline, write down the input, the state and every output, and mark which outputs are transactional, idempotent or neither.
  2. Confirm the producer's effective configuration shows idempotence on and acks=all, and give each worker a stable transactional.id.
  3. Audit every consumer of transactional topics for isolation.level=read_committed.
  4. Add the rewind-after-abort step and a test that kills the worker between send and commit, then checks output counts.
  5. For database sinks, move positions into the database with a compare-and-set update.
  6. Put idempotency keys on every external call, and alert on LSO lag, abort rate and fenced producers.
Key takeaway: Exactly-once means each input's effect is applied once, built from at-least-once delivery plus duplicate-proof output. In Kafka that is the idempotent producer, transactions that commit output and input offsets together, read_committed readers bounded by the last stable offset, and epoch fencing of zombies. Outside Kafka, store the position in the same database transaction as the output and key every external call. The guarantee holds only if every link holds.