[improve][test] Add a catch-up scenario: consumers joining a live stream later - #26826
Merged
Merged
Conversation
…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)
This was referenced Oct 4, 2026
Open
dao-jun
approved these changes
Oct 4, 2026
11 tasks
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
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:
The time a consumer takes to catch up, and its read rate meanwhile, are what an operator sees. No scenario measured them.
Modifications
applications.joinSecondssetting 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.--run-id.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.joinSecondsmust have a value for each application, or none, and each value must be at least 0. The warmup and the latest join must fit withintimeoutSeconds.caughtUpLatencyMillismust be at least 1.iot-telemetry-catch-up.yamlextends the high-rate scenario: 500 gateways publish 30,000 msg/s to one topic for 180 s, and 5 applications consume it.AGENTS.md.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):
iot-application-1iot-application-2iot-application-3iot-application-4numIOThreads8 ([improve][broker] Default numIOThreads to half the available processors, at least 8 #26825), they took 71–76 s and 81–84 s.NullPointerExceptioninCommands.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
Each application's receive rate
Verifying this change
This change added tests and can be verified as follows:
MeasurementStartMarkerTestalso 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.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
This change was prepared with the assistance of Claude Code (claude-opus-5-5); I have reviewed and verified it.