Most explanations of Kafka exactly-once stop at the client: turn on idempotence, give the producer a transactional.id, wrap the work in beginTransaction and commitTransaction, and read with isolation.level=read_committed. That is the right mental model for writing an application, and this site covers it in the exactly-once semantics article. It is the wrong model for operating one, because every guarantee is actually enforced by brokers, and when it goes wrong the symptom is broker-side: a consumer group that stops advancing while the topic keeps growing.
This article is about that side. It walks through the state each broker keeps, the protocol the transaction coordinator runs, how control markers and the last stable offset decide what a read_committed consumer may see, why transactions can hang, and how KIP-890 closes the gap. It ends with a probe you can run against a cluster today and a triage runbook built on the kafka-transactions.sh tool from KIP-664.
Producer state: IDs, epochs and sequence numbers
When an idempotent producer starts, it asks the cluster for a producer ID (PID) and an epoch. Every batch it sends carries the PID, the epoch and a base sequence number that counts per partition from zero. The partition leader keeps a small table keyed by PID: the current epoch and metadata for the last few batches appended, five in current versions, which is why idempotence requires max.in.flight.requests.per.connection of five or fewer.
On each produce request the leader checks three things. If the epoch is older than the one it knows, the writer is a zombie and the request is fenced. If the sequence matches one of the remembered batches, the request is a retry of something already written, so the leader acknowledges it with the original offset and writes nothing. If the sequence is not exactly the next expected number, there is a gap, and the leader rejects the batch with an out-of-order sequence error rather than leaving a hole.
The table survives leader changes because each batch header carries PID, epoch and sequence, so it can be rebuilt from the log, and it is periodically written to producer state snapshot files beside the segments so a restart does not rescan everything. Entries for producers idle longer than producer.id.expiration.ms (one day by default) are forgotten.
The transaction coordinator and its log
Transactions add a second actor. A producer with a transactional.id is assigned a coordinator by hashing that ID onto a partition of the internal topic __transaction_state; the leader of that partition is the coordinator. The topic defaults to 50 partitions with replication factor 3 and a minimum in-sync replica count of 2, and those defaults are worth keeping: the transaction log is the source of truth for every commit decision in the cluster.
The coordinator keeps one small state record per transactional ID: PID, epoch, timeout, the set of partitions in the current transaction, and a state. The states are Empty, Ongoing, PrepareCommit, PrepareAbort, CompleteCommit, CompleteAbort, plus Dead and PrepareEpochFence for expiry and fencing. Every transition is replicated before it is acted on, so a new coordinator can finish whatever its predecessor started.
One commit, request by request
| Step | Request | What the broker does |
|---|---|---|
| 1 | InitProducerId | Coordinator finds or creates the state for the transactional ID, aborts any transaction the previous instance left open, bumps the epoch and returns PID and epoch. The epoch bump is what fences the old instance. |
| 2 | AddPartitionsToTxn | Coordinator adds the partitions to the transaction's set and moves Empty to Ongoing. Under KIP-890 part 2 this is implicit on the first produce. |
| 3 | Produce | Partition leaders append transactional batches. The first such batch fixes the transaction's first offset on that partition. |
| 4 | AddOffsetsToTxn, TxnOffsetCommit | The consumer group's offsets partition joins the transaction; the group coordinator holds the offsets as pending. |
| 5 | EndTxn | Coordinator writes PrepareCommit (or PrepareAbort) to its log. Once that record is replicated the outcome is decided, whatever crashes next. |
| 6 | WriteTxnMarkers | Coordinator sends a COMMIT or ABORT control record to every partition in the set, then writes CompleteCommit and returns the transaction to Empty. |
Two facts follow. First, the commit point is step 5, not step 6: a coordinator that crashes between them is replaced by one that reads PrepareCommit and resends the markers. Second, markers reach partitions independently, so one partition can briefly expose a commit another has not; exactly-once is not a cross-partition snapshot.
Control markers, the LSO and the aborted index
A partition log with transactions contains three kinds of record: ordinary data, transactional data, and control markers. Transactional data is written immediately, before anyone knows whether it will commit, so readers need to know where certainty ends. That boundary is the last stable offset (LSO): the first offset belonging to a transaction that is still open, or the high watermark if none is open.
A read_committed fetch is served only up to the LSO. With the data, the leader returns the list of aborted transactions that overlap the range, looked up in a per-segment aborted-transaction index (the .txnindex file). The consumer then drops batches from those producers between each transaction's first offset and its abort marker. Committed data needs no list, and the markers themselves are skipped. read_uncommitted consumers read up to the high watermark and see everything, aborted data included.
The consequence operators feel is that one open transaction holds back every reader of that partition, including data from other producers written after it. A producer that opens a transaction, writes one record and then stalls for a minute delays every read_committed consumer of that partition by a minute.
Worked example: a timeline on one partition
Two producers share partition orders-0. Producer P (PID 4021) runs transactions; producer Q is a plain idempotent producer.
| Offset | Record | LSO after it |
|---|---|---|
| 100 | Q data | 101 |
| 101 | P txn data (transaction T1 begins) | 101 |
| 102 | Q data | 101 |
| 103 | P txn data | 101 |
| 104 | Q data | 101 |
| 105 | ABORT marker for T1 | 106 |
| 106 | P txn data (T2 begins) | 106 |
| 107 | COMMIT marker for T2 | 108 |
Before offset 105 a read_committed consumer is stuck at 101, even though Q's records 102 and 104 are final. When the abort lands the LSO jumps to 106; the fetch for 101 to 105 returns all five records plus the aborted list entry (PID 4021, first offset 101), and the client discards 101 and 103. The consumer delivers 100, 102 and 104, then 106 after the commit; a read_uncommitted consumer would also have delivered the aborted 101 and 103.
Timeouts, fencing and expiry
Three settings bound how long a transaction can hold the LSO. The producer's transaction.timeout.ms (60 seconds by default) is sent to the coordinator on InitProducerId and must not exceed the broker's transaction.max.timeout.ms (15 minutes by default). When a transaction stays Ongoing past its timeout, the coordinator aborts it itself and bumps the epoch, so the producer's next transactional request fails; whether the client can recover by aborting or must be recreated depends on client and protocol version. Finally, transactional.id.expiration.ms (7 days by default) removes the state of IDs that have not been used, which matters for applications that generate IDs dynamically.
Set the producer timeout just above the longest measured gap between beginTransaction and commitTransaction; every extra second is a second a stuck producer can freeze downstream readers.
Hanging transactions and KIP-890
A hanging transaction is one the partition believes is open but the coordinator does not. The partition never receives a marker, so its LSO never advances and read_committed consumers stall indefinitely. Before KIP-890 the classic path was a race: a produce request delayed in the network reached the leader after the coordinator had already aborted the transaction on timeout and written the marker. The late batch started a new open transaction on that partition that no coordinator knew about.
KIP-890 part 1 added server-side verification: before appending a transactional batch, the leader checks with the coordinator that the partition really is part of an ongoing transaction for that producer and epoch, controlled by transaction.partition.verification.enable. Part 2 changes the protocol so the epoch is bumped at the end of every transaction and returned in the EndTxn response, which makes a late write from the previous transaction detectable as stale, and lets clients add partitions implicitly on first produce. The Kafka 4.0 release announcement describes part 2 as completed in that release. Whether your cluster runs the new protocol depends on its transaction version feature level and on client versions, so confirm both in the upgrade notes for the release you run rather than assuming it is on.
Measuring the LSO gap
The cheapest early warning is the difference between a partition's high watermark and its LSO. Java's consumer exposes both: endOffsets() returns the high watermark for a read_uncommitted consumer and the LSO for a read_committed one.
// Compare log end (high watermark) with the last stable offset for every partition of a topic.
// endOffsets() returns the LSO when the consumer is configured read_committed.
static Map<TopicPartition, Long> lsoGap(String bootstrap, String topic) {
Properties base = new Properties();
base.put("bootstrap.servers", bootstrap);
base.put("key.deserializer", ByteArrayDeserializer.class.getName());
base.put("value.deserializer", ByteArrayDeserializer.class.getName());
Properties committed = (Properties) base.clone();
committed.put("isolation.level", "read_committed");
try (var hw = new KafkaConsumer<byte[], byte[]>(base);
var lso = new KafkaConsumer<byte[], byte[]>(committed)) {
List<TopicPartition> parts = hw.partitionsFor(topic).stream()
.map(pi -> new TopicPartition(topic, pi.partition())).toList();
Map<TopicPartition, Long> high = hw.endOffsets(parts);
Map<TopicPartition, Long> stable = lso.endOffsets(parts);
Map<TopicPartition, Long> gap = new TreeMap<>(Comparator.comparingInt(TopicPartition::partition));
for (TopicPartition tp : parts) gap.put(tp, high.get(tp) - stable.get(tp));
return gap;
}
}Alert on a partition whose LSO has not changed for longer than the broker's maximum transaction timeout while its high watermark has. A non-zero but moving gap is normal; it just means transactions are open. Charting consumer lag separately for read_committed and read_uncommitted groups helps too: a diverging pair points at transactions rather than slow consumers.
Triage runbook for a stuck partition
# 1. Which partitions are stuck, and from which offset? (run against the partition's leader broker)
kafka-transactions.sh --bootstrap-server b1:9092 --find-hanging --broker 3
# 2. Who is the producer behind the open transaction on that partition?
kafka-transactions.sh --bootstrap-server b1:9092 --describe-producers --topic orders --partition 0
# 3. If the producer has a transactional.id, ask its coordinator what state it believes it is in
kafka-transactions.sh --bootstrap-server b1:9092 --describe --transactional-id tx-7
# 4. Last resort, after confirming the owning application is stopped or fenced:
# abort the open transaction that starts at the stuck offset
kafka-transactions.sh --bootstrap-server b1:9092 --abort --topic orders --partition 0 --start-offset 550The commands above are illustrative; flags follow KIP-664, so check --help in your version. Work through them in order. If find-hanging names a partition, describe-producers shows the PID, epoch and start offset of the open transaction, and describe on the transactional ID shows whether the coordinator still thinks it is Ongoing. If the coordinator says Ongoing, the producer is alive or about to time out: fix the application. If the coordinator has no record of that partition in a live transaction, it is a true hang and an abort is justified.
Aborting discards the transaction's data for read_committed consumers, so only do it once the owning application is stopped or fenced. Record the PID, offsets and producer version first.
Failure modes
- Consumers frozen, topic growing. An open or hanging transaction pins the LSO. Measure the gap, then follow the runbook.
- Fencing errors after a pause. The coordinator aborted on timeout and bumped the epoch, or another instance started with the same transactional ID. Handle the error your client version documents, and if in doubt close the producer, create a new one, and reprocess from the last committed offsets.
- Under-replicated transaction log. With fewer than the minimum in-sync replicas on a
__transaction_statepartition, the coordinator cannot write state, and every transactional ID hashed there stops. Monitor it like any critical internal topic. - Duplicates with exactly-once on. Usually a read_uncommitted consumer, an external side effect outside the transaction, or offsets committed with the plain consumer API instead of
sendOffsetsToTransaction.
Operational guidance and trade-offs
Every transaction costs replicated coordinator writes plus one marker per partition, so the commit interval is the main knob: short intervals keep latency low but multiply markers, long ones amortise them but delay every read_committed consumer. Keep transactions short and few-partition, keep transactional IDs stable per logical task so restarts fence their predecessors, and treat any stuck LSO as an incident rather than something to abort and forget.
For the client-side programming model see exactly-once semantics in streaming. For the replication the transaction log depends on see Kafka ISR replication, for offsets and the poll loop see Kafka consumers and consumer groups, for the log model see the Apache Kafka overview, and for a producer that relies on these guarantees see Debezium CDC to Kafka.
What to do next
- Deploy the LSO-gap probe for every topic with transactional producers and alert on an LSO frozen longer than
transaction.max.timeout.ms. - Check
__transaction_stateis at replication factor 3 with min in-sync replicas 2, and add it to under-replication alerts. - Set each producer's
transaction.timeout.msfrom its measured worst-case loop time, not the default. - Confirm every consumer of transactional topics sets
isolation.level=read_committed. - Find out which transaction protocol your brokers and clients run, and plan the upgrade to KIP-890 support.
- Rehearse the kafka-transactions.sh runbook on a staging cluster before you need it.