Skip to content

[improve][test] Add a catch-up scenario: consumers joining a live stream later - #26826

Merged
lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-improve-perf-catch-up
Oct 4, 2026
Merged

lhotari merged 5 commits into
apache:masterfrom
lhotari:lh-improve-perf-catch-up

Conversation

@lhotari

@lhotari lhotari commented Oct 4, 2026

Copy link
Copy Markdown
Member

Motivation

The Pulsar Performance Testing Framework's IoT telemetry scenarios measure consumers that read at the tail of the stream from the start. A common case they don't cover is a consumer that joins a live stream later, such as after a deployment, a scale-up or an outage. Such a consumer has to catch up on a backlog while the producers keep publishing. Catch-up reads exercise other broker paths than tailing reads:

  • Reads from storage: the backlog is read from the bookies instead of the broker's entry cache.
  • Cache inserts: the broker inserts the entries that it read into the cache, for other cursors that read them.
  • Dispatch: the dispatcher reads in batches as fast as the consumers take them.

The time a consumer takes to catch up, and its read rate meanwhile, are what an operator sees. No scenario measured them.

Modifications

  • Joining later: a new applications.joinSeconds setting gives each application a join time, in seconds after the measurement starts. It has a value for each application, and 0 joins at the start, as every application does without the setting.
    • An application that joins later creates its Key_Shared subscription at the earliest position on every topic before the gateways publish, without keeping a consumer. Its backlog then builds up until it joins.
    • It marks the warmup rounds received, so that the gateways don't wait for it.
    • A thread opens its pods at its join time, counted from the gateways' measurement start. The gateways write the measurement start to the coordination directory as a marker file, keyed by the run ID, so a run with late applications requires --run-id.
  • Caught up: an application has caught up when each of its topics has delivered a measured message within applications.caughtUpLatencyMillis (1,000 ms by default) of its publishing. The application summary gets the join time, the catch-up time and the messages received by then.
  • Validation: joinSeconds must have a value for each application, or none, and each value must be at least 0. The warmup and the latest join must fit within timeoutSeconds. caughtUpLatencyMillis must be at least 1.
  • Run report: a Catch-up section lists each late application's join time, its subscription's backlog then (from the sampled topic stats), how long it took to catch up, its catch-up rate and when it received its last message. An application that doesn't catch up shows its rate until its last message instead. The Throughput section notes that late applications dominate its delivered throughput.
  • New scenario: iot-telemetry-catch-up.yaml extends the high-rate scenario: 500 gateways publish 30,000 msg/s to one topic for 180 s, and 5 applications consume it.
    • Joins: applications 1, 2 and 3 join 2 s apart, at 20, 22 and 24 s, as consumers that scale up or restart together do. Application 4 joins alone at 60 s.
    • Entry cache: when the first of the three reads entries from storage, the others are expected to read them within the cache's time to live: 1 s, extended up to 5 times while reads are expected. The 2 s spacing is within that. The last application is too far behind and reads its backlog from storage.
  • Docs: the scenario's documentation, the scenarios README, the run reports documentation and the framework's AGENTS.md.
  • Tests: the validation, the measurement start marker, the catch-up tracking (including listeners that record out of order), the report's Catch-up section, and the scenario file.

Example run

The scenario on master, with the bookies' journal on a tmpfs, on an Intel i9-9980HK (8 cores / 16 threads at a fixed 2.4 GHz):

Application Joined Backlog when it joined Caught up after Catch-up rate Received its last message after
iot-application-1 20.0 s 1,598,275 70.0 s 52,392 msg/s 160.0 s
iot-application-2 22.0 s 1,658,791 67.8 s 54,012 msg/s 158.0 s
iot-application-3 24.0 s 1,719,721 65.8 s 55,648 msg/s 156.0 s
iot-application-4 60.0 s 2,784,161 90.3 s 60,666 msg/s 120.0 s
  • Shared capacity: the 3 applications that join together converge within seconds and then read in step, at about 60k msg/s each. When the fourth joins, the three slow to about 45k msg/s, so catch-up reads share the broker's capacity.
  • Run-to-run variation: this run was among the fastest. In 3 later runs on master with nothing else on the host, the first application caught up after 85–92 s and the one joining at 60 s after 99–107 s. With numIOThreads 8 ([improve][broker] Default numIOThreads to half the available processors, at least 8 #26825), they took 71–76 s and 81–84 s.
  • A broker bug: the scenario's concurrent catch-up reads of shared cached entries found a race in the broker's shared message metadata. In 3 of 14 runs on master, a Key_Shared subscription stopped dispatching for good after a NullPointerException in Commands.resolveStickyKey. [fix][ml] Decode the shared cached message metadata before publishing it to other threads #26824 fixes it.

Each subscription's backlog: building up until its application joins, and falling as it catches up

The backlog of the 4 late applications' subscriptions rising from the warmup, the 3 that join at 20 to 24 s falling together to zero at about 90 s, and the one that joins at 60 s falling to zero at about 150 s

Each application's receive rate

Each application's receive rate: the on-time application at the gateways' 30,000 msg/s, the 3 that join at 20 to 24 s at about 60,000 msg/s each until the fourth joins at 60 s, then about 45,000 msg/s, and the fourth at about 65,000 msg/s once the others caught up

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • MeasurementStartMarkerTest also checks that a measurement that doesn't start times out; the application then doesn't join, and the run's timeout reports it.
  • IotScenarioTest: the join settings, their defaults and their validation, including with a warmup, and the run ID requirement.
  • MeasurementStartMarkerTest: the marker is written atomically and read for its run ID only.
  • CatchUpTrackerTest: catching up when every topic delivered within the threshold after the join, at the latest topic's receipt also when listeners record out of order.
  • RunReportTest.reportsTheCatchUpOfTheApplicationsThatJoinedLater: the Catch-up section, with and without topic stats, an application that didn't catch up, one that didn't receive a measured message, and one that caught up in the millisecond that it joined.
  • ScenarioFilesTest.theCatchUpScenarioJoinsApplicationsAfterTheMeasurementStarts.
  • The scenario ran on master; the table and the charts are from that run.

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

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

…stream after the measurement starts

applications.joinSeconds has applications join after the measurement starts. A late application creates its
subscription at the earliest position before the gateways publish, so that its backlog builds up, skips the warmup,
and opens its pods when the gateways' measurement start, which they write to the coordination directory, plus its
join time has passed. It has caught up when it first receives a measured message within
applications.caughtUpLatencyMillis of its publishing. The run report's Catch-up section shows each late
application's join time, its subscription's backlog then, how long it took to catch up and its catch-up rate.

iot-telemetry-catch-up.yaml publishes the high-rate scenario's 30,000 msg/s for 120 s, with 4 of its 5 applications
joining at 20 s (2 of them), 40 s and 60 s.

Assisted-by: Claude Code (claude-opus-5-5)
… join 3 applications 2 s apart

The applications that join while the gateways publish at 30,000 msg/s may only catch up after the gateways finish,
so the Catch-up section also shows when each received its last message and its read rate since joining. The
catch-up scenario joins 3 applications 2 s apart, close enough for the entries that the first reads from storage to
stay in the broker's entry cache for the others, and the last one alone at 60 s.

Assisted-by: Claude Code (claude-opus-5-5)
… at measurement start, track catch-up

Assisted-by: Claude Code (claude-opus-5-5)
…eceipt, pod cleanup on shutdown, report edge cases

Assisted-by: Claude Code (claude-opus-5-5)
…, catch-up result published together, docs

Assisted-by: Claude Code (claude-opus-5-5)
@lhotari
lhotari merged commit b2f6041 into apache:master Oct 4, 2026
44 checks passed
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