Skip to content

[fix][client] Keep transactional and non-transactional messages in separate batches - #26547

Merged
lhotari merged 2 commits into
apache:masterfrom
SongOf:fix/batch-txn-isolation
Sep 12, 2026
Merged

lhotari merged 2 commits into
apache:masterfrom
SongOf:fix/batch-txn-isolation

Conversation

@SongOf

@SongOf SongOf commented Sep 11, 2026

Copy link
Copy Markdown
Contributor

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:

// BatchMessageContainerImpl.createOpSendMsg
if (currentTxnidMostBits != -1) {
    messageMetadata.setTxnidMostBits(currentTxnidMostBits);
}

hasSameTxn, the admission check meant to enforce that invariant, cannot express it:

public boolean hasSameTxn(MessageImpl<?> msg) {
    if (!msg.getMessageBuilder().hasTxnidMostBits() || !msg.getMessageBuilder().hasTxnidLeastBits()) {
        return true;                                    // (1)
    }
    if (currentTxnidMostBits == -1 || currentTxnidLeastBits == -1) {
        currentTxnidMostBits = msg.getMessageBuilder().getTxnidMostBits();   // (2) mutates the batch
        currentTxnidLeastBits = msg.getMessageBuilder().getTxnidLeastBits();
        return true;
    }
    return currentTxnidMostBits == msg.getMessageBuilder().getTxnidMostBits()   // (3)
            && currentTxnidLeastBits == msg.getMessageBuilder().getTxnidLeastBits();
}

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 hasSameXxx query, and makes a
plain 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) and
newMessage() — legal API usage, with batching enabled by default:

  1. A plain message batched after transactional ones inherits the transaction: it stays invisible until
    commit and is silently discarded if the transaction aborts, although the application never sent
    it inside a transaction.
  2. A transactional message joining a batch that already holds plain messages stamps its transaction
    onto the whole batch, retroactively enrolling those plain messages in it.
  3. The side effect in branch (2) also fires on the duplicate-sequence-id path in ProducerImpl:
    canAddToCurrentBatch adopts the transaction, and the caller then flushes the current batch via
    doBatchSendAndAdd and puts the new message into a fresh one — so a batch of plain messages is
    published stamped with a transaction that none of its messages belongs to.

BatchMessageKeyBasedContainer fails the other way round. The guard reads the outer container's
transaction 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 canAddToCurrentBatch is already the
correct shape: hasSameSchema is a pure query that starts with numMessagesInBatch == 0 and treats
"absent" as an identity that must match.

Modifications

  • AbstractBatchMessageContainer.hasSameTxn is 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 -1
    sentinel 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.add captures 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, batchHasTxn would always be false, and plain messages
    would again be admitted into a transactional container.

Nothing changes at the call site. When the guard rejects a message, ProducerImpl.doBatchSendAndAdd
already 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

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • Added BatchMessageContainerImplTest.testPlainMessageIsNotBatchedWithTransactionalMessages and
    ...testTransactionalMessageIsNotBatchedWithPlainMessages, covering both directions, and
    ...testKeyBasedContainerKeepsTransactionalAndPlainMessagesApart for the key-based batcher. All
    three fail before the change.
  • Added BatchMessageContainerImplTest.testSameTransactionStillBatchesTogether and
    ...testPlainMessagesStillBatchTogether as guards that messages of one transaction still share a
    batch, that two different transactions still do not, and that plain messages still batch together.
  • Added TransactionProduceTest.testAbortedTransactionDoesNotDiscardPlainMessageBatchedWithIt: a
    batching 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".
  • BatchMessageContainerImplTest passes 9/9, the pulsar-client module suite shows no new failures,
    and ./gradlew quickCheck is clean.

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 12, 2026
@lhotari
lhotari merged commit fa05c3e into apache:master Sep 12, 2026
81 of 83 checks passed
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)
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants