Repository navigation
[improve][ml] Allow inline read completion and preserve replication order - #26620
Conversation
|
One possible interleaving in
This appears to leave room for submission order |
|
Another potential ordering case is an inline completion overtaking a queued completion:
This depends on whether issuing another read before the previous callback runs is supported. The expected callback ordering for that usage is not clear from the documentation, particularly whether a later inline completion may overtake an earlier queued one. |
|
A separate design consideration is avoiding recursive read callbacks in continuous readers such as Its current callback can call A per-replicator processing loop might be an alternative: reentrant or concurrent requests are recorded, and follow-up work is processed after the current callback returns. This could retain inline cached reads while avoiding a growing callback stack. For replication, that execution ownership would also need to cover the entire batch submission. Avoiding same-thread recursion alone would still leave the question of ACK-triggered reads overlapping submission on another thread. |
Complete successful cursor reads inline on any completing thread, retaining the existing FastThreadLocal nesting limit and queued stack unwind. Remove the callback opt-in API introduced by the parent change and add keyed multi-consumer profiling coverage. Assisted-by: Codex
b5ba2dd to
b8a3af6
Compare
|
@void-ptr974 Thank you. Good points. |
Serialize cursor invocation and complete batch submission with one processing owner while retaining asynchronous acknowledgments. Coordinate cancellation, rewind and retry backoff with that owner, and settle unsent entries in Geo and Shadow replication. Cap bounded batch-read estimates to the BookKeeper frame limit, clarify callback ownership, and correct raw-entry and asynchronous cache-eviction test fixtures. Add deterministic ordering, recovery, resource-lifetime and batch-range regression coverage. Assisted-by: Codex
|
@void-ptr974 The follow-up commit 93204782decf, on top of b8a3af6, addresses these concerns:
The scoped tests pass locally, including all five invocations of |
|
One clarification relevant to this discussion: BookKeeper does not guarantee FIFO completion of concurrent read operations, including reads on the same ledger. The entries within a returned range are ordered, but that does not imply that separate read futures complete in their invocation order. I checked this against the API and the implementation in BookKeeper 4.18.1, which this PR uses. The following links are pinned to that release's commit:
Consequently, a later read whose responses arrive first can complete before an earlier read still waiting for a response. This follows from the independent operation-completion paths above. Executing callbacks on a common ordered executor serializes the work enqueued there; it does not reconstruct the original read-invocation order when responses become ready in a different order. LAC controls which entries are confirmed and readable, not the completion order of concurrent reads. For this PR, the replicator therefore enforces the ordering it needs itself: one outstanding cursor read, with the next read admitted only after the current batch finishes submission. Cache hits can bypass BookKeeper entirely, so the same coordination also covers inline cached completions. |
Allow storage to recover after three alternating read errors so the test verifies recovery without depending on acknowledgements bypassing the retry backoff. Retain full delivery verification and assert the injected errors occurred. Assisted-by: Codex
Assert a synchronous producer completion cannot submit the next cached batch before the current batch finishes. This also covers same-thread nesting permitted before the generic read-completion handoff removal. Assisted-by: Codex
|
I traced the replication ordering protections before and after #24189. There is a pre-existing ordering hazard before this PR, but the history needs a distinction between removing a guard and exposing the resulting reentrancy. Before #24189, #24189 changed the pending-read condition to #26599 subsequently allowed read completion inline when already on the ledger thread. This permits the following sequence even before #26620:
This sequence follows from the production code; it does not depend on BookKeeper returning reads out of order. #26620 broadens the exposure by also allowing successful callbacks on other threads, including the ordinary concurrent-ACK case raised in this review. The fix in 9320478 makes batch-submission ownership explicit: the next read is admitted only after the current batch finishes submission. ACKs can request more work but cannot re-enter or overlap that submission loop; remote persistence acknowledgements remain asynchronous. For a deterministic check, I added a same-thread nested-completion regression test. With the replicator implementation from |
Remove the unrelated generic frame-size clamp. Pin the nesting-limit handoff with a blocked executor and document inline completion and entry ownership for both callback families. Assisted-by: Codex
Retain producer ACK demand while a cursor read is pending and keep read backoff as a fallback, restoring progress under continuous transient throttling. Recover owner failures without recursive retries, settle unsent entries, and terminate explicitly if retry scheduling is rejected. Restore continuous fault injection and cover the ACK/failure interleaving, callback-then-throw, rewind and shutdown failures, and deduplication-enabled source-position ordering. Assisted-by: Codex
Preserve queued completion for standalone managed-ledger clients, enable inline callbacks explicitly in the broker, and allow callers to supply an executor. Capture policy at ledger open, bound direct-executor nesting, and release undispatched entries when scheduling is rejected. Assisted-by: Codex
Hold read admission throughout test pauses, drive publish-failure recovery through ProducerSendCallback, and assert cursor progress instead of an incidental self-call. This prevents a message escaping the metrics fixture pause and keeps recovery tests aligned with the owner protocol. Assisted-by: Codex
|
For a normal Pulsar broker, inline read completion is enabled by default. The two defaults apply at different layers:
“Inline” means invoking the successful cursor-read callback directly on the thread that has the result. For a cache hit, that can be the thread requesting the read. With the setting A compatibility use case is broker plugins or protocol handlers that expect callbacks on the managed-ledger worker, for example when coordinating their callback-owned state. Deployments whose plugins use the broker-provided managed-ledger configuration can preserve that behavior by setting “Captured at ledger open” means each managed ledger copies the setting when it opens. Later configuration changes do not alter an already-open ledger. The broker setting is not dynamic; apply a changed broker setting through a broker restart. This policy concerns successful ordinary cursor reads; failure callbacks, single-entry reads and replay callbacks retain their existing behavior. The ledger-wide custom callback executor has been removed. Callers requiring a particular execution context can wrap their callback or use the existing future adapter with an explicit async continuation. BookKeeper handles use the managed-ledger ordering key and the same selected worker, so a completion already on that worker does not need another affinity handoff. The nesting limit is configured globally at JVM startup with Historical validation of the earlier configuration change at |
The third transaction log append automatically rolls the ledger. Forcing LedgerOpened while that rollover is pending can start a second creation, corrupt metadata versions, and fail reopen with BadVersion. Wait for the real rollover and trimming, and verify recovery with a fresh interceptor. Let the metadata store own log cleanup so it does not mask init failures. Validation: 20 repeated cases with four JVM-visible/assigned CPUs and retries disabled. Assisted-by: Codex
Read the profiling harness environment through Gradle providers so a cached Test task cannot reuse the previous scenario or output directory. Validation: changing PULSAR_PROFILING_OUTPUT_DIRECTORY invalidates the configuration cache in consecutive profiling task dry runs. Assisted-by: Codex
The compatibility mode must retain the old executor affinity, including inline completion when already on the ledger worker. Resolve that policy once when opening the ledger, while keeping the recursion trampoline and explicit custom-executor behavior unchanged. Cover off-worker queuing, same-worker completion, and the exact nesting limit with a held ledger worker. Clarify the behavior in the API and broker configuration documentation. Assisted-by: Codex
…epth Keep only the captured inline-versus-ledger-affinity policy. Exceptional callers can select an executor at their callback boundary. Configure the global nesting limit through a startup property and retain rejection cleanup coverage for the ledger handoff. Assisted-by: Codex
…urable Expose replicationMaxReadProcessingStepsPerTurn with the existing default of 64 and a minimum of 1. Capture the non-dynamic setting once per replicator, document the fairness tradeoff, and verify yielding and ownership across multiple budgets. Assisted-by: Codex
Coalesce producer acknowledgments into one queued owner turn, including send-failure cancellation and rewind. Preserve ownership through executor rejection cleanup and verify that producer callbacks never drain cursor work under the producer monitor. Assisted-by: Codex
Test the literal depth property, default, clamping, parsing, and asymmetric broker/library defaults independently. Document legacy Exclusive/Failover handoff and depth-one behavior; log the effective depth once at debug level. Assisted-by: Codex
Terminate when retry submission fails, preserve an accepted timer if producer startup throws, and rewind through the read owner before retrying a rejected completion whose cursor may have advanced. Assisted-by: Codex
Identify the storage-read path exercised by deduplicated backlog replication and remove a coarse read-attempt assertion from the throttling smoke test. Deterministic owner-loop tests cover cached reentrancy and ACK-paced retry. Assisted-by: Codex
Use ForkJoinPool for inline-mode depth overflow while retaining ledger affinity in legacy mode and when common parallelism is at most one. Preserve destination depth checks and add affinity, ownership and callback-chain benchmark coverage. This user-requested executor choice avoids returning cached continuations to the per-ledger queue. Treat it as benchmark-neutral: the short idle-queue screen does not establish a stable performance gain at depth ten, increases allocation, and is slower at depth one. Contended broker-level benefits remain unmeasured. Assisted-by: Codex
…ions Retain ACK demand without scheduling a redundant turn while a read is pending. Rewind rejected read completions before another admission even across disconnection, and contain scheduling or startup errors after releasing ownership so they cannot clear a newer owner. Capture failed-send diagnostic state under the task monitor. Replace the live-broker failed-send fixture with controlled executor coverage, retain the throttling fault-injection assertion, and add deterministic recovery and concurrent ownership regressions. Assisted-by: Codex
Document and log common-pool eligibility, callback affinity, cached-chain continuation, and startup property parsing. Pin the initialized depth and its supported radix forms. Verify a fully cached nested chain returns to the ledger worker across a ledger boundary, use managed blocking in the gated common-pool test, and exercise inline-mode rejected completion when common-pool parallelism selects the ledger fallback. Assisted-by: Codex
|
The latest commits tighten replication recovery and callback ownership:
Validation on the pushed head: 79 focused cases and full The concurrency regressions also have negative controls: temporarily removing the pending-read ACK shortcut, disconnected rewind correction and startup-error containment produced the three intended assertion failures. Restoring the fixes passed all four selected cases, including the Started/Disconnected variants. These checks target those interleavings; CI for the newly pushed commits is separate from the earlier runs. |
|
The depth-limit continuation now uses The reason is to keep ready read completions independent of a saturated managed-ledger executor. When publishing work fully occupies that ordered executor, queuing a continuation back onto it at the nesting limit puts read processing behind the same bottleneck again. The common pool lets those continuations progress on other workers while the ledger executor remains occupied. An idle-executor microbenchmark does not measure the queueing delay this avoids. The depth guard still bounds recursive callbacks, and each continuation checks and unwinds its own thread's depth counter. Legacy mode retains ledger-worker affinity; inline mode also falls back to the ledger executor when effective common-pool parallelism is at most one. The existing JVM pool is reused. Callbacks must remain non-blocking, and inline-mode callers cannot rely on Netty or ledger-worker thread affinity. A recursive cached chain can continue on common-pool workers until a cross-ledger read, cache miss, cursor wait or explicit caller handoff changes its execution context. The added test verifies that a fully cached cross-ledger read returns to the ledger worker without making a storage entry read. The PR description now reflects this behavior. Its broker throughput table remains the previously recorded comparison of inline completion versus ledger-thread affinity; it predates this depth-limit executor change and does not isolate an additional ForkJoinPool throughput gain. |
…lease Contain a synchronous termination failure after retry scheduling fails so an old processing turn cannot discard a newer owner's read result or release its ownership. Add a latch-controlled regression that fails on the previous implementation and verifies exactly-once cleanup by the current owner. Assert disconnected rejected-completion recovery leaves no retry timer, and clarify that the fault-injection counter is a ledger-level smoke check. Assisted-by: Codex
Always expose the inline ledger-fallback rejection case in test results, with an explicit skip and fresh-JVM instructions when the common pool is eligible. Clarify initializer smoke coverage and separate the nesting contract from the legacy Exclusive/Failover handoff documentation. Assisted-by: Codex
Share termination-error containment between retry scheduling and ACK handoff failure paths. A throwing hook must not escape into old-owner cleanup or skip send-callback message recycling. Retain the owner-specific finally cleanup and document the separate failed-termination lifecycle limitation. Add a regression that fails on the previous implementation, tighten terminal drain assertions, and clarify nested-read documentation and fallback skip output. Assisted-by: Codex
Rename rawEntryConfig to defaultConfig and exercise inline read callbacks throughout the existing managed-ledger tests. Keep explicit legacy-affinity overrides in the focused compatibility tests. Retain the full entry history in the generated backlog test, which resets to previously appended positions, so background trimming cannot remove reset targets or race sequential backlog assertions. Assisted-by: Codex
Resolve ManagedLedgerBkTest imports and migrate the new upstream cursor-untracking test to defaultConfig. Preserve the upstream cache test synchronization and all assertions.
Motivation
When many producers publish to one persistent topic, publish completions can saturate its ordered
managed-ledger executor. Successful cached cursor reads then queue behind publishing work, although
the entries are already confirmed and available. Consumers cannot keep up, and a large subscription
backlog accumulates. Additional broker cores do not remove that per-topic queue.
#26619 bypassed this handoff for the Exclusive/Failover dispatcher's future adapter. This change
generalizes inline completion and removes that temporary opt-in API. In the fresh
500-producer/20-consumer Key_Shared comparison against master, consumers kept up with publishing:
median steady egress (consume) increased from 5.7k to 153.0k messages/s.
Median whole-run consumer throughput increased by 86.7%, and median sampled
maximum backlog fell from 11.37 million to 17,693 messages.
The performance section records every run and its revision. These results describe this workload
and configuration, not a general capacity multiplier.
The handoff also incidentally serialized replication batch submission. Removing it exposes a
correctness problem: a producer ACK can request another read before the previous batch finishes
submitting messages, allowing later source positions to overtake earlier ones. With destination
deduplication, a lower position can then receive a successful duplicate acknowledgment and be
deleted at the source despite never being stored at the destination. This PR therefore includes
the prerequisite broker fix: one owner coordinates reads, complete batch submission and recovery.
Modifications
OpReadEntry.complete()removes the executor-thread restriction on inline completion. Before#26619's opt-in, the ordinary callback path selected a depth counter only on the ledger thread:
With inline completion enabled, every completing thread uses the
FastThreadLocaldepth counter. Inline completion incrementsit and decrements it in
finally. At the configured nesting limit (ten by default), inline mode queuesone completion through
ForkJoinPool.commonPool()so cached chains can unwind without returning to theper-ledger queue. Legacy mode retains the managed-ledger executor. Common-pool parallelism of one or
lower also uses the ledger executor, avoiding disabled-pool liveness problems. Independent reads do not accumulate depth.
This explicitly reuses the JVM common pool; no additional pool is created. Its workers are daemon
threads shared with other JVM users. Managed-ledger callbacks must remain non-blocking and cannot
rely on Netty or ledger-worker affinity in inline mode. A recursive cached chain may continue on
common-pool workers until a cross-ledger read, cache miss, cursor wait or caller handoff changes its
execution context. Startup DEBUG logging records the effective depth and whether common-pool
unwinding is selected. Common-pool parallelism normally uses available processors minus one
(at least one), configurable with
-Djava.util.concurrent.ForkJoinPool.common.parallelism.The common pool keeps depth-limit continuations independent of a saturated managed-ledger executor.
Queuing them back on that executor would reintroduce its publishing backlog as a bottleneck whenever
the nesting limit is reached. Using the existing common-pool workers lets those ready callbacks
progress while the ledger executor is occupied. The relevant benefit is avoiding that queueing delay;
an idle-queue scheduling benchmark does not measure it. The short JMH comparison has no stable
advantage at depth ten, allocates more, and is slower at depth one; it does not establish the size of
the improvement under executor saturation.
OpScanis a concreterecursive caller. Exclusive/Failover keeps its asynchronous dispatcher continuation and
Shared/Key_Shared coalesces read requests; this change does not claim their normal dispatch paths
reach the nesting limit. Legacy mode remains available when ledger-worker affinity is required.
The depth guard is necessary for callers such as
OpScan, whose callback can immediately requestthe next cached batch.
Pulsar brokers enable inline read completion by default. There are two configuration layers:
managedLedgerReadEntriesCallbackInlinetrueManagedLedgerConfig.readEntriesCallbackInlinefalsepulsar.managedLedger.maxReadCompletionDepth10Inline completion means invoking the successful cursor-read callback directly on the thread that
has the result. For a fully cached read, this can be the thread requesting the read, before that
read method returns. With the setting
false, the callback runs on the managed-ledger executor:it runs directly when already on that worker, within the nesting limit, and queues when completing
on another thread. The nesting guard applies to callbacks that issue another read before returning.
The compatibility setting is useful for broker plugins and protocol handlers that expect the
previous managed-ledger callback-thread affinity. Such code may rely on callbacks running on the
ledger worker when accessing its own state. Deployments using the broker-provided managed-ledger
configuration can retain that behavior with:
managedLedgerReadEntriesCallbackInline=falseA plugin that constructs its own
ManagedLedgerConfigalready gets the library'sfalsedefault,unless it explicitly enables inline completion. This is a compatibility option for code that needs
the previous worker affinity; it does not establish FIFO completion of concurrent reads or promise
that callbacks are always asynchronous. Legacy mode already allows bounded inline callbacks on the
ledger worker.
Each managed ledger copies the completion policy when it opens. Changing the configuration object
later, including through
setConfig, does not change an already-open ledger. The broker setting isnot dynamic; apply a changed broker setting through a broker restart.
The policy covers successful ordinary
asyncReadEntriesandasyncReadEntriesOrWaitcallbacks;failure, single-entry and replay callbacks retain their existing behavior. There are no per-read
options, ledger-wide custom executor or new executor pools. The common inline path checks a final
boolean; the global depth limit is read once at class initialization.
A caller needing a special execution context can wrap its callback or use the existing
ManagedLedgerUtilsfuture adapter with an explicit async continuation such asthenComposeAsync(..., executor). The caller remains responsible for releasing entries, includingif its handoff is rejected. This does not require routing every read through a configurable executor.
BookKeeper handles opened with the managed-ledger name as their ordering key use the same worker
as the managed-ledger executor, including batch-read response processing. The compatibility mode
therefore runs inline when completion is already on that worker. Disabling broker inline completion also
restores the Exclusive/Failover cache-hit handoff that existed before #26619. Cache hits can complete on other
initiating threads; the boolean controls whether they need the ledger handoff. Worker affinity
does not guarantee FIFO completion of concurrent read requests.
The queued continuation checks the depth again on its destination thread. A depth limit of one queues
every subsequent completion in a nested cached chain. Rejected completion scheduling releases
undispatched entries and fails the read once; the cursor can already have advanced. The replicator
requests a rewind through its existing owner before admitting another read after such a rejection.
With the broker's inline default, the changed and unchanged paths are:
asyncReadEntries*, without a BookKeeper entry-read requestbookkeeper-ml-schedulerthreadPendingReadsManagerand the ledger executorShared/Key_Shared therefore avoids the ledger queue for eligible full cache hits. Its dispatch work
can run on the worker initiating the read; this moves CPU work rather than eliminating it. Existing
subscription-thread dispatch configuration and explicit caller continuations still apply.
Exclusive/Failover retains
whenCompleteAsyncon its dispatcher executor. This does not removeevery managed-ledger scheduling boundary or make storage reads complete directly on arbitrary callers.
Replication now records read demand and lets one owner finish submitting each batch before issuing
the next cursor read. ACK callbacks continue to free producer capacity concurrently, but schedule a
coalesced owner turn instead of running cached-read processing while the geo producer monitor is held.
Failed sends publish cancellation and rewind through the same owner. A task is
reusable only after both submission and producer callbacks have completed. Recovery coordinates
cancellation and rewind, discards stale results, and releases entries that were never handed to the
producer. Read failures retain a fallback timer; ACK demand can resume progress before that timer,
including when the ACK arrived while the failed read was pending. A failure alone does not trigger
an immediate retry loop. Recoverable processing errors retain a live replicator; a closed cursor
or failed retry submission terminates it explicitly. A producer-start exception after successful timer
submission leaves that timer intact; it is handled separately from a scheduling failure.
Synchronous termination-hook failures are contained on both retry and ACK handoff failure paths,
preserving current read ownership and allowing send callbacks to recycle their messages.
This is one required correctness fix alongside the completion change, with separate broker fix
commits and tests. The replication fix adds no prefetch/reorder queue or thread pool. The broker
setting
replicationMaxReadProcessingStepsPerTurnlimits work before yielding to the existingbroker executor while retaining ownership. It defaults to 64, requires a value of at least 1 and
is not dynamic. Each step initiates a read, processes a completed batch, or handles cancellation
or rewind; it does not count messages. Lower values favor fairness between tasks, while higher
values reduce scheduling overhead. The setting is documented in both broker and standalone
configuration files. Read sizing and its existing average-entry-size estimate remain unchanged.
The remaining changes remove
ReadEntriesCallback.canExecuteOnAnyThread()and the booleanManagedLedgerUtils.readEntriesWithSkipOrWaitoverload from #26619, document callback ownershipand reentrancy, update raw-entry test fixtures, and add keyed producers and a reproducible
500-producer/20-consumer Key_Shared YAML scenario to the profiling harness.
Correctness, ordering and compatibility
Read visibility does not change. Active-ledger reads remain bounded by the managed ledger's
confirmed position, closed-ledger reads use BookKeeper's LAC, and the topic's maximum readable
position still constrains transaction visibility. Entries within a result remain in cursor order,
and the read position advances before successful completion. None of these properties requires
processing the completed result on a particular thread.
That entry order does not imply FIFO completion of independent asynchronous reads. BookKeeper
does not guarantee completion order across concurrent read operations. Callers advancing one
cursor must coordinate the next read with processing or safe ownership transfer of the previous
result. Dispatchers retain their pending-read guards, conflation, synchronization and replay rules;
the replicator now retains ownership through the entire submission loop, including reentrant ACKs.
The preceding implementation restricted ordinary successful cursor callbacks to the ledger executor,
allowing bounded inline completion when already on that worker. The default standalone
ManagedLedgerConfigpreserves that affinity and bounded inline behavior; brokersexplicitly enable inline completion and can opt out. When inline completion is enabled, external
users of these
LimitedPrivateAPIs can observe callbacks within the initiating call, reentrantly,on a different worker. Some failures
and single-entry cache reads already had inline paths. Callers must initialize pending state
before invoking the read, avoid blocking callbacks, release returned entries, and schedule their
own continuation when executor affinity is required. The interface and future-adapter Javadocs
state those responsibilities; this is not a promise supporting overlapping active cursor reads.
Verifying this change
Full
./gradlew quickCheckpassed. Focused correctness runs use four JVM-visibleprocessors and four assigned CPUs, with test retries disabled. They cover cached completion and depth-counter unwinding, the exact nesting-limit
handoff through the common pool in inline mode and the ledger executor in legacy mode, explicit
asynchronous continuations, and cached scans. A fully cached chain crossing two ledgers verifies
that the cross-ledger continuation completes on the ledger worker without a storage entry read. The inline overflow test gates its callback, verifies
its destination, exactly-once delivery and released entry ownership, and allows an independent
read to complete. A low-parallelism common pool uses the ledger fallback.
Manual fresh-JVM runs with
JAVA_TOOL_OPTIONS, outside the automated suite, test common-poolparallelism properties zero and one, and a nesting limit of one.
Parallelism one exercises inline-mode ledger-fallback rejection cleanup. The zero-property run
still uses the common pool on JDK 25 because
CompletableFutureinitialization raises its effectiveparallelism to two. The inline-fallback test is an explicit skipped case when the test JVM has
common-pool parallelism greater than one, so normal CI reports this coverage gap.
Replication tests use a controlled executor for failed-publish permit recovery,
retain ACK demand without queued work while a read is pending, rewind rejected completion across
disconnection without a leftover timer, and prevent an old producer-start or termination error
from releasing a newer owner's state or discarding its read result.
The ACK-handoff rejection test also injects a termination-hook error and verifies message recycling
and entry release without cursor work on the producer callback thread.
The completion-policy tests cover legacy ledger-thread affinity, including inline completion on the ledger worker and a queued
depth-limit handoff, fixed policy across configuration updates in either mode, and buffer release
on rejected ledger scheduling without double-decrementing pending-read accounting. Earlier manual
fresh-JVM runs also passed with a configured depth of two and with zero or negative values clamped to one.
Replication turn-budget tests cover 1, 2, 3, 64 and 128 steps, verifying exact yield boundaries,
retained ownership while queued, ordered delivery and non-recursive cached reads. Configuration
loading rejects zero and negative turn budgets. Replication tests cover concurrent and reentrant
ACKs, submission order, task reuse, cancellation,
rewind and failures, stale/unsent entry cleanup, and remote deduplication with a payload and source
ledger/entry-position oracle. The continuous read-throttling fixture retains alternating failures
throughout the 10,000-message transfer. Recovery fixtures invoke the actual failed send callback
and verify cursor activity and end-to-end delivery; the metrics pause helper keeps admission
closed until resume. The previously failing metrics case passed five repeated invocations. Transaction sequence recovery
now awaits automatic ledger rollover and trimming before reopening with a fresh sequence generator;
both recovery variants passed ten repeated invocations each. The profiling task also tracks scenario
and output environment overrides as Gradle configuration-cache inputs.
Two supporting test fixes are intentionally retained in this PR: the profiling task declares its
environment overrides as configuration-cache inputs so repeated runs use the selected scenario,
and transaction sequence-recovery tests await real ledger rollover/trimming instead of mutating
private state through reflection. They are independent of the production callback scheduling change.
Performance comparison
Fresh comparison on September 18, 2026, against master production code, with the latest PR implementation including the replication fixes, ACK-turn scheduling and ForkJoinPool depth-limit handoff. All five runs completed.
origin/masterwhen the series started)ea97c3c7e50ba9a368a076ade24ab502415e176ea257bf8eb905d2d9f2220ea3eedbd6195daabd87dfc6bf1e9be5d5895bad351e6340dab8e54c5c42The baseline adds only this PR's profiling harness changes to master; no baseline production code is modified. Both revisions use the same scenario and configuration defaults, including
managedLedgerCacheEvictionExtendTTLOfRecentlyAccessed=false. This compares the complete PR to master, rather than toggling completion policy on the PR branch.The workload uses one persistent topic, 500 isolated V4 producer clients/connections, 20 isolated consumer clients on one Key_Shared subscription, shared client resources, random keys, unbatched 128-byte messages, unrestricted production and 40 outstanding sends per producer. Every completed run received 12 million messages with zero failed ACKs. Asynchronous sends overshot that target slightly.
Runs alternate baseline → PR → baseline → PR → baseline. The table below uses medians of the completed runs in each group (the median of two runs is their midpoint).
Individual runs, in execution order:
Steady rates use topic counter deltas after a 20-second warmup from the first snapshot with all 500 producers, and end before producers disconnect. Only complete intervals with all 500 producers at both endpoints are included; interval durations weight the reported rate. Egress measures dispatch to consumers. Whole-run rates use each client's aggregate timer, including startup, reporting and backlog drain. Sampled maximum backlog uses snapshots roughly six seconds apart and can miss intervening peaks.
Before every run, the host had less than 5% CPU work over a five-second sample, more than 12,000 MiB available memory, no running workload containers, performance CPU governors and temperature below 65°C. Gradle daemons were stopped between runs. The host has 8 cores / 16 logical CPUs with SMT and turbo enabled; benchmark JVMs were not restricted to four CPUs. The four-worker Gradle setting limits build workers only. All completed runs used zero swap and recorded no broker OOM or restart.
Thermal throttling was observed in both baseline and PR runs. Alternating runs and cooling between them help reveal drift; they do not eliminate throttling. These results show the effect of the saturated per-topic executor in this workload, not a general broker capacity multiplier or a fixed-rate latency comparison.
The final baseline's steady window recorded no increase in its package thermal-throttle counter, yet steady egress remained 5.8k messages/s and sampled maximum backlog reached 11.42 million. The repeated gap therefore persists in a baseline window without recorded throttling. Throttle counters do not measure the duration or severity of throttling.
The scenario and inherited harness defaults include ensemble/write/ACK quorum one, disabled bookie journal syncing,
dispatcherMaxReadBatchSize=1000,preciseDispatcherFlowControl=true, and both dispatcher retry-backoff settings at zero. The complete configuration is inkey-shared-500x20.yaml, its inheritedread-completion-isolation.yaml, andPulsarProfilingConfig.defaults().Reproduction
Create separate worktrees at the production revisions above. Copy only the following harness files from the PR revision onto the master worktree, leaving its production sources unchanged:
Build the profiling image once in each worktree before collecting measurements, using a distinct
RUN_IMAGE_TAGfor baseline and PR to prevent image reuse across revisions:After the host is idle and cooled, run the following in the corresponding worktree, changing
RUN_NAMEfor each repetition. Repeat baseline, PR, baseline, PR, baseline, and let the host cool between runs:These options collect CPU, lock and allocation recordings without wall-clock sampling. Profile analysis ran only after all five workloads completed.
JFR findings
After all five workload runs finished, the recordings were analyzed with Jafar MCP
jfr_diagnose, CPU/lock/parkjfr_stackprofile, and windowed queries. The figures below use a 30-second window starting six seconds into each run's steady interval; independentjfrconvand JDK JFR exports were used to check the interpretation. Whole-recording diagnostic and stack-profile views also contain startup and shutdown, so they are not used as steady-window measurements. CPU samples use a 10 ms interval; allocation figures are sampled estimates.The publishing ledger worker remains close to one fully occupied CPU in both modes. The improvement comes from allowing eligible cached-read completion to proceed without waiting for that worker. In the baseline, all sampled
OpReadEntry.completeNowstacks in these windows run on the ledger worker; with this PR, most run onpulsar-ioworkers, with smaller ledger-worker and managed-ledger scheduler contributions. This is consistent with consumers keeping up while the publishing worker remains busy. No read-completion stack in these windows was sampled on a common-pool worker, so these results do not isolate a benefit from the depth-limit ForkJoinPool fallback.The higher dispatch rate comes with more allocation per second and more dispatcher-monitor contention, principally at
internalConsumerFlowon broker I/O threads. Cumulative contention sums waits across threads and can exceed the 30-second observation window. The baseline dispatches very few messages during that window, so it is not an equal-work comparison and does not show a per-message allocation or locking regression. The change removes the ledger-queue dependency for eligible reads; it does not eliminate dispatcher contention. Cache-eviction and GC-worker CPU samples decrease while consumers keep up. Maximum individual stop-the-world pauses were below 0.13 ms in all five selected windows, so long GC pauses do not explain the baseline's dispatch deficit.A separate, earlier executor-handoff JMH screen measured 7.21 → 4.74 µs/op and 280 → 152 B/op with a direct continuation (one fork, 2 × 1s warmup, 3 × 1s measurement, GC profiler). It measures idle executor/future overhead, not storage latency or the throughput improvement above.
Does this pull request potentially affect one of the following parts: