Repository navigation
[fix][broker] Fix dispatcher permit accounting race with recycled batch index acks - #26809
Merged
merlimat merged 1 commit intoOct 2, 2026
Conversation
…ch index acks The Shared and Key_Shared dispatchers read EntryBatchIndexesAcks.getTotalAckedIndexCount() after Consumer.sendMessages() returned, to update the dispatcher-level TOTAL_AVAILABLE_PERMITS. By then sendMessages has handed the object to the consumer channel's event loop, which recycles it (nulling every slot) once the entries are written. When the dispatcher runs on another thread, the read can see 0 or a partial count, so the dispatcher subtracts all the messages of a partially acked batch instead of only the unacked ones, and its permit total drifts below the sum of the consumers' permits. Read the count once before calling sendMessages in the three dispatch paths (Shared, Shared chunked and Key_Shared), and use that value for the permit update and the debug log. Consumer.sendMessages already computes its own count before the handoff, so consumer-level permits were not affected.
11 tasks
void-ptr974
approved these changes
Oct 2, 2026
lhotari
pushed a commit
that referenced
this pull request
Oct 2, 2026
…ch index acks (#26809) (cherry picked from commit 1046481) [branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers (PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic), which were removed on master by #26687. The test covers both implementations.
lhotari
pushed a commit
that referenced
this pull request
Oct 2, 2026
…ch index acks (#26809) (cherry picked from commit 1046481) [branch-4.x] Also applied to the classic Shared and Key_Shared dispatchers (PersistentDispatcherMultipleConsumersClassic, PersistentStickyKeyDispatcherMultipleConsumersClassic), which were removed on master by #26687. The test covers both implementations.
dao-jun
pushed a commit
to ascentstream/pulsar
that referenced
this pull request
Oct 5, 2026
…ch index acks (apache#26809) (cherry picked from commit 1046481)
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 Shared or Key_Shared dispatcher keeps
totalAvailablePermits, which should always equal the sum of itsconsumers' permits. It uses this total to decide how much to read. On every send, the dispatcher and the consumer
each subtract the number of messages delivered. Batch-index acks are on by default. When a partially acknowledged
batch is redelivered, only its unacked messages are delivered, so both subtract
messages in batch - acked indexes.The acked-index count comes from a pooled
EntryBatchIndexesAcksobject.The dispatcher read that count after
Consumer.sendMessages()returned. By thensendMessages()has handed theobject to the consumer connection's event loop. That task writes the messages and then recycles the object, which
clears every slot. The dispatcher runs on a different thread, so nothing orders the two:
When the event loop gets there first, the dispatcher subtracts every message in the batch. Its total then falls
below the sum of the consumers' permits, by 4 in this example. The error adds up with each occurrence and is only
reset when all consumers disconnect. Until then, reads are smaller than the consumers can take. The dispatcher is
also reading, without synchronization, an object that another thread is clearing.
The same pattern is in three places:
PersistentDispatcherMultipleConsumers.trySendMessagesToConsumersPersistentDispatcherMultipleConsumers.sendChunkedMessagesToConsumers, used for any read that contains a chunkPersistentStickyKeyDispatcherMultipleConsumers.trySendMessagesToConsumersConsumer-level permits are correct, because
Consumer.sendMessages()reads the count before the handoff.Modifications
batchIndexesAcks.getTotalAckedIndexCount()into a local variable beforecalling
sendMessages(). Use it for the permit update and the debug log, so the dispatcher no longer touches theobject after the handoff. Without a race the result is unchanged, because
sendMessages()never modifies theobject.
SharedDispatcherPermitAccountingTest.testRedeliveredPartiallyAckedBatchDoesNotLosePermits. The test'sconnection stub recycles
batchIndexesAcksbeforesendMessages()returns. That is the losing interleaving, sothe race reproduces every time. The test dispatches a redelivered batch of 10 messages with 4 acked indexes
through each of the three paths. It then checks that the dispatcher total equals the sum of the consumers'
permits. Without the fix, the total is 4 permits short on every path.
Verifying this change
This change added tests and can be verified as follows:
SharedDispatcherPermitAccountingTest.testRedeliveredPartiallyAckedBatchDoesNotLosePermitscovers the Shared,Shared chunked and Key_Shared paths. It fails without the fix (for example,
expected: 194 but was: 190) andpasses with it.
PersistentDispatcherMultipleConsumersTest,PersistentStickyKeyDispatcherMultipleConsumersTest,KeySharedSubscriptionTest,BatchMessageWithBatchIndexLevelTest,BatchMessageIndexAckTestandMessageChunkingSharedTest.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes