8.3 Transactions
Transactions provide both.
3.1 What they were built for — and the naming distinction
"transactions were added to Kafka to guarantee the correctness of applications developed using Kafka Streams. In order for a stream processing application to generate correct results, each input record must be processed exactly one time, and its processing result will be reflected exactly one time, even in case of failure."
NOTE — vocabulary that matters
"Transactions is the name of the underlying MECHANISM. Exactly-once semantics or exactly-once guarantees is the BEHAVIOR of a stream processing application. Kafka Streams uses transactions to implement its exactly-once guarantees. Other stream processing frameworks, such as Spark Streaming or Flink, use DIFFERENT MECHANISMS to provide their users with exactly-once semantics."
The scope, stated up front:
"transactions in Kafka were developed specifically for stream processing applications. And therefore they were built to work with the 'consume-process-produce' pattern... the processing of each input record will be considered complete after the application's internal state has been updated AND the results were successfully produced to output topics."
3.2 Use cases
"Transactions are useful for any stream processing application where accuracy is important, and especially where stream processing includes AGGREGATION and/or JOINS."
Why filtering/mapping doesn't need them:
"If the stream processing application only performs single record transformation and filtering, there is no internal state to update, and even if duplicates were introduced, it is fairly straightforward to filter them out of the output stream. When the stream processing application aggregates several records into one, it is much more difficult to check whether a result record is wrong because some input records were counted more than once; it is impossible to correct the result without reprocessing the input."
"Financial applications are typical examples... However, because it is rather trivial to configure any Kafka Streams application to provide exactly-once guarantees, we've seen it enabled in more mundane use cases, including, for instance, CHATBOTS."
3.3 The two problems transactions solve
- · lack of heartbeat → rebalance (a few seconds)
- · the new consumer starts at THE LAST COMMITTED OFFSET
- · “all the records processed between the last committed offset and the crash will be PROCESSED AGAIN, and the results will be WRITTEN TO THE OUTPUT TOPIC AGAIN — resulting in DUPLICATES.”
- · missed heartbeats → assumed dead → partitions reassigned
- · the frozen instance “can do all that BEFORE it polls Kafka for records or sends a heartbeat and DISCOVERS THAT IT IS SUPPOSED TO BE DEAD.”
- · ► “A consumer that is dead but doesn’t know it is called a ZOMBIE… without additional guarantees, ZOMBIES CAN PRODUCE DATA TO THE OUTPUT TOPIC AND CAUSE DUPLICATE RESULTS.”
Note that these need different solutions:
- Problem 1 is an atomicity problem → solved by atomic multipartition writes.
- Problem 2 is an identity/authority problem → solved by zombie fencing with epochs.
Transactions provide both.
3.4 The mechanism: atomic multipartition writes
The insight that makes it possible:
"Exactly-once processing means that consuming, processing, and producing are done ATOMICALLY. Either the offset of the original message is committed and the result is successfully produced, or neither of these things happen."
"committing offsets AND producing results BOTH involve writing messages to partitions. However, the results are written to an output topic, and offsets are written to the
__consumer_offsetstopic. If we can open a transaction, write both messages, and commit if both were written successfully — or abort to retry if they were not — we will get the exactly-once semantics we are after."
► The trick: OFFSET COMMITS ARE JUST WRITES TO A TOPIC. So “atomically commit offsets + produce results” reduces to “atomically write to multiple partitions.”
The transactional producer
"A transactional producer is simply a Kafka producer that is configured with a
transactional.idand has been initialized usinginitTransactions()."
producer.id | transactional.id | |
|---|---|---|
| Origin | "generated automatically by Kafka brokers" | "part of the producer configuration" |
| Lifetime | New on every init | "expected to persist between restarts" |
| Purpose | dedup within one producer's life | "the main role of the transactional.id is to IDENTIFY THE SAME PRODUCER ACROSS RESTARTS" |
"Kafka brokers maintain
transactional.id→producer.idmapping, so ifinitTransactions()is called again with an existingtransactional.id, the producer will also be assigned the SAMEproducer.idinstead of a new random number."
This is exactly the hole from §2.5 Case A, now plugged. A stable transactional.id yields a stable producer.id, which makes cross-restart zombie detection possible.
Zombie fencing via epochs
"The usual way of fencing zombies — using an epoch — is used here. Kafka increments the epoch number associated with a
transactional.idwheninitTransaction()is invoked. Send, commit, and abort requests from producers with the sametransactional.idbut LOWER EPOCHS will be rejected with theFencedProducererror. The older producer will not be able to write to the output stream and will be forced toclose(), preventing the zombie from introducing duplicate records."
Same pattern as the Controller epoch in Ch. 6 §2.2: you cannot stop a frozen process from waking up — you can only make its writes Rejectable.
"In Apache Kafka 2.5 and later, there is also an option to add consumer group metadata to the transaction metadata. This metadata will also be used for fencing, which will allow producers with DIFFERENT transactional IDs to write to the same partitions while still fencing against zombie instances." (→ §3.7)
3.5 The consumer side — isolation.level
"Transactions are a producer feature for the most part... However, this isn't quite enough — records written transactionally, EVEN ONES THAT ARE PART OF TRANSACTIONS THAT WERE EVENTUALLY ABORTED, ARE WRITTEN TO PARTITIONS JUST LIKE ANY OTHER RECORDS. Consumers need to be configured with the right isolation guarantees, otherwise we won't have the exactly-once guarantees we expected."
isolation.level = read_uncommitted — the default | isolation.level = read_committed |
|---|---|
| “Will return ALL records, INCLUDING those that belong to OPEN or ABORTED transactions.” | Returns messages that were part of a successfully COMMITTED transaction, or written NONTRANSACTIONALLY. |
| ► Your exactly-once producer + a default consumer = NOT exactly-once. | Does not return messages that were part of an ABORTED transaction, or part of a transaction that is STILL OPEN. |
⚠️ Aborted records are physically in the log. Transactions don't prevent the write; they mark it. Filtering is the consumer's job, and it's off by default.
What read_committed does NOT give you
"Configuring
read_committedmode does NOT guarantee that the application will get ALL messages that are part of a specific transaction. It is possible to subscribe to only a SUBSET of topics that were part of the transaction and therefore get a subset of the messages. In addition, the application CAN'T KNOW WHEN TRANSACTIONS BEGIN OR END, or WHICH MESSAGES ARE PART OF WHICH TRANSACTION."
This limitation is the root cause of most of §3.6's failure cases. Consumers see committed records, but have no transaction boundary information whatsoever.
The Last Stable Offset (LSO) — and the latency it costs
"To guarantee that messages will be read in order,
read_committedmode will not return messages that were produced AFTER the point when the FIRST STILL-OPEN TRANSACTION BEGAN (known as the Last Stable Offset, or LSO). Those messages will be withheld until that transaction is committed or aborted by the producer, or until they reachtransaction.timeout.ms(default of 15 minutes) and are aborted by the broker.""Holding a transaction open for a long duration will introduce higher end-to-end latency by delaying consumers."
► ONE long-running transaction blocks EVERY LATER RECORD in the partition for read_committed consumers. Worst case: transaction.timeout.ms = 15 min.
read_committed consumers ALWAYS LAG BEHIND read_uncommitted ones. Transaction duration is therefore a LATENCY BUDGET, not just a correctness detail. Commit often.
The guarantee you actually get
"Our simple stream processing job will have exactly-once guarantees on its output even if the input was written NONTRANSACTIONALLY. The atomic multipartition produce guarantees that if the output records were committed to the output topic, the offset of the input records was ALSO committed for that consumer, and as a result the input records will NOT BE PROCESSED AGAIN."