Skip to content

[improve][ml] Batch managed-ledger adds across the thread boundary to the ledger executor - #26717

Merged
lhotari merged 14 commits into
masterfrom
lh-ml-mpsc-add-handoff
Sep 28, 2026
Merged

lhotari merged 14 commits into
masterfrom
lh-ml-mpsc-add-handoff

Conversation

@lhotari

@lhotari lhotari commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

Motivation

ManagedLedgerImpl.asyncAddEntry submits one executor task per entry to the managed ledger's ordered executor. With hundreds of producer connections on one topic, the connection threads contend on the executor's queue lock (BookKeeper's GrowableBatchedArrayBlockingQueue.offer → ReentrantLock.lock) once per message. The ledger thread also runs one task per message, and other work on that thread, such as waiting-cursor wake-ups for tailing dispatchers, queues behind those tasks.

Off-CPU profiling of 500 producers and a 20-member Key_Shared subscription shows this queue lock as the largest blocked time with an application frame in the broker after the dispatcher Flow contention, which #26716 removes.

Modifications

Batch the adds across the thread boundary between the publishing threads and the ledger executor, through a new package-private BatchingExecutorWrapper in the managed ledger module:

  • asyncAddEntry hands the add to the wrapper, which appends it to a JCTools MpscUnboundedArrayQueue (already used by RangeCacheRemovalQueue). The thread that finds no batch task scheduled submits one to the ledger executor, and that task runs the adds queued by then. Concurrent publishers contend on the executor's queue once per batch instead of once per add, and other executor work no longer waits behind one task per published message.
  • A batch stops after managedLedgerAddEntryHandoverMaxBatchItems adds (default 1024), or once the entries it has run add up to managedLedgerAddEntryHandoverMaxBatchBytesSize bytes (default 5 MB), and schedules the next batch for the rest, so add completions and other executor tasks keep running under load. The byte limit keeps a managed ledger with large entries from holding the executor thread, which is shared with other managed ledgers, for as long as a full batch of large entries would. A batch always runs at least one add, even one larger than the byte limit. The wrapper is generic: a task implementing BatchingExecutorWrapper.WeightedRunnable weighs getWeight(), and the managed ledger's add task weighs the readable bytes of its entry. That task replaces the lambda that captured the same fields, so it allocates nothing more. A failing add is logged and does not stop the batch.
  • The thread that sets the batch flag is the only consumer of the queue until it clears it, including while its batch runs, and it clears the flag before checking the queue again, so no add is left behind. If the executor rejects a batch task, that thread takes the queued adds: the caller's own add does not run and the caller fails with the executor's exception, as before this change, and every other queued add is failed through its callback with its buffer released. asyncAddEntry also releases the buffer it retained for an add that the executor rejects.
  • The queue is created by the first add, with a compare-and-set on a volatile field, so managed ledgers that are never written to do not allocate it; it is never replaced, so racing first adds all use the same queue. It uses 512-entry chunks and links another chunk only when a batch backs up beyond that: about 2.7 KB per managed ledger that has been written to.
  • Adds are still created and processed on the ledger executor, so the ledger's threading and each thread's add order are unchanged.
  • New dynamic broker settings managedLedgerAddEntryHandoverMaxBatchItems (default 1024) and managedLedgerAddEntryHandoverMaxBatchBytesSize (default 5242880), passed to the managed ledger as ManagedLedgerConfig.addEntryHandoverMaxBatchItems and addEntryHandoverMaxBatchBytesSize. A managed ledger captures the values when it opens, so the add path reads no shared configuration; an update applies to managed ledgers opened after it. A max batch items of 0 or 1 disables batching, and each add is then submitted to the executor as a task of its own, as before this change. A max batch bytes size of 0 removes the byte limit. Negative values are rejected.

Trade-offs of the handover batch limits

A batch runs to completion on the ledger's executor thread before any other task on that thread. Larger limits reduce scheduling overhead and contention between publishing threads at high publish rates, but keep the executor thread occupied for longer per batch, which can delay add completions, reads and cursor notifications for the managed ledgers that share the thread. Smaller limits favor that latency over add throughput. The item limit bounds the per-add overhead of a batch, and the byte limit bounds the work of a batch of large entries, so that on a heavily loaded broker or bookies a managed ledger with very large entries does not starve managed ledgers with small entries that share its executor thread. Both are tunable so that this balance can be adapted to the workload.

Measurements

The IoT telemetry max-rate (iot-telemetry-max-rate.yaml, 500 producers and a 20-member Key_Shared subscription, formerly iot-key-shared-500x20.yaml) and high-rate (iot-telemetry-high-rate.yaml) performance scenarios of tests/performance, run with the performance launcher on one host (Intel i9-9980HK, 8 cores, 16 hardware threads) running the broker, 3 bookies and the clients, on JDK 25.0.4.1 with ZGC. The host was configured with the performance testing environment setup of tests/performance/environment, whose TuneD profile disables turbo, so that the CPU ran at a fixed 2.4 GHz in every run. Baseline is the parent of this change, master at f925a25884e6 (#26715), and the candidate is the head of this change, 87c32c1213c9, with the default limits of 1024 adds and 5 MB per handover batch. Both were built from clean checkouts into separate images. The builds were interleaved run by run: three unprofiled max-rate runs (baseline and candidate, then candidate and baseline, then baseline and candidate), two unprofiled high-rate runs (baseline first, then candidate first) and one profiled max-rate run with --extends configs/profile-broker of each build. thermald was stopped, and the launcher let the CPU package cool down to 55 °C before each run and again before each measurement (-Pperformance.cooldownTemperature=55), so every run started at 48–53 °C. None of the runs throttled thermally. The metrics stack of tests/performance scraped the broker, the bookies and ZooKeeper every 5 s during each run, with the broker's stats periods set to the same 5 s.

Setting IoT telemetry max rate IoT telemetry high rate
Cluster 1 broker, 3 bookies 1 broker, 3 bookies
Ledger replication E=1, W=1, A=1 broker defaults
Producers 500 gateways × 1 topic 500 gateways × 1 topic
Applications 1 × 20 consumers, Key_Shared 5 × 10 consumers, Key_Shared
Messages 4,000,000 measured, 1,000,000 warmup 5,000,000 measured, 1,000,000 warmup
Payload 128 bytes, batching off 64 bytes, batching on
Rate limit none 30,000 msg/s

Every run delivered all messages to every application with no duplicates, no ordering violations and no invalid messages.

IoT telemetry max rate, unprofiled (runs in the order they ran):

Baseline (#26715) This change Change
Producer throughput (msg/s) 91,127 · 90,509 · 88,483 (mean 90,040) 99,923 · 100,805 · 99,012 (mean 99,913) +11.0 %
Publish latency p50 (ms) 1,070.1 · 1,078.3 · 1,094.7 993.3 · 987.6 · 1,004.5 −7.9 %
Publish latency p99 (ms) 1,442.8 · 1,421.3 · 1,520.6 1,152.0 · 1,091.6 · 1,123.3 −23.2 %
End-to-end latency p50 (ms) 1,352.7 · 1,382.4 · 1,379.3 1,000.4 · 1,003.0 · 1,019.4 −26.5 %
End-to-end latency p99 (ms) 1,909.8 · 1,956.9 · 2,105.3 1,161.2 · 1,104.9 · 1,140.7 −43.0 %
Sampled maximum backlog 61,892 · 59,365 · 57,182 36,333 · 52,369 · 47,169 −23.9 %

IoT telemetry high rate, unprofiled (end-to-end latency and backlog are the ranges over the 5 applications, and their change is that of the mean over the applications):

Baseline (#26715) This change Change
Producer throughput (msg/s) 30,000 · 30,000 (mean 30,000) 30,000 · 30,000 (mean 30,000) +0.0 %
Publish latency p50 (ms) 2.1 · 2.2 2.1 · 2.1 −2.3 %
End-to-end latency p50 (ms) 8 · 8 8 · 8 +0.0 %
End-to-end latency p99 (ms) 23 · 19–20 20–21 · 24–25 +7.1 %
Sampled maximum backlog per application 2,625–2,902 · 2,212–2,466 3,530–5,242 · 2,380–2,727 +43.1 %

Broker, profiled max-rate run (JFR recording, flame graphs and off-CPU profile of the measurement period, 4,000,000 messages on each side):

Baseline (#26715) This change Change
Producer throughput (msg/s) 88,192 95,001 +7.7 %
Blocked in GrowableBatchedArrayBlockingQueue.offer (off-CPU, observed) 1.835 s 0.002 s
Blocked time with an application frame (off-CPU, observed) 5.96 s 3.19 s
Ledger-thread CPU samples per million messages 1,037 958 −7.6 %
Broker CPU samples per million messages 6,774 6,237 −7.9 %

In the baseline, the connection threads blocking on the ledger executor's queue lock from ServerCnx.handleSend are 30.8 % of the broker's blocked time with an application frame, the largest single wait, next to the dispatcher monitor waits in internalConsumerFlow (56 % over its two rows) that #26716 removes; with this change that wait is gone.

The managed-ledger thread (BookKeeperClientWorker-OrderedExecutor-12-0) stays the serial stage, so throughput follows its cost per message, which drops with fewer executor tasks. The max-rate throughput improves by 11.0 %, with every run of this change faster than every baseline run. End-to-end latency now follows the publish latency closely (the max-rate end-to-end p50 is on average 13 ms above the publish p50, against 290 ms in the baseline), most likely because dispatch work on the ledger executor no longer queues behind one task per published message.

The high-rate scenario publishes at a fixed 30,000 messages per second, which this host sustains on both sides, so it shows the latency at a rate below saturation rather than the maximum throughput. Both sides deliver every message at that rate with a publish p50 of 2.1–2.2 ms and an end-to-end p50 of 8 ms. The end-to-end p99 of 19–25 ms and the backlogs of at most 5,300 messages vary between the runs of each side by more than the difference between the sides: the first run of this change had the highest backlogs, with a p99 within the baseline's range, and its second run the highest p99, with backlogs within the baseline's range.

In an earlier session, an eagerly created queue with 256-entry chunks and the final lazily created queue with 512-entry chunks were compared in 3 interleaved unprofiled runs each (means 110.7k and 108.7k msg/s, within the run-to-run spread).

Charts of the median unprofiled max-rate run of each side (baseline 90,509 msg/s, this change 99,923 msg/s), as the run report renders them:

Baseline (#26715) This change
Latency by percentile of the baseline Latency by percentile with this change
Throughput over time of the baseline Throughput over time with this change
Backlog over time of the baseline Backlog over time with this change

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • BatchingExecutorWrapperTest (unit tests with a hand-driven delegate):
    • adds queued before a batch runs share one batch;
    • a batch stops at the item limit, and at the weight limit, and the rest runs in the next batch;
    • a task heavier than the weight limit runs in a batch of its own, and unweighted tasks are limited only by the item limit;
    • a task queued while a batch runs is not left behind, and a failing task does not stop the batch;
    • the queue is created by the first task;
    • a rejected caller fails and its task does not run later; tasks queued meanwhile by other threads, and those left for a rejected follow-up batch, go to the rejected task handler;
    • invalid limits are rejected.
  • BatchingExecutorWrapperTest stress tests: 8 threads submitting 10,000 tasks each keep their per-thread order (with item limits 2 and 1024, and with a weight limit); against a delegate that always rejects, every task is rejected and none runs; against a delegate that rejects every third batch, every task runs or is rejected exactly once.
  • ManagedLedgerTest:
    • testAddEntryHandoverLimitsAreCapturedWhenOpened, testAddEntryHandoverWithoutByteLimit and testAddEntryHandoverBatchingDisabledUsesTheExecutor: a ledger passes its limits to the wrapper, a byte limit of 0 removes the weight limit, max batch items of 0 or 1 hands adds straight to the executor, and setConfig does not change the limits of an open ledger.
    • testAddEntryWithAddEntryHandoverBatchingDisabled and testAsyncAddEntriesWithSmallAddEntryHandoverMaxBatchBytesSize: adds complete in order with batching disabled, and with a byte limit that gives every add a batch of its own.
    • testRejectedAddEntryReleasesTheRetainedBuffer: with batching enabled and disabled, an add rejected by a shut down executor fails its caller and releases the buffer it retained.
    • testAddEntryHandoverMaxBatchLimitsRejectNegativeValues: defaults and validation of the configuration.
  • BrokerServiceTest.testManagedLedgerAddEntryHandoverMaxBatchItemsConfiguration, testManagedLedgerAddEntryHandoverMaxBatchBytesSizeConfiguration and testManagedLedgerAddEntryHandoverMaxBatchItemsDynamicUpdate: the broker settings reach the managed ledger configuration, and dynamic updates are validated and applied.
  • ManagedLedgerErrorsTest (22 cases) passes locally; it exercises add failures and recovery through the changed hand-off.

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: adds reach the managed ledger's executor in batches bounded by a number of adds and a total entry size, through an MPSC queue run by one batch task at a time, instead of one executor task per add. The adds still run on the same executor thread, in the same per-thread order.

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

@lhotari
lhotari added this pull request to stack #26718 September 25, 2026 19:28
@lhotari
lhotari marked this pull request as draft September 25, 2026 19:29
@lhotari
lhotari marked this pull request as ready for review September 25, 2026 19:30
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from 09e083c to a7036c2 Compare September 25, 2026 21:37
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from a7036c2 to 1e80d28 Compare September 25, 2026 22:22
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from 1e80d28 to 9d0081d Compare September 26, 2026 11:33
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from 9d0081d to 77dcc9c Compare September 26, 2026 12:02
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from 77dcc9c to ec23e49 Compare September 26, 2026 12:32
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch 2 times, most recently from 3754c95 to cf3ff5a Compare September 26, 2026 18:18
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch 2 times, most recently from 0ed9907 to c369e43 Compare September 26, 2026 20:18
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from 81e9053 to 5edbb51 Compare September 27, 2026 22:02
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch 3 times, most recently from 8784fd5 to e319d7d Compare September 28, 2026 13:32
Base automatically changed from lh-key-shared-500x20-scenario to master September 28, 2026 15:07
… MPSC queue

Motivation

ManagedLedgerImpl.asyncAddEntry submitted one executor task per entry. With
hundreds of producer connections on one topic, the connection threads contend
on the ledger executor's queue lock (BookKeeper GrowableBatchedArrayBlockingQueue)
once per message: off-CPU profiling with jonoffcpu shows it as the largest
remaining blocked time with an application frame in a 500-producer Key_Shared
load, and 4.9 s of blocked time in the IoT high-rate scenario.

Modifications

- asyncAddEntry offers the add to a JCTools MPSC queue and schedules a drain
  task only when none is scheduled, so concurrent publishers contend on the
  executor's queue once per batch of adds.
- The drain task processes at most 1024 adds and reschedules itself for the
  rest, so add completions and other executor tasks keep running under load.
  A failing add is logged and does not stop the drain; a rejected drain task
  clears the scheduled flag and fails the caller as before.
- The add is still created and processed on the executor, so the ledger's
  threading and each thread's add order are unchanged.

Assisted-by: Claude Code (claude-opus-5-5)
…it does

Batch adds across the thread boundary under names that say so: the add batch
queue, scheduleAddBatch and runAddBatch replace the hand-off and drain names,
and the comments describe the batching.

Create the queue on the first add with a compare-and-set on a volatile field,
so managed ledgers that are never written to do not allocate it. The queue is
never replaced, so racing first adds all use the queue that won the
compare-and-set. With that, use 512-entry chunks instead of 256, so that a
batch backing up during a burst links new chunks less often: an idle ledger
costs nothing and a ledger that has been written to about 2.7 KB.

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

The number of adds that the managed ledger's executor thread processes in one batch was fixed at 1024. A larger
batch reduces scheduling overhead and contention between publishing threads, but occupies the executor thread for
longer, which can delay add completions, reads and cursor notifications for the ledgers that share the thread.
Make it tunable so that the trade-off can be adapted to the workload.

- New dynamic broker setting managedLedgerMaxAddBatchSize (default 1024), passed to ManagedLedgerConfig as
  maxAddBatchSize. A managed ledger captures the value when it opens, so the add path reads no shared state;
  updates apply to managed ledgers opened afterwards.
- 0 disables batching: each add is submitted to the executor as a task of its own, as before the batching.
- Negative values are rejected by the dynamic configuration validator and by ManagedLedgerConfig.

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

Name the setting after what it limits, the adds handed over to the managed ledger's executor thread in one batch,
and rename the related ManagedLedgerConfig property, fields, methods and tests to match: the add batch queue becomes
the add entry handover queue.

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

Lead with the add entry handover it configures, so that the setting reads as the maximum batch size of the handover,
and rename the ManagedLedgerConfig property, fields, methods and tests to match.

Assisted-by: Claude Code (claude-opus-5-5)
…efore running it

Take up to addEntryHandoverMaxBatchSize adds out of the handover queue with a single drain call and then run them,
instead of polling the queue once per add. The list is local to the batch run on the ledger's executor thread. An
add that is still being offered when the drain stops is picked up by the next batch, which the existing emptiness
check schedules.

Assisted-by: Claude Code (claude-opus-5-5)
@lhotari
lhotari force-pushed the lh-ml-mpsc-add-handoff branch from e319d7d to 74b3d1f Compare September 28, 2026 15:07
… tests

Move the add entry handover tests from ManagedLedgerTest to BatchingExecutorWrapperTest as unit tests, and
remove the ManagedLedgerTest ones that relied on accessors removed by the extraction.

Assisted-by: Claude Code (claude-opus-5-5)
…h size is greater than 1

BatchingExecutorWrapper now requires a max batch size greater than 1, since a batch of one gives nothing over
handing each add over on its own. ManagedLedgerImpl still wrapped the executor for a batch size of 1, so opening
a ledger with that valid configuration failed. Treat 0 and 1 alike as disabling batching, document it, and test
that ledgers accept adds with either value.

Assisted-by: Claude Code (claude-opus-5-5)
@lhotari
lhotari marked this pull request as draft September 28, 2026 17:55
…heir entries

A handover batch ran up to managedLedgerAddEntryHandoverMaxBatchSize adds regardless of their size, so on a
loaded broker a ledger with very large entries could hold the executor thread for a long time, delaying add
completions, reads and cursor notifications for the ledgers that share the thread.

BatchingExecutorWrapper now takes a maxWeight in addition to maxItems: a batch stops taking tasks once the
weights of those it has run add up to maxWeight, and always runs at least one task. Tasks implementing
BatchingExecutorWrapper.WeightedRunnable carry a weight; other tasks weigh 0. ManagedLedgerImpl hands adds over
as an AddEntryHandover that weighs the readable bytes of its entry, in place of a lambda capturing the same fields.

The new limit is configured with managedLedgerAddEntryHandoverMaxBatchBytesSize (default 5 MB, 0 for no limit),
and managedLedgerAddEntryHandoverMaxBatchSize is renamed to managedLedgerAddEntryHandoverMaxBatchItems. Both are
dynamic and apply to managed ledgers opened after a change.

Assisted-by: Claude Code (claude-opus-5-5)
…d by a rejected handover batch

BatchingExecutorWrapper allocated its queue when it was created, so every managed ledger paid for it even when
it was never written to. The queue is now created by the first submitted task.

When the delegate rejected a handover batch, the task of the rejected caller stayed queued and could still run
with a later batch, although its caller had been failed, and tasks queued meanwhile by other threads were left
in the queue with no batch to run them. The thread that sets the handover flag is now the only consumer of the
queue until it clears it, including while a batch runs, so a rejected thread can take the queued tasks: its own
task is dropped and the rejection is thrown as before, and every other task is passed to a rejected task
handler. ManagedLedgerImpl fails such adds through their callback and releases their buffer.

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

Assisted-by: Claude Code (claude-opus-5-5)
…er limits in tests

Address review findings:
- asyncAddEntry retained the entry buffer before handing the add over, and leaked that reference when the
  executor rejected the add and the exception was thrown to the caller. Release it before rethrowing.
- Test that a ledger passes its handover limits to BatchingExecutorWrapper, that 0 bytes removes the weight
  limit, that 0 or 1 items hands adds straight to the executor, and that setConfig does not change the limits
  of an open ledger.
- Stress BatchingExecutorWrapper with a delegate that rejects some handover batches, checking that every task
  runs or is rejected exactly once.

Assisted-by: Claude Code (claude-opus-5-5)
@lhotari
lhotari marked this pull request as ready for review September 28, 2026 19:00
@lhotari
lhotari merged commit 4580b0e into master Sep 28, 2026
45 of 47 checks passed
@lhotari
lhotari deleted the lh-ml-mpsc-add-handoff branch October 1, 2026 00:02
@lhotari lhotari added this to the 5.0.0 milestone Oct 1, 2026
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 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.

2 participants