Kafka Exactly-Once Semantics with Transactions

Kafka is one of the few systems where "exactly-once" is a real, supported guarantee — within carefully drawn limits. This guide configures it for a consume-transform-produce pipeline as part of Exactly-Once vs At-Least-Once Delivery in Queue Fundamentals & Architecture, and is explicit about where the guarantee ends: at the edge of Kafka.

Problem Statement

A payments platform reads payment-authorized events, computes fees, and writes fee-calculated events that a ledger service consumes. During a broker failover, the fee processor produced its output, crashed before committing its consumer offset, restarted, and reprocessed the same input — writing the fee event twice. The ledger booked the fee twice for 1,400 payments. You want each input event to produce its output exactly once in the output topic, even across crashes, rebalances, and broker failovers, and downstream consumers to never see output from an attempt that did not complete.

Prerequisites

  • Kafka 2.5+ brokers (3.x recommended) with at least 3 brokers, so the transaction state log can be replicated (transaction.state.log.replication.factor=3, transaction.state.log.min.isr=2).
  • A client that supports transactions: the Java client, librdkafka-based clients (confluent-kafka for Python/Go/.NET), or Kafka Streams with processing.guarantee=exactly_once_v2.
  • Input and output both in Kafka — the pattern covers Kafka-to-Kafka processing.
  • Downstream consumers you can configure with isolation.level=read_committed.

Step 1 — Understand What the Transaction Covers

A Kafka transaction atomically commits three things: the output records written to one or more partitions, and the consumer offsets for the input records that produced them. Either all become visible, or none do. That atomic pairing is what turns at-least-once processing into exactly-once effect within Kafka: if the process crashes before commit, the transaction is aborted, the offsets are not advanced, and the input is reprocessed — but the aborted output is invisible to read_committed consumers.

begin transaction
  consume  payment-authorized @ offset 5812
  produce  fee-calculated     (payment 91, fee 2.30)
  sendOffsetsToTransaction(payment-authorized: 5813)   # offset is part of the txn
commit transaction   -> output + offset visible together, or neither
Output and offsets commit together The processor reads from the input topic, writes fee records to the output topic, and adds the consumer offset to the same transaction. The transaction coordinator commits them atomically. If the processor crashes before commit, the transaction aborts: output records are marked aborted and hidden from read_committed consumers, and the offset is not advanced, so the input is processed again. One atomic unit: output + offset payment-authorized offset 5812 transaction produce fee record + commit offset 5813 fee-calculated visible on commit Crash before commit: output aborted and hidden, offset not advanced, input reprocessed.

Step 2 — Configure a Transactional Producer

A transactional producer needs a transactional.id that is stable across restarts of the same logical processor instance. Kafka uses it to fence "zombie" instances — an old process that is still alive after its replacement started — so they cannot commit.

# fee_processor.py — confluent-kafka
from confluent_kafka import Producer, Consumer, KafkaException

INSTANCE = os.environ["POD_NAME"]           # stable per StatefulSet replica: fee-proc-0, fee-proc-1

producer = Producer({
    "bootstrap.servers": BOOTSTRAP,
    "transactional.id": f"fee-processor-{INSTANCE}",
    "enable.idempotence": True,             # implied by transactional.id
    "acks": "all",
    "transaction.timeout.ms": 60000,        # abort if a txn stays open longer than this
})
producer.init_transactions()                # fences any older producer with the same id

Deploy the processor as a Kubernetes StatefulSet (or anything else that gives each replica a stable identity) so fee-proc-0 always restarts with the same transactional id. With exactly_once_v2 (Kafka 2.5+ and Kafka Streams), fencing uses the consumer group generation instead, and a per-instance id is less critical — but a stable id remains the safest choice for hand-written processors.

Step 3 — Consume with Manual Offsets and read_committed

The consumer must not auto-commit — offsets are committed through the producer's transaction. It must also read only committed data from upstream transactional producers.

consumer = Consumer({
    "bootstrap.servers": BOOTSTRAP,
    "group.id": "fee-processor",
    "enable.auto.commit": False,            # offsets go into the transaction instead
    "isolation.level": "read_committed",    # never read aborted upstream output
    "auto.offset.reset": "earliest",
    "partition.assignment.strategy": "cooperative-sticky",
})
consumer.subscribe(["payment-authorized"])

read_committed is not optional for any consumer of transactional output — including the ledger service downstream. A read_uncommitted consumer sees records from aborted transactions, which is exactly the duplicate the transaction was supposed to hide.

Step 4 — Run the Read-Process-Write Loop in Transactions

Process a batch per transaction: begin, produce outputs, send the consumer's offsets into the transaction, commit. On a retriable error, abort and let the loop re-read from the last committed offsets.

def run() -> None:
    while True:
        msgs = consumer.consume(num_messages=500, timeout=1.0)
        if not msgs:
            continue
        producer.begin_transaction()
        try:
            for m in msgs:
                if m.error():
                    raise KafkaException(m.error())
                fee = compute_fee(json.loads(m.value()))
                producer.produce("fee-calculated", key=m.key(), value=json.dumps(fee))
            producer.send_offsets_to_transaction(
                consumer.position(consumer.assignment()),         # next offsets to read
                consumer.consumer_group_metadata())
            producer.commit_transaction()
        except KafkaException as e:
            err = e.args[0]
            if err.txn_requires_abort():
                producer.abort_transaction()
                rewind_to_committed(consumer)                     # re-read the batch
            elif err.retriable():
                continue                                          # commit will be retried
            else:
                raise                                             # fenced or fatal: exit, restart

rewind_to_committed seeks each assigned partition back to its committed offset, so the aborted batch is processed again inside a new transaction. Batch size trades latency for throughput: each commit costs a round trip to the transaction coordinator and writes transaction markers to every touched partition.

Fencing the zombie Instance A with transactional id fee-processor-0 pauses for a long garbage collection. Kubernetes starts a replacement B with the same id, which calls init_transactions and bumps the producer epoch. When A resumes and tries to commit its open transaction, the coordinator sees the old epoch and rejects it with a fencing error, so A's output is never visible. Same transactional.id, higher epoch wins instance A txn open, epoch 7 40 s GC pause commit: ProducerFenced instance B init_transactions: epoch 8, aborts A's txn, reprocesses Only one instance per transactional id can ever commit; the zombie's output stays invisible.

Step 5 — Know Where Exactly-Once Stops

The transaction covers Kafka writes and Kafka offsets. It does not cover a database write, an HTTP call, or an email sent from inside the loop. If the fee processor also inserted into Postgres and then the transaction aborted, the Postgres row would remain and the reprocessed batch would insert it again.

# WRONG: external side effect inside the Kafka transaction
producer.begin_transaction()
db.execute("INSERT INTO fees ...")          # not rolled back if the Kafka txn aborts
producer.produce("fee-calculated", ...)
producer.commit_transaction()

The boundary is easy to lose track of because the code looks transactional: a begin_transaction at the top and a commit_transaction at the bottom suggest that everything between them is covered. Only calls on the transactional producer are. Every other side effect inside the block happens immediately and permanently, whether or not the Kafka transaction later commits, and happens again when the aborted batch is reprocessed.

Inside and outside the guarantee Inside the Kafka transaction boundary are records produced to Kafka topics and the consumer offsets sent to the transaction. Outside it, even when called within the same code block, are database inserts, HTTP calls to other services, and emails. Those happen immediately and are repeated when an aborted batch is reprocessed. What exactly-once actually covers inside the transaction producer.produce(topic, ...) send_offsets_to_transaction outside, even in the same block INSERT INTO fees, HTTP POST send email, write to S3

Two correct designs: keep the processor Kafka-only and let a separate consumer (with read_committed) write to the database idempotently; or store the consumer offset in the database transaction alongside the row, and seek from the database-stored offset on startup — the database becomes the source of truth for progress. The idempotent-consumer side is covered in idempotent consumers with Postgres unique constraints.

Step 6 — Consider Kafka Streams for Most Pipelines

Kafka Streams implements all of the above with one setting and handles state stores, rebalances, and fencing for you.

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "fee-processor");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100);    // txn per 100 ms of work

StreamsBuilder b = new StreamsBuilder();
b.stream("payment-authorized", Consumed.with(Serdes.String(), paymentSerde))
 .mapValues(FeeCalculator::compute)
 .to("fee-calculated", Produced.with(Serdes.String(), feeSerde));
new KafkaStreams(b.build(), props).start();

commit.interval.ms controls transaction frequency; lower values reduce end-to-end latency for read_committed consumers (who see output only after commit), higher values improve throughput.

Verification

Kill the processor mid-transaction repeatedly and check that the output contains exactly one record per input.

# Produce 100k payments, run the processor, kill it every 20-40 s with SIGKILL
python produce_payments.py --count 100000
for i in $(seq 1 15); do kubectl delete pod fee-proc-0 --grace-period=0 --force; sleep $((RANDOM % 20 + 20)); done

# Count outputs per payment id with a read_committed consumer
kafka-console-consumer.sh --bootstrap-server $BOOTSTRAP --topic fee-calculated \
  --isolation-level read_committed --from-beginning --timeout-ms 30000 \
  --property print.key=true | cut -f1 | sort | uniq -d | wc -l      # expect 0 duplicates

Run the same count with --isolation-level read_uncommitted to see the aborted duplicates that exist in the log but are hidden from committed readers — a useful demonstration of why downstream isolation matters.

Gotchas & Edge Cases

Transaction timeouts. A transaction open longer than transaction.timeout.ms is aborted by the coordinator. Long processing per batch needs smaller batches, not a longer timeout, because open transactions also block read_committed consumers at the last stable offset.

Latency for downstream consumers. read_committed consumers cannot read past the first open transaction in a partition. A processor holding transactions open for seconds adds seconds of latency downstream.

Many transactional ids. Each id holds state in the transaction coordinator. Do not generate a new random id per process start; that defeats fencing and leaks coordinator state.

Exactly-once is per pipeline stage. Each stage that reads and writes Kafka needs its own transactional configuration; one transactional stage does not make the whole chain exactly-once.

FAQ

Does enable.idempotence alone give exactly-once? It prevents duplicates caused by producer retries within one producer session. It does not cover reprocessing after a consumer crash; transactions are needed for that.

Can I use transactions with SQS or RabbitMQ? No — this is a Kafka mechanism. Other brokers rely on idempotent consumers, which is also what you need at Kafka's edges. See preventing duplicate job execution with idempotency.

How do I monitor transactional processors? Watch the transaction abort rate (client metrics such as txn-abort counts, or the broker's transaction coordinator metrics), consumer lag per partition, and the gap between the log end offset and the last stable offset. A rising abort rate usually means processing exceeds the transaction timeout or instances are being fenced by rapid restarts; a growing gap to the last stable offset means an open transaction is holding back every read_committed consumer on that partition.

What does it cost in throughput? With reasonable batch sizes (hundreds of records or 100 ms per transaction), overhead is typically a few percent to low tens of percent. Tiny transactions per record are expensive.

Related