Repository navigation
[fix][client] Keep transactional and non-transactional messages in separate batches - #26547
Merged
Merged
Conversation
lhotari
requested changes
Sep 11, 2026
lhotari
left a comment
Member
There was a problem hiding this comment.
Thanks for fixing this transaction isolation issue and adding the abort regression test. The normal batching paths look correct. There is one remaining allocation-recovery path that can still mix transactional and plain messages, plus a small cleanup issue in the new unit tests.
lhotari
approved these changes
Sep 12, 2026
lhotari
left a comment
Member
There was a problem hiding this comment.
LGTM. The follow-up resets the key-based transaction identity on every first add, including a plain message after a failed allocation, and adds the requested recovery regression test. The container tests also release their batch buffers in finally blocks.
nodece
added a commit
to nodece/pulsar
that referenced
this pull request
Sep 14, 2026
Conflict with apache#26547 (transactional/non-transactional batch isolation) in BatchMessageContainerImplTest: both sides appended tests at the end of the class. Resolution keeps both sets — apache#26547's transaction-batch tests and this PR's buffer-ownership tests — with the import block rebuilt as the sorted union. Verified: BatchMessageContainerImplTest 21/21, ProducerImplTest 23/23, spotless + checkstyle clean.
lhotari
pushed a commit
that referenced
this pull request
Sep 14, 2026
…parate batches (#26547) Co-authored-by: maxlisongsong <[email protected]> (cherry picked from commit fa05c3e)
lhotari
pushed a commit
that referenced
this pull request
Sep 14, 2026
…parate batches (#26547) Co-authored-by: maxlisongsong <[email protected]> (cherry picked from commit fa05c3e)
3 of 4 tasks
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
…parate batches (apache#26547) Co-authored-by: maxlisongsong <[email protected]>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Motivation
A batch carries a single transaction id in its metadata, so every message inside it inherits that
transaction — there is no per-message override:
hasSameTxn, the admission check meant to enforce that invariant, cannot express it:Branch (1) inspects only the message, never the batch, so a plain message is admitted into any batch
including a transactional one. Branch (2) performs a write inside a
hasSameXxxquery, and makes aplain batch adopt the transaction of the first transactional message merely offered to it. Only
branch (3) actually compares, and it is reachable only when both sides carry a transaction. The check
can answer "are these two transactions the same?" but never "is one of them absent?".
Three consequences, all reachable from a single batching producer used for both
newMessage(txn)andnewMessage()— legal API usage, with batching enabled by default:commit and is silently discarded if the transaction aborts, although the application never sent
it inside a transaction.
onto the whole batch, retroactively enrolling those plain messages in it.
ProducerImpl:canAddToCurrentBatchadopts the transaction, and the caller then flushes the current batch viadoBatchSendAndAddand puts the new message into a fresh one — so a batch of plain messages ispublished stamped with a transaction that none of its messages belongs to.
BatchMessageKeyBasedContainerfails the other way round. The guard reads the outer container'stransaction state, while the batch metadata is built by the per-key inner container from its own
first message. A transactional message landing in a key bucket whose first message was plain is
therefore published with no transaction id at all — outside the transaction, and not rolled back when
it aborts.
For reference, the sibling admission check on the same line of
canAddToCurrentBatchis already thecorrect shape:
hasSameSchemais a pure query that starts withnumMessagesInBatch == 0and treats"absent" as an identity that must match.
Modifications
AbstractBatchMessageContainer.hasSameTxnis now a pure query comparing transaction identities, with"no transaction" as an identity of its own that is incompatible with any transaction. An empty batch
(
numMessagesInBatch == 0) still accepts anything; testing the message count rather than the-1sentinel is what distinguishes "no identity yet" from "identity is no transaction". Adopting the
identity is left solely to
add, where it already happened — which also removes consequence 3 above.BatchMessageKeyBasedContainer.addcaptures the transaction id from the container's first message.This is required rather than cosmetic: with the side effect gone, the outer container's transaction
state would otherwise stay unset forever,
batchHasTxnwould always be false, and plain messageswould again be admitted into a transactional container.
Nothing changes at the call site. When the guard rejects a message,
ProducerImpl.doBatchSendAndAddalready flushes the current batch and opens a new one with that message, exactly as it does today when
the batch is full or the schema differs.
Applications that mix both kinds of send on one producer will see batches split more often. That is the
necessary cost: the two cannot share a batch without one of them inheriting the wrong transaction
identity. There is no wire format or API change, and no behaviour change for applications that do not
mix.
Verifying this change
This change added tests and can be verified as follows:
BatchMessageContainerImplTest.testPlainMessageIsNotBatchedWithTransactionalMessagesand...testTransactionalMessageIsNotBatchedWithPlainMessages, covering both directions, and...testKeyBasedContainerKeepsTransactionalAndPlainMessagesApartfor the key-based batcher. Allthree fail before the change.
BatchMessageContainerImplTest.testSameTransactionStillBatchesTogetherand...testPlainMessagesStillBatchTogetheras guards that messages of one transaction still share abatch, that two different transactions still do not, and that plain messages still batch together.
TransactionProduceTest.testAbortedTransactionDoesNotDiscardPlainMessageBatchedWithIt: abatching producer sends a transactional and a plain message into the same batch, the transaction is
aborted, and the plain message must still be readable. It fails before the change with "the plain
message was discarded by the aborted transaction".
BatchMessageContainerImplTestpasses 9/9, thepulsar-clientmodule suite shows no new failures,and
./gradlew quickCheckis clean.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes