Skip to content

[improve][broker] Apply Shared and Key_Shared Flow permits without the dispatcher monitor - #26716

Draft
lhotari wants to merge 1 commit into
masterfrom
lh-dispatcher-lockfree-flow-master
Draft

lhotari wants to merge 1 commit into
masterfrom
lh-dispatcher-lockfree-flow-master

Conversation

@lhotari

@lhotari lhotari commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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 readMoreEntriesAsync already 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

  • 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 the consumer's removal, and then requests a coalesced read. internalConsumerFlow, its executor task and the membership check under the dispatcher monitor are removed.
  • 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.

Measurements

The Key_Shared 500×20 (iot-key-shared-500x20.yaml) and IoT telemetry high-rate (iot-telemetry-high-rate.yaml) performance scenarios of tests/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 of tests/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. thermald was 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 of tests/performance scraped the broker, the bookies and ZooKeeper every 5 s during each run, with the broker's stats periods set to the same 5 s.

Setting Key_Shared 500×20 IoT telemetry high rate
Cluster 1 broker, 3 bookies 1 broker, 3 bookies
Ledger replication E=1, W=1, A=1 broker defaults
Producers 500 gateways × 1 topic 500 gateways × 1 topic
Applications 1 × 20 consumers, Key_Shared 5 × 10 consumers, Key_Shared
Messages 4,000,000 measured, 1,000,000 warmup 5,000,000 measured, 1,000,000 warmup
Payload 128 bytes, batching off 64 bytes, batching on
Rate limit none 30,000 msg/s

Every run delivered all messages to every application with no duplicates, no ordering violations and no invalid messages.

Key_Shared 500×20, unprofiled:

Baseline (#26717) This change Change
Producer throughput (msg/s) 99,102 · 95,648 · 95,797 (mean 96,849) 91,425 · 91,491 · 90,580 (mean 91,165) −5.9 %
Publish latency p50 (ms) 1,002.0 · 1,047.0 · 1,035.8 1,091.6 · 1,080.3 · 1,097.7 +6.0 %
End-to-end latency p50 (ms) 1,015.3 · 1,067.0 · 1,052.7 1,109.0 · 1,118.2 · 1,121.3 +6.8 %
End-to-end latency p99 (ms) 1,200.1 · 1,470.5 · 1,484.8 1,726.5 · 1,726.5 · 1,571.8 +20.9 %
Sampled maximum backlog 51,077 · 53,854 · 61,784 75,571 · 77,039 · 60,902 +28.1 %

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):

Baseline (#26717) This change Change
Producer throughput (msg/s) 30,000 · 30,000 (mean 30,000) 30,000 · 30,000 (mean 30,000) +0.0 %
Publish latency p50 (ms) 2.5 · 2.4 2.4 · 2.5 +0.0 %
End-to-end latency p50 (ms) 9 · 9 9 · 9 +0.0 %
End-to-end latency p99 (ms) 23 · 26–27 33–35 · 31–33 +32.8 %
Sampled maximum backlog per application 2,207–2,587 · 2,635–3,229 2,481–3,100 · 2,452–2,602 +1.7 %

Broker, profiled Key_Shared run (JFR recording, flame graphs and off-CPU profile of the measurement period, 5,000,000 messages on each side):

Baseline (#26717) This change Change
Producer throughput (msg/s) 92,187 86,512 −6.2 %
Blocked time with an application frame (off-CPU, observed) 4.00 s 0.58 s
internalConsumerFlow monitor waits (off-CPU, observed) 2.92 s none
Broker CPU samples per million messages 6,330 6,802 +7.5 %
CPU samples in trySendMessagesToConsumers (Key_Shared) 824 1,042 +26.5 %
CPU samples in PulsarCommandSenderImpl.sendMessagesToConsumer 90 291 +223 %
CPU samples in EpollIoHandler.wakeup (eventfd writes) 336 615 +83.0 %

In the baseline, executor threads waiting for the dispatcher monitor in internalConsumerFlow are 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 in epoll_wait. In the profiled runs, which moved the same 5,000,000 messages, the CPU samples in sendMessagesToConsumer rise from 90 to 291 and those in EpollIoHandler.wakeup from 336 to 615. In the baseline, the Flow tasks ran on the same pulsar-io event 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:

Baseline (#26717) This change
Latency by percentile of the baseline Latency by percentile with this change
Throughput over time of the baseline Throughput over time with this change
Backlog over time of the baseline Backlog over time with this change

Verifying this change

  • Make sure that the change passes the CI checks.

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

  • 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

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.

@lhotari
lhotari added this pull request to stack #26718 September 25, 2026 19:28
@lhotari
lhotari marked this pull request as draft September 25, 2026 19:29
@lhotari
lhotari marked this pull request as ready for review September 25, 2026 19:30
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from 697d172 to cd9726b Compare September 25, 2026 22:22
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from cd9726b to c05fae5 Compare September 26, 2026 11:33
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from ab178a9 to 55dcef5 Compare September 26, 2026 12:32
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 55dcef5 to 947eec3 Compare September 26, 2026 17:19
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 947eec3 to 51f887a Compare September 26, 2026 18:18
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 51f887a to c53384f Compare September 26, 2026 18:36
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from c53384f to 296523e Compare September 26, 2026 20:18
@lhotari
lhotari removed this pull request from stack #26718 September 26, 2026 20:22
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 296523e to 04efe8e Compare September 26, 2026 20:22
@lhotari
lhotari changed the base branch from lh-key-shared-500x20-scenario to lh-ml-mpsc-add-handoff September 26, 2026 20:22
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from 2e0b911 to 598bd52 Compare September 27, 2026 12:30
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 598bd52 to 542f95e Compare September 27, 2026 12:57
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 542f95e to d107cf5 Compare September 27, 2026 13:45
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from e3a6f09 to 3754d1e Compare September 27, 2026 17:18
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 3754d1e to 726af4f Compare September 27, 2026 21:51
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch from 726af4f to 473c841 Compare September 27, 2026 22:02
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from 48a7983 to bd7075d Compare September 28, 2026 13:28
@lhotari
lhotari force-pushed the lh-dispatcher-lockfree-flow-master branch 2 times, most recently from 941d5a3 to 9ae5a9b Compare September 28, 2026 15:07
Base automatically changed from lh-ml-mpsc-add-handoff to master September 28, 2026 23:03
…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)
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.

1 participant