Repository navigation
[improve][ml] Batch managed-ledger adds across the thread boundary to the ledger executor - #26717
Merged
Merged
Conversation
lhotari
added this pull request to stack #26718
September 25, 2026 19:28
lhotari
marked this pull request as draft
September 25, 2026 19:29
lhotari
marked this pull request as ready for review
September 25, 2026 19:30
1 of 11 tasks
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 25, 2026 21:37
09e083c to
a7036c2
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 25, 2026 22:22
a7036c2 to
1e80d28
Compare
1 of 11 tasks
Draft
1 of 11 tasks
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 26, 2026 11:33
1e80d28 to
9d0081d
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 26, 2026 12:02
9d0081d to
77dcc9c
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 26, 2026 12:32
77dcc9c to
ec23e49
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
2 times, most recently
from
September 26, 2026 18:18
3754c95 to
cf3ff5a
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
2 times, most recently
from
September 26, 2026 20:18
0ed9907 to
c369e43
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 27, 2026 22:02
81e9053 to
5edbb51
Compare
lhotari
force-pushed
the
lh-ml-mpsc-add-handoff
branch
3 times, most recently
from
September 28, 2026 13:32
8784fd5 to
e319d7d
Compare
… 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
force-pushed
the
lh-ml-mpsc-add-handoff
branch
from
September 28, 2026 15:07
e319d7d to
74b3d1f
Compare
merlimat
approved these changes
Sep 28, 2026
… 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
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
marked this pull request as ready for review
September 28, 2026 19:00
Radiancebobo
pushed a commit
to Radiancebobo/pulsar
that referenced
this pull request
Oct 8, 2026
… the ledger executor (apache#26717)
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
ManagedLedgerImpl.asyncAddEntrysubmits 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'sGrowableBatchedArrayBlockingQueue.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
BatchingExecutorWrapperin the managed ledger module:asyncAddEntryhands the add to the wrapper, which appends it to a JCToolsMpscUnboundedArrayQueue(already used byRangeCacheRemovalQueue). 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.managedLedgerAddEntryHandoverMaxBatchItemsadds (default 1024), or once the entries it has run add up tomanagedLedgerAddEntryHandoverMaxBatchBytesSizebytes (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 implementingBatchingExecutorWrapper.WeightedRunnableweighsgetWeight(), 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.asyncAddEntryalso releases the buffer it retained for an add that the executor rejects.managedLedgerAddEntryHandoverMaxBatchItems(default 1024) andmanagedLedgerAddEntryHandoverMaxBatchBytesSize(default 5242880), passed to the managed ledger asManagedLedgerConfig.addEntryHandoverMaxBatchItemsandaddEntryHandoverMaxBatchBytesSize. 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 of0or1disables 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 of0removes 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, formerlyiot-key-shared-500x20.yaml) and high-rate (iot-telemetry-high-rate.yaml) performance scenarios oftests/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 oftests/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,masteratf925a25884e6(#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-brokerof each build.thermaldwas 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 oftests/performancescraped the broker, the bookies and ZooKeeper every 5 s during each run, with the broker's stats periods set to the same 5 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):
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):
Broker, profiled max-rate run (JFR recording, flame graphs and off-CPU profile of the measurement period, 4,000,000 messages on each side):
GrowableBatchedArrayBlockingQueue.offer(off-CPU, observed)In the baseline, the connection threads blocking on the ledger executor's queue lock from
ServerCnx.handleSendare 30.8 % of the broker's blocked time with an application frame, the largest single wait, next to the dispatcher monitor waits ininternalConsumerFlow(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:
Verifying this change
This change added tests and can be verified as follows:
BatchingExecutorWrapperTest(unit tests with a hand-driven delegate):BatchingExecutorWrapperTeststress 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,testAddEntryHandoverWithoutByteLimitandtestAddEntryHandoverBatchingDisabledUsesTheExecutor: 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, andsetConfigdoes not change the limits of an open ledger.testAddEntryWithAddEntryHandoverBatchingDisabledandtestAsyncAddEntriesWithSmallAddEntryHandoverMaxBatchBytesSize: 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,testManagedLedgerAddEntryHandoverMaxBatchBytesSizeConfigurationandtestManagedLedgerAddEntryHandoverMaxBatchItemsDynamicUpdate: 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
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.