Skip to content

[improve][broker] Raise the default dispatcherMaxReadBatchSize from 100 to 500 - #26744

Merged
merlimat merged 1 commit into
apache:masterfrom
lhotari:lh-improve-read-batch-500
Sep 29, 2026
Merged

merlimat merged 1 commit into
apache:masterfrom
lhotari:lh-improve-read-batch-500

Conversation

@lhotari

@lhotari lhotari commented Sep 29, 2026

Copy link
Copy Markdown
Member

Motivation

In the IoT telemetry max-rate performance scenario (500 producers publishing 128-byte unbatched messages to one topic, one application consuming it with a 20-member Key_Shared subscription), the dispatcher runs behind the tail and every read returns the maximum of 100 entries, about 5 entries for each consumer per read. Larger reads raise the throughput and lower the latency, up to a plateau at 250–500 entries.

Modifications

  • ServiceConfiguration.dispatcherMaxReadBatchSize: default 500 instead of 100, and its documentation.
  • conf/broker.conf and deployment/terraform-ansible/templates/broker.conf: the same default.

The setting also bounds the replicator's reads (PersistentReplicator) and the batch of PersistentSubscription's backlog analysis; the read size in bytes stays limited by dispatcherMaxReadSizeBytes (5 MB).

Measurements

Master c5a001313e8e (this PR is the same change rebased onto a later master), performance launcher of tests/performance on one host (Intel i9-9980HK, 8 cores, 16 hardware threads, turbo off, fixed 2.4 GHz, thermald stopped, cool-down to 55 °C before each run), no profiling. The setting was varied with --set cluster.brokers.env.dispatcherMaxReadBatchSize=<n> on one image. 100 and 500 were interleaved in one series, 1000 ran right after it, and 250 and 750 were interleaved in a later series. Every run delivered every message with no duplicates, ordering violations or invalid messages, and none throttled.

iot-telemetry-max-rate (4,000,000 measured messages, 128 bytes, no rate limit), means of 3 runs:

dispatcherMaxReadBatchSize Producer throughput (msg/s) Range Publish p50 (ms) Publish p99 (ms) End-to-end p99 (ms)
100 (current default) 99,383 98,005–100,388 996.5 1,126.7 1,229.5
250 106,625 106,315–106,824 926.4 1,120.3 1,144.2
500 106,068 103,882–107,697 940.7 1,052.8 1,062.6
750 102,422 100,489–103,601 970.9 1,094.0 1,108.3
1000 98,542 97,891–99,718 1,005.9 1,190.6 1,201.8

Charts of the median unprofiled max-rate run of each side, as the run report renders them (master c5a001313e8e with the current default of 100: 99,756 msg/s; with 500: 106,624 msg/s):

Note

Each chart has its own scale, set by that run's values, so the baseline's and this change's charts of the same measure have different axes. Compare their values, not the heights of the lines.

Latency by percentile, baseline (master, 100)

Latency by percentile of the baseline

Latency by percentile, this change (500)

Latency by percentile with this change

Throughput over time, baseline (master, 100)

Throughput over time of the baseline

Throughput over time, this change (500)

Throughput over time with this change

Backlog over time, baseline (master, 100)

Backlog over time of the baseline

Backlog over time, this change (500)

Backlog over time with this change

iot-telemetry-high-rate (30,000 msg/s, 5 applications): 100 → 500 lowers the publish p99 from 14.2 to 8.7 ms and the end-to-end p99 from 23.2 to 16.8 ms (means of 2 runs; one 100 run is an outlier at 27–28 ms, the other is at 19 ms), with the same end-to-end p50 of 8 ms.

On #26716 (Flow permits without the dispatcher monitor), 500 raised the max-rate throughput by 9.7 % over 100.

Why 750 and 1000 are slower: profiling shows that the managed-ledger thread is the serial stage (91 % busy at 100, 97 % at 500 and 1000), and that a read asking for more entries than are confirmed is extended with further reads on that thread, which then also runs the Key_Shared dispatch. That work grows with the batch size (3, 31 and 83 CPU samples per million messages at 100, 500 and 1000). #26743 completes such reads at the last confirmed entry instead; with it, 500 and 1000 measure the same (104.8k and 104.9k msg/s, +4.2 % over master at 500 in that series).

Larger payloads (iot-telemetry-max-rate with the payload size, measured messages scaled to 4–5 GB per run, the gateways' in-flight messages to about 100 MB, and more direct memory for the workload containers), 100 against 500, 3 interleaved runs each:

Payload 100: msg/s mean (range) 500: msg/s mean (range) Change
4 KB 60,454 (58,112–61,891) 59,765 (57,835–62,364) −1.1 %
8 KB 36,591 (34,553–37,622) 37,041 (34,693–38,227) +1.2 %
64 KB 5,545 (5,191–5,764) 5,501 (5,228–5,797) −0.8 %
128 KB 2,437 (2,079–2,622) 2,814 (2,650–2,897) +15.5 %
128 KB, 120,000 messages (47 s) 2,545 (2,494–2,603) 2,464 (2,372–2,578) −3.2 %

From 4 KB up the workload is bound by bytes (about 240–380 MB/s), and the setting makes no difference. At 64 and 128 KB, dispatcherMaxReadSizeBytes (5 MB) limits a read to about 80 and 40 entries, fewer than either setting. The +15.5 % of the first 128 KB series, whose measurements lasted only 14–17 s, did not repeat in 3 runs each of 120,000 messages (−3.2 %, within the spread). The p99 stalls of about 400–850 ms seen at every payload size, with either setting, come from the storage: with one topic and E=W=A=1, the topic's only ledger lives on one bookie, which flushes its ledger storage almost continuously and throttles adds (bookie_throttled_write); its throttled-write time peaks during the stalls. Striping the ledger over the three bookies (E=3, W=1, A=1) did not raise the throughput (all bookies share the host's disk).

Verifying this change

  • Make sure that the change passes the CI checks.

This change is already covered by existing tests: the tests that depend on the read batch size set it explicitly.

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

Default values of configurations: dispatcherMaxReadBatchSize changes from 100 to 500. It is a dynamic setting, so a broker can be set back to 100 without a restart.

This PR was prepared with AI assistance (Claude Code) and reviewed by a human contributor.

…00 to 500

With 500 producers and a 20-member Key_Shared subscription on one topic (iot-telemetry-max-rate performance
scenario), every dispatcher read returns the maximum of 100 entries, since the dispatcher runs behind the tail.
A sweep on master c5a0013 measured the producer throughput at 99.4k msg/s with 100, 106.6k with 250,
106.1k with 500, 102.4k with 750 and 98.5k with 1000 (means of 3 runs each), with lower publish and end-to-end
latency at 250 and 500.

Change the default in ServiceConfiguration, conf/broker.conf and the terraform-ansible template.

Assisted-by: Claude Code (claude-opus-5-5)
@merlimat
merlimat merged commit a56538b into apache:master Sep 29, 2026
44 checks passed
@lhotari lhotari added this to the 5.0.0 milestone Oct 1, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants