Repository navigation
Conversation
Merged
1 of 11 tasks
lhotari
added this pull request to stack #26718
September 25, 2026 19:28
lhotari
marked this pull request as draft
September 25, 2026 19:29
lhotari
marked this pull request as ready for review
September 25, 2026 19:30
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 25, 2026 22:22
697d172 to
cd9726b
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 11:33
cd9726b to
c05fae5
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 26, 2026 12:32
ab178a9 to
55dcef5
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 17:19
55dcef5 to
947eec3
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 18:18
947eec3 to
51f887a
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 18:36
51f887a to
c53384f
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 20:18
c53384f to
296523e
Compare
lhotari
removed this pull request from stack #26718
September 26, 2026 20:22
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 26, 2026 20:22
296523e to
04efe8e
Compare
lhotari
changed the base branch from
lh-key-shared-500x20-scenario
to
lh-ml-mpsc-add-handoff
September 26, 2026 20:22
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 27, 2026 12:30
2e0b911 to
598bd52
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 27, 2026 12:57
598bd52 to
542f95e
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 27, 2026 13:45
542f95e to
d107cf5
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 27, 2026 17:18
e3a6f09 to
3754d1e
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 27, 2026 21:51
3754d1e to
726af4f
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 27, 2026 22:02
726af4f to
473c841
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 28, 2026 13:28
48a7983 to
bd7075d
Compare
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
2 times, most recently
from
September 28, 2026 15:07
941d5a3 to
9ae5a9b
Compare
…e dispatcher monitor Motivation Each Flow command on a Shared or Key_Shared subscription scheduled a task on the broker executor, and the task took the dispatcher monitor only to add the permits to the dispatcher total and request a read, which readMoreEntriesAsync already coalesces. Under a 500-producer, 20-consumer Key_Shared load, off-CPU profiling with jonoffcpu shows these tasks as the largest blocked time with an application frame: executor threads waiting for the dispatcher monitor in internalConsumerFlow. Modifications - Consumer tracks whether the dispatcher accounts its Flow updates, from addConsumer until removal, under its existing flow permit accounting lock. - PersistentDispatcherMultipleConsumers.consumerFlow completes the pending Flow update on the calling thread: under the consumer's accounting lock it adds the permits to the total atomically while the consumer is accounted, so the update is ordered with removal, then requests a coalesced read. internalConsumerFlow, its executor task and the membership check under the monitor are gone. - Removal stops the consumer's accounting and debits its balance with an atomic update, as the other writers of the total already do. - SharedDispatcherPermitAccountingTest defers Flow updates through a test subscription instead of holding the dispatcher monitor, and checks that a Flow no longer waits for the monitor. Assisted-by: Claude Code (claude-opus-5-5)
lhotari
force-pushed
the
lh-dispatcher-lockfree-flow-master
branch
from
September 28, 2026 23:03
9ae5a9b to
2375de6
Compare
This was referenced Sep 29, 2026
lhotari
marked this pull request as draft
September 30, 2026 21:07
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
Each Flow command on a Shared or Key_Shared subscription schedules a task on the broker executor. The task takes the dispatcher monitor only to add the permits to the dispatcher's total and request a read, and
readMoreEntriesAsyncalready coalesces read requests. While a dispatch or a read completion holds the monitor, these tasks block broker executor threads.Off-CPU profiling of 500 producers and a 20-member Key_Shared subscription shows these tasks as the largest blocked time with an application frame in the broker: executor threads waiting for the dispatcher monitor in
PersistentDispatcherMultipleConsumers.internalConsumerFlow(2.5 s observed, about 10.7 s estimated per run).Modifications
Consumertracks whether the dispatcher accounts its Flow updates, fromaddConsumeruntil removal, under its existing flow permit accounting lock.PersistentDispatcherMultipleConsumers.consumerFlowcompletes the pending Flow update on the calling thread. Under the consumer's accounting lock it adds the permits to the total atomically while the consumer is accounted, so the update is ordered with the consumer's removal, and then requests a coalesced read.internalConsumerFlow, its executor task and the membership check under the dispatcher monitor are removed.SharedDispatcherPermitAccountingTestdefers Flow updates through a test subscription instead of holding the dispatcher monitor, and checks that a Flow no longer waits for the monitor.Measurements
The Key_Shared 500×20 (
iot-key-shared-500x20.yaml) and IoT telemetry high-rate (iot-telemetry-high-rate.yaml) performance scenarios oftests/performance, run with the performance launcher on one host (Intel i9-9980HK, 8 cores, 16 hardware threads) running the broker, 3 bookies and the clients, on JDK 25.0.4.1 with ZGC. The host was configured with the performance testing environment setup oftests/performance/environment, whose TuneD profile disables turbo, so that the CPU ran at a fixed 2.4 GHz in every run. Baseline is the parent of this change (#26717,d7fe05346c4f, which hands adds off to the managed-ledger executor in batches), candidate is this change (3754d1e860a3). The measurements come from one rotation over the builds of the stack, with the builds interleaved run by run: three unprofiled Key_Shared runs, two unprofiled high-rate runs and one profiled Key_Shared run (iot-key-shared-500x20-profile.yaml) of each build.thermaldwas stopped, and the launcher let the CPU package cool down to 55 °C before each run and again before each measurement (-Pperformance.cooldownTemperature=55), so every run started at 48–52 °C. None of the runs throttled thermally. The metrics stack oftests/performancescraped the broker, the bookies and ZooKeeper every 5 s during each run, with the broker's stats periods set to the same 5 s.Every run delivered all messages to every application with no duplicates, no ordering violations and no invalid messages.
Key_Shared 500×20, unprofiled:
IoT telemetry high rate, unprofiled (end-to-end latency and backlog are the ranges over the 5 applications, and their change is that of the mean over the applications):
Broker, profiled Key_Shared run (JFR recording, flame graphs and off-CPU profile of the measurement period, 5,000,000 messages on each side):
internalConsumerFlowmonitor waits (off-CPU, observed)trySendMessagesToConsumers(Key_Shared)PulsarCommandSenderImpl.sendMessagesToConsumerEpollIoHandler.wakeup(eventfd writes)In the baseline, executor threads waiting for the dispatcher monitor in
internalConsumerFloware 73 % of the broker's blocked time with an application frame; with this change that wait is gone, and the blocked time with an application frame drops from 4.0 s to 0.6 s.The contention is gone, but on its own the change lowers the Key_Shared producer throughput by 5.9 % in this CPU-saturated benchmark, with every run of this change slower than every baseline run. The profiled broker uses 7.5 % more CPU per message, and the extra CPU is spent handing messages to the consumer connections and waking up event loops. It isn't a change in read batching: counters added to the dispatcher in an experiment show that on both sides every normal read returns the maximum of 100 entries (
dispatcherMaxReadBatchSize), since the dispatcher runs behind the tail, and each read sends about 5 entries to each of the 20 consumers. When the dispatcher hands a consumer's messages to the consumer connection's event loop (PulsarCommandSenderImpl.sendMessagesToConsumer), Netty writes to the loop's eventfd if the loop is asleep inepoll_wait. In the profiled runs, which moved the same 5,000,000 messages, the CPU samples insendMessagesToConsumerrise from 90 to 291 and those inEpollIoHandler.wakeupfrom 336 to 615. In the baseline, the Flow tasks ran on the samepulsar-ioevent loops (BrokerService.executor()) and blocked them on the dispatcher monitor, so a loop was less often asleep when the dispatcher handed it work; with this change the loops no longer block, sleep between their own events, and more hand-offs have to wake them.The high-rate scenario publishes at a fixed 30,000 messages per second, which this host sustains on both sides, so it shows the latency at a rate below saturation. Both sides deliver every message at that rate with the same publish p50 of 2.4–2.5 ms and end-to-end p50 of 9 ms, but this change's end-to-end p99 is higher, 31–35 ms against 23–27 ms, with both of its runs above both baseline runs, consistent with the extra CPU per message.
Together with #26717 below it, the stack reaches a mean of 91,165 msg/s in the Key_Shared scenario in the same rotation, +5.3 % over #26715.
Charts of the median unprofiled Key_Shared run of each side (baseline 95,797 msg/s, this change 91,425 msg/s), as the run report renders them:
Verifying this change
This change added tests and can be verified as follows:
SharedDispatcherPermitAccountingTest(16 cases) covers Flow updates racing consumer removal, a replacement consumer with the same identity, blocked-consumer permits, negative removal balances, and a new test that a Flow command updates the dispatcher total while another thread holds the dispatcher monitor.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
Threading model: Shared and Key_Shared Flow permits are applied on the connection thread that received the Flow command, without the dispatcher monitor and without a broker executor task. Reads are still scheduled through the existing coalescing
readMoreEntriesAsync.This PR was prepared with AI assistance (Claude Code) and reviewed by a human contributor.