Skip to content

[improve][ml] Complete a read at the last confirmed entry instead of extending it on the ledger executor - #26743

Merged
merlimat merged 2 commits into
apache:masterfrom
lhotari:lh-improve-ml-read-at-lac
Sep 29, 2026
Merged

merlimat merged 2 commits into
apache:masterfrom
lhotari:lh-improve-ml-read-at-lac

Conversation

@lhotari

@lhotari lhotari commented Sep 29, 2026

Copy link
Copy Markdown
Member

Motivation

When a cursor read asks for more entries than the managed ledger has confirmed, ManagedLedgerImpl.internalReadFromLedger reads the current ledger only up to the last confirmed entry. OpReadEntry.checkReadCompletion then sees that more entries were confirmed while the read ran (cursor.hasMoreEntries()), and schedules another read on the managed ledger's executor to fill the requested count, repeating until the count is reached. Its comment says it is meant for continuing into the next ledger. These reads are served from the entry cache, so the whole read completion, including the dispatcher's work (Key_Shared's filterAndGroupEntriesForDispatching, sending to the consumers), then runs on the managed ledger's executor thread, which is also the thread that adds entries.

Profiling of the IoT telemetry max-rate performance scenario on master c5a001313e8e (500 producers, 128-byte unbatched messages, one topic, 20-member Key_Shared subscription) shows that thread, BookKeeperClientWorker-OrderedExecutor-12-0, as the serial stage: 91 % busy with dispatcherMaxReadBatchSize 100 and 97 % with 500 or 1000. Its CPU samples per million messages in these extended reads (under OpReadEntry.checkReadCompletion → ManagedLedgerImpl.asyncReadEntries) grow with the batch size, and account for most of the growth of the thread's cost per message:

dispatcherMaxReadBatchSize 100 500 1000
Managed-ledger thread, CPU samples per million messages 966 1,004 1,061
of which extended reads under checkReadCompletion 3 31 83
of which Key_Shared dispatch (trySendMessagesToConsumers) 2 22 57
pulsar-io threads, CPU samples per million messages 4,553 4,296 4,435
Producer throughput (profiled runs, mean of 2) 94,198 96,394 91,969

This is why the throughput falls when dispatcherMaxReadBatchSize is raised beyond 500 (unprofiled means on master: 99.4k msg/s with 100, 106.6k with 250, 106.1k with 500, 102.4k with 750, 98.5k with 1000).

Modifications

  • OpReadEntry.readUpToLastConfirmedEntry: whether the last ledger read of the operation ended at the managed ledger's last confirmed entry; ManagedLedgerImpl.internalReadFromLedger sets it for each ledger read, and when the read position is already past the last confirmed entry of the current ledger.
  • OpReadEntry.checkReadCompletion completes such a read with the entries it already has. Reads that cross into the next ledger, and reads that have returned no entries yet (for example because all of them were filtered as individually deleted), continue as before.

The entries confirmed after the read are read by the dispatcher's next read, from the thread that issues it.

Measurements

tests/performance launcher 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). Baseline master c5a001313e8e, candidate this change on top of it (9f888c13858e; this PR is the same change rebased onto a later master), each built into its own image; dispatcherMaxReadBatchSize set with --set cluster.brokers.env.dispatcherMaxReadBatchSize=<n>. The runs were interleaved in one series (master@500, change@500, change@1000, change@100, three rounds in rotating order), then the high-rate runs and two profiled runs. Every run delivered every message with no duplicates, ordering violations or invalid messages, and none throttled.

iot-telemetry-max-rate (4,000,000 measured 128-byte messages, 500 producers, 1 × 20 Key_Shared consumers, no rate limit), 3 unprofiled runs each:

Producer throughput (msg/s) Change Publish p50 (ms) Publish p99 (ms) End-to-end p99 (ms) Sampled max backlog
master, 500 99,583 · 100,109 · 102,071 (mean 100,588) 992.6 1,104.9 1,117.9 20,111
this change, 500 106,181 · 103,347 · 104,952 (mean 104,827) +4.2 % 948.4 1,051.0 1,057.3 16,477
this change, 1000 105,744 · 104,999 · 103,928 (mean 104,890) +4.3 % 949.2 1,050.1 1,051.6 12,597
this change, 100 97,299 · 97,579 · 96,067 (mean 96,982) 1,024.2 1,149.3 1,157.1 55,850

Every run of this change at 500 or 1000 is faster than every master run at 500, and 1000 no longer falls behind 500 (on master in an earlier series: 106.1k msg/s at 500, 98.5k at 1000). With the default of 100 a read rarely asks for more entries than are confirmed, so the change has little to act on; there was no master@100 run in this series (99.4k msg/s in an earlier series, when master@500 measured 106.1k, so the host ran about 5 % slower in this series).

Charts of the median unprofiled max-rate run of each side at dispatcherMaxReadBatchSize 500, as the run report renders them (master c5a001313e8e: 100,109 msg/s; this change: 104,952 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, 500)

Latency by percentile of the baseline

Latency by percentile, this change (500)

Latency by percentile with this change

Throughput over time, baseline (master, 500)

Throughput over time of the baseline

Throughput over time, this change (500)

Throughput over time with this change

Backlog over time, baseline (master, 500)

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 × 10 Key_Shared consumers), 2 runs each at 500: publish p50 1.9 ms (master) and 1.8 ms (this change), end-to-end p50 8 ms on both, end-to-end p99 15–18 ms (master) and 17–19 ms (this change), within the run-to-run spread.

Profiled broker (max-rate, --extends configs/profile-broker), CPU samples per million messages (10 ms interval); master from 2 runs per size of an earlier series, this change from 1 run per size:

master 500 master 1000 this change 500 this change 1000
Managed-ledger thread BookKeeperClientWorker-OrderedExecutor-12-0 1,004 1,061 975 920
of which extended reads under OpReadEntry.checkReadCompletion 31 83 0 0
of which Key_Shared dispatch (trySendMessagesToConsumers) 22 57 0 1
pulsar-io threads 4,296 4,435 4,148 3,915

The extended reads and the dispatch they carried are gone from the managed-ledger thread. That thread stays the serial stage (about 95 % busy in the profiled runs of this change), now with the cost of the adds alone.

Verifying this change

  • Make sure that the change passes the CI checks.

This change is already covered by existing tests, such as ManagedCursorTest, NonDurableCursorTest, ManagedLedgerTest, ManagedLedgerErrorsTest and OpReadEntryConfigTest, which pass locally (436 tests). There is no new unit test for the tail-read case yet: it needs entries to be confirmed between the start and the completion of a cache read.

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: a read that ends at the last confirmed entry completes on the thread that performed it, instead of being extended on the managed ledger's executor, so the dispatcher's read completion no longer runs on the thread that adds entries in that case.

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

…extending it on the ledger executor

When a cursor read asks for more entries than the managed ledger has confirmed, the read of the current ledger
stops at the last confirmed entry. OpReadEntry.checkReadCompletion then sees that more entries were confirmed while
the read ran, and schedules another read on the managed ledger's executor to fill the requested count, repeating
until the count is reached. Those reads are served from the entry cache, so the whole read completion, including the
dispatcher's work such as Key_Shared's filterAndGroupEntriesForDispatching and sending to consumers, then runs on the
managed ledger's executor thread, which is also the thread that adds entries.

In the iot-telemetry-max-rate performance scenario (500 producers, 20-member Key_Shared subscription, one topic),
that thread is the serial stage: it is 91 % busy with dispatcherMaxReadBatchSize 100 and 97 % with 500 or 1000, and
its CPU samples per million messages in these extended reads grow from 3 (100) to 31 (500) and 83 (1000).

Record whether the last ledger read of an OpReadEntry ended at the last confirmed entry, and complete such a read
with the entries it already has. Reads that cross into the next ledger, and reads that have returned no entries yet
(for example because all of them were filtered as individually deleted), continue as before.

Assisted-by: Claude Code (claude-opus-5-5)
Assisted-by: Claude Code (claude-opus-5-5)
@merlimat
merlimat merged commit 372a92a into apache:master Sep 29, 2026
43 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.

4 participants