Kafka courseLesson 9 of 10
Kafka course · Lesson 9 of 10
Kafka Delivery Semantics: At-Most-Once, At-Least-Once and Exactly-Once
What at-most-once, at-least-once and exactly-once mean in Kafka, how transactions and read_committed work, and how to get exactly-once effects in external sinks.
On this page
- Sample setup
- Where loss and duplicates come from
- At-most-once
- How you get it
- When it is acceptable
- Pitfalls
- In interviews
- At-least-once
- How you get it
- Making duplicates harmless
- Pitfalls
- In interviews
- Exactly-once semantics (EOS)
- The building blocks
- What EOS does not cover
- Pitfalls
- In interviews
- Transactions in Kafka
- How it works
- Worked example: commit, abort, read
- Zombie fencing
- Pitfalls
- In interviews
- Idempotence end to end
- The chain
- Worked example: offsets stored in the sink
- Pitfalls
- In interviews
- Read-process-write pattern
- The loop
- Worked example: a crash before commit
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
“Exactly-once” is the most misunderstood phrase in streaming. Kafka offers a precise, real guarantee, but it covers a specific scope: reading from Kafka, processing, and writing back to Kafka. Everything outside that scope (a database, an API, a file) needs its own design. This lesson starts from the two simple guarantees, shows where each one breaks, then builds exactly-once with transactions on a real cluster and extends it to external sinks.
Sample setup
The examples use a single-partition transfers topic, an output topic transfer-audit, and a three-broker Kafka 4.3.1 cluster. Transactions need the internal __transaction_state topic, whose defaults (transaction.state.log.replication.factor=3, transaction.state.log.min.isr=2) assume at least three brokers; on a single development broker set both to 1.
bin/kafka-topics.sh --bootstrap-server localhost:19092 --create --topic transfers \
--partitions 1 --replication-factor 3
bin/kafka-topics.sh --bootstrap-server localhost:19092 --create --topic transfer-audit \
--partitions 1 --replication-factor 3
Where loss and duplicates come from
Every pipeline has three places where a failure can strike:
| Step | Failure | Effect without protection |
|---|---|---|
| Producer to broker | response lost, producer retries | duplicate record in the topic |
| Broker | leader fails before replication | acknowledged record lost (with acks=1) |
| Consumer processing | crash between processing and committing | record reprocessed, or skipped |
The semantics below are statements about all three steps together, not about one setting.
At-most-once
What it is. Each record is processed zero or one times: it may be lost, but never duplicated.
How you get it
- Consumer: commit the offset before processing (or use librdkafka’s default auto offset store with auto-commit, which marks records as consumed when they are handed to your code).
- Producer:
acks=0, or no retries.
A crash after the commit and before processing skips the rest of the batch. The consumers lesson showed this with a simulation: committing first lost payment-6 and payment-7.
When it is acceptable
Metrics, logs and telemetry where a small gap is cheaper than duplicates or extra latency, for example a live dashboard that recomputes every few seconds anyway.
Pitfalls
- Getting it by accident. Auto-commit combined with asynchronous processing is at-most-once even if you meant at-least-once.
- Assuming loss is rare. Every deployment and rebalance is a possible loss point.
In interviews
Define it, show the commit-before-process order, and give an example where it is the right choice. Interviewers like candidates who say it is sometimes correct.
At-least-once
What it is. Each record is processed one or more times: nothing is lost, but duplicates are possible. It is the practical default for data pipelines.
How you get it
- Producer:
acks=all, retries on (the defaults), idempotence on (the default) to avoid producer-side duplicates. - Topic: replication factor 3,
min.insync.replicas=2, unclean leader election off. - Consumer: process the batch, write it to the sink, then commit offsets.
A crash between the sink write and the commit makes the next owner of the partition reprocess everything since the last commit. Rebalances do the same if offsets are not committed in onPartitionsRevoked.
Making duplicates harmless
At-least-once becomes “effectively once” when processing a record twice has the same effect as processing it once:
| Sink | Idempotent write |
|---|---|
| Relational database | INSERT ... ON CONFLICT (event_id) DO UPDATE / MERGE keyed by a business ID |
| Data lake table (Delta, Iceberg) | MERGE on a key, or overwrite of a deterministic partition |
| Key-value store | PUT key value (last write wins) |
| Search index | index by document ID |
| Another Kafka topic | keyed, compacted topic, or transactions (below) |
Appends without a key (INSERT into a log table, sending an email, charging a card) are not idempotent; they need a deduplication table or an idempotency key at the receiver.
Pitfalls
- Deduplicating on Kafka offsets alone. Offsets change if data is replayed into a new topic or after an unclean election. Use a business key or an event ID carried in the payload.
- Non-deterministic processing. If reprocessing produces a different result (current time, random IDs), duplicates are not harmless even with upserts.
In interviews
Expect “how do you handle duplicates in an at-least-once pipeline?” The strong answer is idempotent sinks keyed by an event ID, plus the reasons at-least-once is preferred over at-most-once for data you must not lose.
Exactly-once semantics (EOS)
What it is. In Kafka, exactly-once semantics means that, for a pipeline that reads from Kafka topics and writes to Kafka topics, the effect of each input record appears exactly once in the outputs and in the consumer offsets, even when producers retry, brokers fail and consumers crash. Records may physically be written more than once; consumers using read_committed simply never see the extra copies.
The building blocks
- Idempotent producer: no duplicates from retries within a producer session (see producers).
- Transactions: a producer writes to several partitions and commits consumer offsets atomically; all of it becomes visible, or none of it.
isolation.level=read_committedon consumers: they skip records from aborted transactions and wait for open ones to finish.- Zombie fencing with a stable
transactional.id: an old instance that comes back after being replaced cannot commit.
What EOS does not cover
| In scope | Out of scope |
|---|---|
| Kafka topic to Kafka topic (Kafka Streams, transactional consumers and producers) | writes to databases, files, APIs or email |
| Offsets committed in the same transaction | side effects inside your processing code (an HTTP call, a log line) |
Consumers with read_committed |
consumers left on the default read_uncommitted |
Kafka Connect sink connectors and stream processors such as Spark or Flink reach exactly-once to external systems by combining their own checkpoints with idempotent or transactional sinks, not through Kafka transactions alone.
Pitfalls
- “We enabled idempotence, so we have exactly-once.” Idempotence covers producer retries only.
- Forgetting
read_committeddownstream. Consumers default toread_uncommittedand see aborted records. - Expecting zero cost. Transactions add commit latency and broker work; very small transactions hurt throughput.
In interviews
Be precise about the scope. A strong answer names the four building blocks, says EOS is for Kafka-to-Kafka read-process-write, and explains how to extend it to external sinks with idempotence or offsets stored in the sink.
Transactions in Kafka
What it is. A Kafka transaction groups writes to any number of partitions, plus consumer offset commits, into one atomic unit.
How it works
- The producer sets a
transactional.id. OninitTransactions(), the transaction coordinator (a broker, chosen by hashing the ID onto__transaction_state) assigns a producer ID and bumps its epoch, fencing any older producer with the same ID. - Records written inside the transaction go into the normal partition logs, flagged as transactional.
commitTransaction()orabortTransaction()makes the coordinator write a COMMIT or ABORT control marker into every partition the transaction touched.read_committedconsumers read only up to the last stable offset (the first offset of any still-open transaction) and filter out aborted records.- Open transactions that are never finished are aborted after
transaction.timeout.ms(60 seconds by default; brokers cap it withtransaction.max.timeout.ms, 15 minutes). - Kafka 4.0 enabled a strengthened transaction protocol (KIP-890, the
transaction.version=2feature): 4.x producers bump the epoch on every transaction, so a stray retry cannot leak into the next transaction.
Worked example: commit, abort, read
static KafkaProducer<String, String> txProducer(String txId) {
Properties p = new Properties();
p.put("bootstrap.servers", "localhost:19092");
p.put("key.serializer", StringSerializer.class.getName());
p.put("value.serializer", StringSerializer.class.getName());
p.put("transactional.id", txId); // implies enable.idempotence=true and acks=all
return new KafkaProducer<>(p);
}
try (KafkaProducer<String, String> p = txProducer("transfer-writer-1")) {
p.initTransactions();
p.beginTransaction();
p.send(new ProducerRecord<>("transfers", "t1", "debit A 100"));
p.send(new ProducerRecord<>("transfers", "t1", "credit B 100"));
p.commitTransaction();
p.beginTransaction();
p.send(new ProducerRecord<>("transfers", "t2", "debit A 999"));
p.flush(); // the record is in the log...
p.abortTransaction(); // ...but marked aborted
p.beginTransaction();
p.send(new ProducerRecord<>("transfers", "t3", "debit C 5"));
p.send(new ProducerRecord<>("transfers", "t3", "credit D 5"));
p.commitTransaction();
}
Reading the topic from the beginning with each isolation level (values shown as offset:value):
read_uncommitted: [0:debit A 100, 1:credit B 100, 3:debit A 999, 5:debit C 5, 6:credit D 5]
read_committed: [0:debit A 100, 1:credit B 100, 5:debit C 5, 6:credit D 5]
The aborted debit A 999 is physically in the log at offset 3 and a read_uncommitted consumer sees it. Offsets 2, 4 and 7 are missing from both lists: they hold the control markers, which kafka-dump-log.sh shows (columns trimmed):
baseOffset: 0 lastOffset: 1 producerId: 10000 isTransactional: true isControl: false
| offset: 0 ... key: t1 payload: debit A 100
| offset: 1 ... key: t1 payload: credit B 100
baseOffset: 2 lastOffset: 2 producerId: 10000 isTransactional: true isControl: true
| offset: 2 ... endTxnMarker: COMMIT coordinatorEpoch: 0
baseOffset: 3 lastOffset: 3 producerId: 10000 isTransactional: true isControl: false
| offset: 3 ... key: t2 payload: debit A 999
baseOffset: 4 lastOffset: 4 producerId: 10000 isTransactional: true isControl: true
| offset: 4 ... endTxnMarker: ABORT coordinatorEpoch: 0
baseOffset: 5 lastOffset: 6 producerId: 10000 isTransactional: true isControl: false
| offset: 5 ... key: t3 payload: debit C 5
| offset: 6 ... key: t3 payload: credit D 5
baseOffset: 7 lastOffset: 7 producerId: 10000 isTransactional: true isControl: true
| offset: 7 ... endTxnMarker: COMMIT coordinatorEpoch: 0
Zombie fencing
A “zombie” is an old instance that was presumed dead (a long GC pause, a network partition) but is still running. When its replacement calls initTransactions() with the same transactional.id, the epoch is bumped, the zombie’s open transaction is aborted, and the zombie can no longer commit:
KafkaProducer<String, String> zombie = txProducer("transfer-writer-2");
zombie.initTransactions();
zombie.beginTransaction();
zombie.send(new ProducerRecord<>("transfers", "t4", "debit E 1")).get();
KafkaProducer<String, String> replacement = txProducer("transfer-writer-2");
replacement.initTransactions(); // bumps the epoch, aborts the zombie's open transaction
zombie.commitTransaction(); // throws
zombie commit: ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one.
The zombie’s debit E 1 stayed in the log with an ABORT marker after it, and the coordinator’s view confirms both IDs:
$ bin/kafka-transactions.sh --bootstrap-server localhost:19092 list
TransactionalId Coordinator ProducerId TransactionState
transfer-writer-1 2 10000 CompleteCommit
transfer-writer-2 3 11000 Empty
Pitfalls
- Random
transactional.idper start disables fencing: the new instance cannot fence the old one. Use a stable ID per logical writer (for example per input partition or per pod ordinal). - Treating
ProducerFencedExceptionas retriable. It is fatal for that producer instance: close it and stop. - Long-running transactions hold back the last stable offset, so
read_committedconsumers of those partitions stall until the transaction ends or times out.kafka-transactions.shcan list and, as a last resort, abort hanging transactions. - Lag never quite reaches zero on transactional topics in tools that count offsets, because control markers and aborted records occupy offsets.
In interviews
Describe the coordinator, the producer ID and epoch, control markers, the last stable offset and read_committed. Explain fencing with a concrete zombie story. Mention transaction.timeout.ms and what a hanging transaction does to consumers.
Idempotence end to end
What it is. “End to end” means the whole path from the source to the final sink tolerates retries and replays without changing the result. Kafka’s idempotent producer is only one link.
The chain
| Link | Threat | Protection |
|---|---|---|
| Source to producer | the application resends after a crash | idempotency key in the event (event ID), outbox pattern, or CDC from the database log |
| Producer to Kafka | network retries | idempotent producer (default since 3.0) |
| Kafka to Kafka | crashes in a processor | transactions with read_committed, or Kafka Streams with exactly_once_v2 |
| Kafka to external sink | reprocessing after a crash or replay | upsert by event ID, or store offsets atomically with the data |
Worked example: offsets stored in the sink
If the sink is a transactional database, you can write the results and the consumer’s next offset in one database transaction, and seek to that stored offset on start-up instead of using Kafka’s committed offsets. A crash can then never separate “data written” from “offset recorded”. The simulation below uses SQLite as the sink and crashes once while writing the second batch:
import sqlite3
# The "topic": one partition of payment events (offset -> event)
events = [{"payment_id": f"p{i}", "amount": 10 * i} for i in range(1, 9)]
db = sqlite3.connect(":memory:")
db.executescript("""
CREATE TABLE payments (payment_id TEXT PRIMARY KEY, amount INTEGER);
CREATE TABLE kafka_offsets (topic TEXT, part INTEGER, next_offset INTEGER, PRIMARY KEY (topic, part));
INSERT INTO kafka_offsets VALUES ('payments', 0, 0);
""")
def next_offset():
return db.execute("SELECT next_offset FROM kafka_offsets WHERE topic='payments' AND part=0").fetchone()[0]
def process(batch_size, crash_at=None):
"""Read from the offset stored in the sink; write rows and the new offset in ONE database transaction."""
pos = next_offset()
while pos < len(events):
batch = events[pos:pos + batch_size]
try:
with db: # BEGIN ... COMMIT, or ROLLBACK on exception
for i, e in enumerate(batch):
if crash_at is not None and pos + i == crash_at:
raise RuntimeError(f"crash while writing offset {crash_at}")
db.execute("INSERT INTO payments VALUES (?, ?) "
"ON CONFLICT (payment_id) DO UPDATE SET amount = excluded.amount",
(e["payment_id"], e["amount"]))
db.execute("UPDATE kafka_offsets SET next_offset = ? WHERE topic='payments' AND part=0",
(pos + len(batch),))
except RuntimeError as err:
print(" ", err, "-> transaction rolled back")
return
pos += len(batch)
print("run 1 (crashes mid-batch):")
process(batch_size=3, crash_at=4)
print(" rows:", db.execute("SELECT COUNT(*) FROM payments").fetchone()[0], "| stored next offset:", next_offset())
print("run 2 (restart, seek to stored offset):")
process(batch_size=3)
rows, total = db.execute("SELECT COUNT(*), SUM(amount) FROM payments").fetchone()
print(" rows:", rows, "| sum:", total, "| stored next offset:", next_offset())
run 1 (crashes mid-batch):
crash while writing offset 4 -> transaction rolled back
rows: 3 | stored next offset: 3
run 2 (restart, seek to stored offset):
rows: 8 | sum: 360 | stored next offset: 8
The half-written batch was rolled back together with its offset, and the restart resumed at offset 3. All eight payments appear once (10 + 20 + … + 80 = 360). In a real consumer you would call seek(partition, stored_offset) in onPartitionsAssigned. The upsert adds a second layer of protection in case the same event is replayed into the topic.
Pitfalls
- Exactly-once at the source is often ignored. If the service that produces events can send the same event twice after its own crash, every downstream guarantee only preserves that duplicate faithfully. Carry an event ID from the origin.
- Storing offsets in the sink but also auto-committing to Kafka creates two sources of truth. Pick one.
In interviews
Walk the chain link by link. Interviewers often ask “the consumer writes to Postgres; how do you avoid duplicates?” Strong answers: upsert by event ID, or store offsets in the same database transaction and seek on assignment.
Read-process-write pattern
What it is. The read-process-write (consume-transform-produce) loop reads from input topics, transforms records, writes to output topics, and commits the input offsets inside the same transaction as the outputs. It is the pattern Kafka’s exactly-once was designed for, and what Kafka Streams does internally.
The loop
// consumer: enable.auto.commit=false, isolation.level=read_committed
// producer: transactional.id=transfer-auditor-0 (stable per instance)
producer.initTransactions(); // fences older instances, aborts their open txn
consumer.subscribe(List.of("transfers"));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
if (records.isEmpty()) continue;
producer.beginTransaction();
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (ConsumerRecord<String, String> r : records) {
producer.send(new ProducerRecord<>("transfer-audit", r.key(), "audited: " + r.value()));
offsets.put(new TopicPartition(r.topic(), r.partition()), new OffsetAndMetadata(r.offset() + 1));
}
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction(); // outputs and input offsets: all or nothing
}
sendOffsetsToTransaction with consumer.groupMetadata() writes the offsets to __consumer_offsets as part of the transaction, and lets the coordinator fence a zombie that lost its partitions in a rebalance.
Worked example: a crash before commit
Instance 1 processed the four committed transfers, sent four audit records and the offsets, then “crashed” before commitTransaction(). Instance 2 started with the same transactional.id, which aborted instance 1’s open transaction, and processed the same input again:
instance 1: processed 4 records, offsets [7]
instance 1: crash before commitTransaction()
instance 2: processed 4 records, offsets [7]
instance 2: committed
audit records seen with read_uncommitted: 8
audit records seen with read_committed: 4
The output topic physically holds eight audit records, but a read_committed consumer sees exactly four, and the group’s committed offset is 7 (the offset after the last data record). Note that the input itself was read with read_committed, so the aborted debit A 999 was never processed.
Pitfalls
- Committing offsets with
consumer.commitSync()in this loop breaks atomicity; offsets must go throughsendOffsetsToTransaction. - One transaction per record is correct but slow. Commit one transaction per poll, or per time interval, so each commit covers many records.
- Side effects in “process”. A call to an external API inside the loop runs again on retry; transactions cannot undo it.
- Error handling. On
ProducerFencedException(and other fatal errors) close the producer; on abortable errors callabortTransaction(), then seek the consumer back to the last committed offsets before retrying, because its position has moved past the aborted batch.
In interviews
Write the loop from memory: initTransactions, beginTransaction, sends, sendOffsetsToTransaction(offsets, groupMetadata), commitTransaction, with read_committed and auto-commit off. Then say that Kafka Streams gives you this with one setting, processing.guarantee=exactly_once_v2 (see Kafka Streams).
Practice questions
A consumer writes each record to a database and then commits the offset. What guarantee is that, and how do you make it exactly-once in effect?
At-least-once: a crash after the database write and before the commit causes reprocessing. Make the effect exactly-once by upserting with a unique event ID, or by storing the consumer offset in the same database transaction as the data and seeking to it when partitions are assigned.
Your team enabled transactions on the producer, but a downstream consumer still sees records from failed jobs. Why?
The consumer is using the default isolation.level=read_uncommitted, so it reads records from aborted transactions. Set isolation.level=read_committed, which skips aborted records and only reads up to the last stable offset.
What is zombie fencing and how does transactional.id enable it?
A zombie is an old instance still running after being replaced. Each initTransactions() with a given transactional.id bumps the producer epoch at the transaction coordinator and aborts any open transaction for that ID; requests from the older epoch are rejected with ProducerFencedException. A stable ID per logical writer is required; a random ID per start would let zombies keep committing.
Why can a read_committed consumer stall even though new data is arriving?
It can only read up to the last stable offset, the start of the earliest open transaction in that partition. A producer that began a transaction and has not committed or aborted (for example because it hung) blocks the partition for read_committed readers until it finishes or transaction.timeout.ms aborts it.
Is exactly-once possible when the sink is an external REST API?
Not through Kafka transactions. You need the API to support idempotency keys (send the event ID so the server ignores repeats), or a deduplication store you check before calling it. Otherwise the best you can do is at-least-once with duplicates or at-most-once with loss.
What does sendOffsetsToTransaction do that commitSync does not?
It writes the consumer group’s offsets as part of the producer’s transaction, so the input offsets become committed if and only if the output records are committed. Passing consumer.groupMetadata() also lets the coordinator reject commits from a consumer that has lost its partitions in a rebalance. commitSync() commits independently of the transaction, so a crash could commit offsets for outputs that were then aborted.
Key takeaways
- At-most-once commits before processing and can lose data; at-least-once commits after and can duplicate. Most pipelines choose at-least-once with idempotent sinks.
- Kafka’s exactly-once covers Kafka-to-Kafka read-process-write, built from the idempotent producer, transactions,
read_committedand fencing. - Transactions write COMMIT or ABORT markers into every touched partition;
read_committedconsumers read up to the last stable offset and skip aborted records. - A stable
transactional.idfences zombies;ProducerFencedExceptionis fatal for that instance. - For external sinks, use upserts keyed by an event ID or store offsets in the same sink transaction.
- The read-process-write loop commits outputs and input offsets together with
sendOffsetsToTransaction.
Progress is saved in this browser only. No account needed.

