Skip to content

[improve][broker] Default numIOThreads to half the available processors, at least 8 - #26825

Open
lhotari wants to merge 1 commit into
apache:masterfrom
lhotari:lh-improve-io-threads-default
Open

lhotari wants to merge 1 commit into
apache:masterfrom
lhotari:lh-improve-io-threads-default

Conversation

@lhotari

@lhotari lhotari commented Oct 4, 2026 •

Copy link
Copy Markdown
Member

Motivation

The broker's Netty I/O threads (numIOThreads, the pulsar-io event loops) default to 2 * Runtime.getRuntime().availableProcessors(): 16 on an 8-processor host and 32 on a 16-processor one. In the Pulsar Performance Testing Framework's IoT telemetry scenarios, that many event loops cost CPU and throughput:

  • Less batching: each loop gets less work per wakeup, so it sleeps in epoll_wait more often and batches less per read and write system call.
  • More wakeups: a topic's managed-ledger thread hands every publish receipt to the producer connection's event loop (ServerCnx.execute), and every bookie write to the BookKeeper client's channel. With many mostly idle loops, most of those handovers wake a sleeping loop with an eventfd write.
  • The limit at small messages: at 128 B and the maximum rate, the managed-ledger thread is the broker's limit, at 95 % of a core. With 8 I/O threads instead of 32, its work per message fell by about a third.

Too few I/O threads hurt too: they can't keep up with dispatching to the consumers. So the default needs a floor.

Modifications

  • New default: half the available processors, but no fewer than 8 unless that exceeds twice the processors: Math.min(2 * n, Math.max(8, n / 2)) for n available processors.

    Available processors 1 2 4 8 16 32 64 128
    Current default 2 4 8 16 32 64 128 256
    New default 2 4 8 8 8 16 32 64

    Up to 4 processors, the default doesn't change.

  • Code: ServiceConfiguration.defaultNumIOThreads(int) computes the default from one read of the available processors. A test covers the default for 1 to 128 processors.

  • Docs: the documentation in ServiceConfiguration, broker.conf, standalone.conf and the Terraform template states the formula, and names the symptom of too few I/O threads: the busiest pulsar-io thread stays busy and consumers' backlogs grow.

  • Buffer share: ServerCnx divides maxMessagePublishBufferSizeInMB between the I/O threads, so each thread's share grows with fewer threads; the total is unchanged.

  • Unchanged: the proxy's numIOThreads, the BookKeeper client's bookkeeperClientNumIoThreads and the HTTP server's threads. An explicitly configured numIOThreads still wins.

This changes a default for deployments with more than 4 available processors, so it needs a release note.

Measurements

The framework's IoT telemetry scenarios, in the test cluster's single broker:

  • iot-telemetry-max-rate: 500 gateways publish unbatched messages to one topic without a rate limit, with up to 100,000 in flight, and 20 consumers on one Key_Shared subscription receive them.
  • iot-telemetry-high-rate: 30,000 msg/s to 30 topics, with 5 subscriptions of 10 consumers each.

Setup:

Means of 2 (min–max) Current default This change Change
16 hardware threads: I/O threads 32 8
Max rate 128 B, throughput 105.5k msg/s (105.0k–106.1k) 119.9k msg/s (118.1k–121.7k) +13.6 %
Max rate 128 B, broker CPU per million messages 55.7 s 40.3 s −27.6 %
Max rate 128 B with the gateways connecting concurrently (#26827), throughput 91.7k msg/s (90.1k–93.2k) 101.9k msg/s (101.7k–102.2k) +11.2 %
Max rate 8 KB, throughput 53.6k msg/s (53.2k–54.0k) 54.8k msg/s (54.7k–54.9k) +2.3 %
High rate, broker CPU per million messages 154.3 s 127.5 s −17.4 %
High rate, end-to-end p99 15.0 ms 14.5 ms
Catch-up (#26826), gateways connecting one at a time, means of 3: the first of 3 applications joining together caught up after 92.7 s (89.9–94.9) 73.8 s (71.4–75.9) −20.4 %
Catch-up: the application joining alone at 60 s caught up after 107.6 s (101.0–111.2) 82.4 s (80.5–83.8) −23.4 %
Catch-up, broker CPU per million messages 155.4 s 132.2 s −14.9 %
Catch-up, gateways connecting concurrently (#26827), means of 3: the first of 3 applications joining together caught up after 145.8 s (142.6–148.2) 101.7 s (98.9–105.2) −30.2 %
Catch-up, gateways connecting concurrently: the application joining alone at 60 s caught up after not in any run 120.2 s (119.5–120.7)
Catch-up, gateways connecting concurrently: broker CPU per million messages 162.1 s 139.9 s −13.7 %
8 hardware threads (broker pinned): I/O threads 16 8
Max rate 128 B, delivered to the consumers 97.4k msg/s (93.8k–101.0k) 113.6k msg/s (113.5k–113.7k) +16.6 %
Max rate 128 B, broker CPU per million messages 48.0 s 40.4 s −15.8 %
High rate, end-to-end p99 18 ms 19 ms
High rate, broker CPU per million messages 127.9 s 118.9 s −7.0 %
4 hardware threads (broker pinned): I/O threads 8 8 unchanged
  • 16 hardware threads: 8 I/O threads publish 13.6 % more at 128 B, with no overlap between the runs, at 28 % less broker CPU per message. The managed-ledger thread went from 95 % busy to 67–77 % at the higher rate. 8 KB gains 2.3 %, where the bookie limits the rate. The high-rate latency is the same at 17 % less broker CPU per message.
  • Catch-up: in the catch-up scenario ([improve][test] Add a catch-up scenario: consumers joining a live stream later #26826), applications join a live stream later and read their backlog while the gateways keep publishing. With 8 I/O threads they caught up about 20 % sooner, without overlap between the runs, at 15 % less broker CPU per message.
  • With connections spread: the framework's gateways used to connect one at a time, which put every data connection on every other I/O thread. With them connecting concurrently ([improve][test] Create the IoT gateways' producers concurrently in a random order #26827), the load spreads over all the I/O threads, and 8 still beat 32 by 11.2 %.
  • 8 hardware threads (broker pinned to 4 cores): 8 I/O threads deliver 16.6 % more than 16 at 128 B, at 16 % less broker CPU per message. The high-rate latency, with its 5 subscriptions per topic, is within a millisecond.
  • Why the floor of 8:
    • On 8 hardware threads, 4 I/O threads publish faster but deliver only 66k–85k msg/s, so the consumers fall behind.
    • On 4 hardware threads, half the processors would be 2 I/O threads. They deliver only about 62k msg/s against 106k with 8, and the end-to-end p99 reaches 32–34 s.
    • The floor keeps 8 on those hosts.
  • 8 isn't the measured optimum on 16 hardware threads, where 6 did a little better in an earlier sweep. It keeps a margin above the point where dispatch can't keep up.
    • Catch-up with fewer threads: with the gateways connecting concurrently, 6 and 4 I/O threads caught up faster still, at a higher tail latency for the application reading at the tail. The table has means of 2 runs.

      I/O threads 8 6 4
      The first application caught up after 94.2 s 89.5 s 78.5 s
      Broker CPU per million messages 137.3 s 132.8 s 124.1 s
      The tailing application's end-to-end p99 43 ms 52 ms 60 ms
    • Why the floor stays at 8: on smaller brokers, 4 aren't enough at the maximum rate (above).

Max rate 128 B on a broker with 8 hardware threads

Published and delivered throughput with 16, 8 and 4 I/O threads on a broker with 8 hardware threads

Max rate 128 B on a broker with 4 hardware threads: why the default doesn't go below 8

Published and delivered throughput with 8, 4 and 2 I/O threads on a broker with 4 hardware threads

Throughput at max rate 128 B on 16 hardware threads, numIOThreads 32 (A) and 8 (B), with the same axes

Throughput over time with 32 I/O threads above and 8 below

Extrapolated: above 17 processors, the new default is half the processors, which nothing measured. Up to 17, the measurements support the floor of 8.

Shared with other uses: the broker's I/O event loop group also serves the broker's internal client for replication and lookups. Protocol handlers with a dedicated worker group get numIOThreads threads too.

Not measured: TLS (handshakes and encryption run on the event loops), 128 KB entries, thousands of connections or topics, a reconnect storm, geo-replication, and protocol handlers with dedicated worker groups sized from numIOThreads. The test host runs the clients and bookies beside the broker, so part of the gain on 16 hardware threads may come from less CPU contention, which a broker on a host of its own wouldn't have. The pinned runs model a broker on cores of its own.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • ServiceConfigurationDefaultsTest:
    • The default for 1, 2, 3, 4, 5, 7, 8, 16, 17, 18, 32 and 128 available processors.
    • That numIOThreads defaults to it.

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
    • numIOThreads: half the available processors, but no fewer than 8 unless that exceeds twice the processors, instead of twice the processors.
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

This change was prepared with the assistance of Claude Code (claude-opus-5-5); I have reviewed and verified it.

…rs, at least 8

Assisted-by: Claude Code (claude-opus-5-5)

@dao-jun dao-jun 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

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.

2 participants