Learn Labs
8. Exactly-Once Semantics

8.5 How transactions work internally

Three new producers per second is nothing — that's a modest serverless workload or a badly-written client that constructs a producer per request.

"We can use transactions by calling the APIs without understanding how they work. But having some mental model of what is going on under the hood will help us troubleshoot applications that do not behave as expected."

5.1 The algorithm

"The basic algorithm for transactions in Kafka was inspired by Chandy-Lamport snapshots, in which 'marker' control messages are sent into communication channels, and consistent state is determined based on the arrival of the marker."

The problem markers alone don't solve: "what happens if the producer crashes after only writing commit messages to a subset of the partitions?"

The answer: two-phase commit + a transaction log.

  1. Log the existence of an ongoing transaction, including the partitions involved
  2. Log the intent to commit or abort
    ► “once this is logged, We are doomed to commit or abort eventually”
  3. Write all the transaction markers to all the partitions
  4. Log the completion of the transaction

The transaction log: an internal topic called __transaction_state.

5.2 Walking the API calls

producertxn coordinator__transaction_statepartitionsinitTransactions()register ID, bump epochAddPartitionsToTxnRequestrecorded in the logsend() recordssendOffsetsToTransactionEndTransactionRequest① logs the INTENTION③ COMMIT MARKER to all④ commit completed
  • · initTransactions() is sent to a broker that becomes the transaction coordinator for this producer. “Each broker is the transaction coordinator for a SUBSET of the producers, just like each broker is the consumer group coordinator for a subset of the consumer groups.” “The transaction coordinator for each transactional ID is THE LEADER OF THE PARTITION OF THE TRANSACTION LOG the transactional ID is MAPPED TO.” It registers a new transactional ID, or increments the epoch of an existing one “in order to fence off previous producers that may have become zombies” — and “when the epoch is incremented, PENDING TRANSACTIONS WILL BE ABORTED.”
  • · ⚠ beginTransaction() “ISN’T PART OF THE PROTOCOL — it simply tells the PRODUCER that there is now a transaction in progress. THE TRANSACTION COORDINATOR ON THE BROKER SIDE IS STILL UNAWARE THAT THE TRANSACTION BEGAN.” Instead, “each time the producer detects that it is sending records to A NEW PARTITION, it will also send AddPartitionsToTxnRequest informing the broker that there is a transaction in progress and that additional partitions are part of it. This information is RECORDED IN THE TRANSACTION LOG.”
  • · sendOffsetsToTransaction(offsets, groupMetadata) — “committing offsets CAN BE DONE AT ANY TIME but MUST BE DONE BEFORE THE TRANSACTION IS COMMITTED.” It sends offsets plus the consumer group ID to the coordinator, which “will use the consumer group ID to FIND THE GROUP COORDINATOR and commit the offsets AS A CONSUMER GROUP NORMALLY WOULD.”
  • · commitTransaction() / abortTransaction() send EndTransactionRequest. After ① the intention is logged, ② “it is THE TRANSACTION COORDINATOR’S RESPONSIBILITY TO COMPLETE the commit (or abort) process” — ③ write a commit marker to all partitions involved, ④ record in the transaction log that the commit completed.
  • · Crash recovery: “if the transaction coordinator shuts down or crashes after logging the INTENTION to commit and before completing, A NEW TRANSACTION COORDINATOR WILL BE ELECTED, PICK UP THE INTENT FROM THE TRANSACTION LOG, AND COMPLETE THE PROCESS.”
  • · Timeout: “if a transaction is not committed or aborted within transaction.timeout.ms, THE TRANSACTION COORDINATOR WILL ABORT IT AUTOMATICALLY.”
Figure 8.5.15.2 Walking the API calls

Two things worth noticing:

  1. beginTransaction() is purely client-side. The broker learns about the transaction lazily, per-partition, via AddPartitionsToTxnRequest. So a transaction that never sends anything never exists on the broker.
  2. Once the intent is logged, completion is guaranteed by the coordinator, not by your process. That's what makes it a real two-phase commit rather than best-effort marker writing.

5.3 ⚠️ The producer-state memory leak — a genuine production hazard

"Each broker that receives records from transactional or idempotent producers will store the producer/transactional IDs IN MEMORY, together with related state for each of the last five batches sent by the producer: sequence numbers, offsets, and such. This state is stored for transactional.id.expiration.ms milliseconds AFTER the producer stopped being active (SEVEN DAYS BY DEFAULT). This allows the producer to resume activity without running into UNKNOWN_PRODUCER_ID errors."

"It is possible to cause something similar to a MEMORY LEAK in the broker by creating new idempotent producers or new transactional IDs at a very high rate but NEVER REUSING THEM."

The arithmetic — do this calculation for your own workload:

3 new idempotent producers per second, accumulated over one week (the default expiration):

Producer state entries1,800,000
Batch metadata entries stored (5 per producer)9,000,000
RAM≈ 5 GB

► “This can cause OUT-OF-MEMORY or SEVERE GARBAGE COLLECTION ISSUES on the broker.”

Three new producers per second is nothing — that's a modest serverless workload or a badly-written client that constructs a producer per request.

The recommendations:

PRIMARY
“architect the application to initialize A few long-lived producers when the application starts up, and then Reuse them for the lifetime of the application”
FALLBACK
“If this isn’t possible (Function as a service makes this difficult), we recommend Lowering transactional.id.expiration.ms so the IDs will expire faster, and therefore old state that will never be reused won’t take up a significant part of the broker memory”

Note the FaaS callout explicitly. Lambda-style architectures that spin up a producer per invocation are the canonical way to hit this.


On this page