Skip to content

[improve][ml] Allow inline read completion and preserve replication order - #26620

Merged
merlimat merged 25 commits into
masterfrom
lh-improve-inline-read-completion
Sep 18, 2026
Merged

merlimat merged 25 commits into
masterfrom
lh-improve-inline-read-completion

Conversation

@lhotari

@lhotari lhotari commented Sep 17, 2026 •

Copy link
Copy Markdown
Member

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:

int[] depth = executor.isCurrentThread() ? INLINE_COMPLETION_DEPTH.get() : null;
if (depth != null && depth[0] < MAX_NESTED_INLINE_COMPLETIONS) {
    depth[0]++;
    try {
        completeNow(ctx);
    } finally {
        depth[0]--;
    }
} else {
    executor.execute(() -> completeNow(ctx));
}

With inline completion enabled, every completing thread uses the FastThreadLocal depth counter. Inline completion increments
it and decrements it in finally. At the configured nesting limit (ten by default), inline mode queues
one completion through ForkJoinPool.commonPool() so cached chains can unwind without returning to the
per-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. OpScan is a concrete
recursive 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 request
the next cached batch.

Pulsar brokers enable inline read completion by default. There are two configuration layers:

Configuration Default Behavior
Broker managedLedgerReadEntriesCallbackInline true The broker copies this value into the managed-ledger configuration for its topics, enabling inline completion by default
ManagedLedgerConfig.readEntriesCallbackInline false The library default preserves the previous callback-thread affinity for applications that construct their own managed-ledger configuration; broker-created configurations use the broker setting above
JVM property pulsar.managedLedger.maxReadCompletionDepth 10 Maximum nested inline completions across managed ledgers; set at JVM startup; values below 1 use 1

Inline 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=false

A plugin that constructs its own ManagedLedgerConfig already gets the library's false default,
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 is
not dynamic; apply a changed broker setting through a broker restart.

The policy covers successful ordinary asyncReadEntries and asyncReadEntriesOrWait callbacks;
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
ManagedLedgerUtils future adapter with an explicit async continuation such as
thenComposeAsync(..., executor). The caller remains responsible for releasing entries, including
if 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:

Read path Completion thread / scheduling
Full cache hit with ledger handle and permits available The caller can complete the callback within asyncReadEntries*, without a BookKeeper entry-read request
Full cache hit from a delayed new-entry recheck Can complete on the factory's bookkeeper-ml-scheduler thread
Storage read or partial cache hit Still passes through PendingReadsManager and the ledger executor
Permit deferral, waiting-read wakeup and cross-ledger continuation Retain their existing asynchronous boundaries
Nesting-limit handoff Queues through the common ForkJoinPool; uses the ledger executor in legacy mode or when common-pool parallelism is one or lower

Shared/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 whenCompleteAsync on its dispatcher executor. This does not remove
every 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 replicationMaxReadProcessingStepsPerTurn limits work before yielding to the existing
broker 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 boolean
ManagedLedgerUtils.readEntriesWithSkipOrWait overload from #26619, document callback ownership
and 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
ManagedLedgerConfig preserves that affinity and bounded inline behavior; brokers
explicitly enable inline completion and can opt out. When inline completion is enabled, external
users of these LimitedPrivate APIs 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

  • Make sure that the change passes the CI checks.

Full ./gradlew quickCheck passed. Focused correctness runs use four JVM-visible
processors 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-pool
parallelism 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 CompletableFuture initialization raises its effective
parallelism 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.

Source Revision
Baseline production (origin/master when the series started) ea97c3c7e50ba9a368a076ade24ab502415e176e
Baseline with matching profiling harness only a257bf8eb905d2d9f2220ea3eedbd6195daabd87
PR dfc6bf1e9be5d5895bad351e6340dab8e54c5c42

The 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).

Metric Master baseline PR Improvement
Steady ingress (produce) 139,464 msg/s 152,971 msg/s +9.7%
Steady egress (consume) 5,738 msg/s 153,002 msg/s +2566.6%
Whole-run producer throughput 132,828 msg/s 132,798 msg/s -0.02%
Whole-run consumer throughput 76,983 msg/s 143,732 msg/s +86.7%
Sampled maximum backlog 11,366,885 messages 17,693 messages 99.84% lower

Individual runs, in execution order:

Run Steady ingress (produce), msg/s Steady egress (consume), msg/s Whole-run consumer throughput, msg/s Sampled maximum backlog, messages Steady window
1. Baseline 1 139,150 5,567 76,983 11,250,179 67s
2. PR 1 152,406 152,548 143,454 17,913 49s
3. Baseline 2 139,464 5,738 75,228 11,366,885 61s
4. PR 2 153,536 153,456 144,010 17,473 49s
5. Baseline 3 140,575 5,770 81,217 11,415,246 61s
image

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 in key-shared-500x20.yaml, its inherited read-completion-isolation.yaml, and PulsarProfilingConfig.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:

git restore --source=dfc6bf1e9be5d5895bad351e6340dab8e54c5c42 -- \
  tests/integration/build.gradle.kts \
  tests/integration/src/test/java/org/apache/pulsar/tests/integration/profiling/AbstractPulsarProfilingTest.java \
  tests/integration/src/test/java/org/apache/pulsar/tests/integration/profiling/PulsarProfilingConfig.java \
  tests/integration/src/test/java/org/apache/pulsar/tests/integration/profiling/PulsarProfilingConfigTest.java \
  tests/performance/README.md \
  tests/performance/scenarios/key-shared-500x20.yaml

Build the profiling image once in each worktree before collecting measurements, using a distinct RUN_IMAGE_TAG for baseline and PR to prevent image reuse across revisions:

RUN_IMAGE_TAG=pr26620-baseline # use pr26620-inline in the PR worktree
./gradlew --no-daemon --max-workers=4 \
  :tests:java-test-image:dockerBuildWithAsyncProfiler \
  -Pdocker.tag="$RUN_IMAGE_TAG"

After the host is idle and cooled, run the following in the corresponding worktree, changing RUN_NAME for each repetition. Repeat baseline, PR, baseline, PR, baseline, and let the host cool between runs:

RUN_NAME=baseline-1 # then pr-1, baseline-2, pr-2, baseline-3
./gradlew --stop
PULSAR_PROFILING_CONFIG="$PWD/tests/performance/scenarios/key-shared-500x20.yaml" \
PULSAR_PROFILING_OUTPUT_DIRECTORY="$PWD/tests/integration/build/pulsar-profiling/$RUN_NAME" \
./gradlew --no-daemon --max-workers=4 \
  :tests:integration:profilingIntegrationTest --tests '*PulsarProfilingV4Test' \
  -x :tests:java-test-image:dockerBuildWithAsyncProfiler \
  -Pdocker.tag="$RUN_IMAGE_TAG" \
  -Pinttest.asyncprofiler.opts=event=cpu,interval=10ms,lock=0,alloc=2m,jfrsync=profile

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/park jfr_stackprofile, and windowed queries. The figures below use a 30-second window starting six seconds into each run's steady interval; independent jfrconv and 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.

Broker profile metric (median across runs) Master baseline PR
Active ledger-worker CPU samples / 30s 2,966 2,915
Average process CPU, logical-core equivalents 6.03 5.21
Cache-eviction worker CPU samples / 30s 1,690 908
GC-thread CPU samples / 30s 3,851 2,240
Estimated allocation rate 127 MiB/s 292 MiB/s
Cumulative dispatcher-monitor contention / 30s 0.073 thread-seconds 54.12 thread-seconds
Total stop-the-world GC pause time / 30s 0.94 ms 1.21 ms

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.completeNow stacks in these windows run on the ledger worker; with this PR, most run on pulsar-io workers, 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 internalConsumerFlow on 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:

  • Dependencies (add or upgrade a dependency)
  • The public API — LimitedPrivate callback scheduling/ownership contract and removal of [improve][broker] Avoid redundant read-completion executor handoffs for Exclusive/Failover #26619's temporary opt-in API.
  • The schema
  • The default values of configurations — broker inline completion enabled; standalone managed-ledger executor affinity preserved.
  • The threading model — bounded inline read completion and serialized replication submission.
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment — a non-dynamic broker switch can restore legacy executor affinity.

@void-ptr974

Copy link
Copy Markdown
Contributor

One possible interleaving in PersistentReplicator, shared by Geo and Shadow replication:

  1. The current read callback sets inFlightTask.entries and submits entry 0 from batch [0, 1], but has not yet submitted entry 1.
  2. If the ACK for 0 triggers readMoreEntries() at this point, available permits and a cached read may allow the next batch [2, 3] to complete inline on the ACK thread.
  3. That callback may submit entries 2 and 3 before the original callback submits entry 1.

This appears to leave room for submission order 0 → 2 → 3 → 1, without reaching the nesting limit. The concern is whether the existing guards prevent a later batch from being submitted while the current batch is still being submitted.

@void-ptr974

Copy link
Copy Markdown
Contributor

Another potential ordering case is an inline completion overtaking a queued completion:

  1. After ten nested inline completions, the next read advances the cursor, but its callback C11 is queued on the managed-ledger executor.
  2. Before C11 runs, another thread issues a read on the same cursor.
  3. If that read completes inline from cache, C12 may run before C11.

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.

@void-ptr974

Copy link
Copy Markdown
Contributor

A separate design consideration is avoiding recursive read callbacks in continuous readers such as PersistentReplicator.

Its current callback can call readMoreEntries() directly. If the next read completes synchronously from cache, the next callback can start before the current one returns. Consecutive cached batches, including batches of expired messages skipped by Geo replication, may repeatedly follow this path.

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.

Base automatically changed from lh-improve-read-completion to master September 17, 2026 16:58
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
@lhotari
lhotari force-pushed the lh-improve-inline-read-completion branch from b5ba2dd to b8a3af6 Compare September 17, 2026 16:58
@lhotari
lhotari marked this pull request as draft September 17, 2026 17:22
@lhotari

lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

@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
@lhotari

lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

@void-ptr974 The follow-up commit 93204782decf, on top of b8a3af6, addresses these concerns:

  1. An ACK starts the next batch before the current batch has finished submitting. PersistentReplicator now has one processing owner covering both cursor invocation and the complete batch-submission loop. Concurrent ACKs record read demand; they cannot start another submission loop. The next read is admitted only after the preceding batch finishes submission, while its producer ACKs can remain outstanding. A separate submissionComplete flag also prevents the final ACK from making a task recyclable before the submission loop returns. Geo and Shadow replication both use this coordination.

    The deterministic regression test pauses batch [0,1] after submitting 0, then delivers its ACK concurrently with a cached next batch [2,3]. Against b8a3af6, it reproduces submission of 2 and 3 while the first batch is still paused before 1. With the fix, submission remains 0,1,2,3.

  2. An inline completion overtakes a queued completion. The replicator retains one outstanding cursor read and keeps that reservation until its result has been processed through batch submission. A queued callback therefore keeps the next read from being admitted; an ACK cannot issue a second read around it. Cancellation and rewind also go through the same owner, including cancellation requested between reserving a read and invoking the cursor.

    This does not introduce a general FIFO guarantee for concurrent reads on a cursor. The ManagedCursor Javadoc now makes the caller's coordination responsibility explicit: callbacks can complete inline or on another thread, and entry order or callback-thread affinity does not establish callback ordering or exclusive access to caller state.

  3. Use a processing loop to avoid recursive continuous reads. Adopted this approach. Inline callbacks publish their result and return to the existing owner; concurrent and reentrant requests are coalesced into pending demand. The processing loop keeps ownership across the whole submission, with a bounded turn that yields to the existing broker executor for fairness. The managed-ledger nesting guard remains for other callers.

    The cached-read test exercises 256 consecutive inline completions and verifies callback depth remains one, active tasks are not recycled, and only one cursor read remains pending. Additional tests cover stale results, cancellation, termination, retry backoff, and unsent-entry cleanup.

The scoped tests pass locally, including all five invocations of testReplicationWithSchema. The new CI run is still in progress.

@lhotari

lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

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
@lhotari

lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

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, havePendingRead stayed true until replicateEntries(entries) returned, covering submission of the entire batch. Producer callbacks also checked that flag before starting another read. This prevented a later batch from overtaking the current batch without requiring BookKeeper to complete concurrent reads in FIFO order. See the old batch-completion guard and producer callback check.

#24189 changed the pending-read condition to inFlightTask.entries == null, while setting those entries at callback entry, before replicateEntries. Another read could therefore start while the current batch was still being submitted. See callback entry and read admission. However, successful read callbacks were still always queued on the managed-ledger executor. That serialized batch submission, so #24189 alone does not establish this particular reordering bug at its merge revision.

#26599 subsequently allowed read completion inline when already on the ledger thread. This permits the following sequence even before #26620:

  1. Batch A contains entries [0, 1] and starts submission on the ledger thread.
  2. Entry 0 exceeds the replication producer's message-size limit. The producer rejects it with InvalidMessageException and invokes the send callback synchronously. This is possible with different message-size limits at the source and destination.
  3. The replicator treats this exception as terminally handled and can request another read, provided queue permits and producer writability allow it. This path does not rewind.
  4. Cached batch B [2, 3] completes inline and submits 2 and 3 before A resumes and submits 1. The valid entries are therefore submitted as 2, 3, 1. Limiting recursion depth does not prevent this shallow nesting.

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 a846a45953d8 (the parent of the original #26620 change), it fails with submission order [0, 2, 3, 1]; with the fix, it passes with [0, 1, 2, 3]. The test uses a controlled inline cursor callback and synchronous completion to isolate batch ordering. It is not an end-to-end reproduction of the message-size rejection; that trigger is established by the source paths above.

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
@lhotari lhotari changed the title [improve][ml] Remove managed-ledger executor handoff for read completion [improve][ml] Remove read-completion handoffs and preserve replication order Sep 17, 2026
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
@lhotari

lhotari commented Sep 17, 2026 •

Copy link
Copy Markdown
Member Author

For a normal Pulsar broker, inline read completion is enabled by default. The two defaults apply at different layers:

  • The managed-ledger library's ManagedLedgerConfig.readEntriesCallbackInline defaults to false. This preserves the previous callback-thread affinity for applications that construct their own managed-ledger configuration.
  • The broker's managedLedgerReadEntriesCallbackInline defaults to true. The broker copies this value into the managed-ledger configuration used for its topics, enabling the new behavior by default.

“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 false, the callback stays on the managed-ledger executor: it can run directly when already on that worker, within the nesting limit, and otherwise it is queued there.

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 managedLedgerReadEntriesCallbackInline=false. Plugins that construct their own ManagedLedgerConfig already get false unless they explicitly enable inline completion. This preserves worker affinity; it does not guarantee FIFO completion of concurrent reads or make every callback asynchronous.

“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 -Dpulsar.managedLedger.maxReadCompletionDepth=10. The default is ten, with values below one clamped to one. Both completion modes bound callbacks that recursively request another read. In the current implementation, inline-mode overflow uses the common ForkJoinPool when its effective parallelism exceeds one, keeping those continuations independent of a saturated managed-ledger executor. Legacy mode and the low-parallelism fallback use the ledger executor.

Historical validation of the earlier configuration change at 6fcf863a8ffe: 13 focused default cases and quickCheck passed with four CPU cores and retries disabled. Manual fresh-JVM runs also passed 11 policy cases each with depth-property values 2, 0 and -1, using fixed expected depths 2, 1 and 1. Rejection tests exercised the ledger handoff, including the depth-limit queue, and checked released entries and balanced read accounting. All 40 main-workflow jobs and both flaky-workflow jobs passed on that revision, including CodeQL; the Proxy job passed on an unchanged-commit rerun after a native Conscrypt finalizer crash. These are historical results, not a CI claim for later commits.

@lhotari lhotari changed the title [improve][ml] Remove read-completion handoffs and preserve replication order [improve][ml] Allow inline read completion and preserve replication order Sep 17, 2026
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
@lhotari
lhotari marked this pull request as ready for review September 18, 2026 06:08
@lhotari
lhotari marked this pull request as draft September 18, 2026 09:30
…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
@lhotari

lhotari commented Sep 18, 2026

Copy link
Copy Markdown
Member Author

The latest commits tighten replication recovery and callback ownership:

  • Producer ACKs record demand and queue a coalesced processing turn, so cached-read processing and cursor recovery do not run while the geo producer's monitor is held. If a read is already pending and no cancellation or rewind is needed, its completion consumes the recorded demand without an extra executor task. See c765056cebbb and 20b8a96a9daa.
  • Failed sends publish cancellation and rewind together. A rejected read-completion handoff requests a rewind before another read is admitted, including when the replicator has become disconnected. Retry scheduling failures terminate explicitly; a producer-start failure preserves an already accepted retry timer. Errors after releasing processing ownership cannot clear a newer owner's state. See a3f2c3429397 and 20b8a96a9daa.
  • Replaced the timing-dependent failed-publish test with a controlled executor fixture. It checks entry release, task completion and restored permits before draining recovery, then verifies rewind precedes the next read. Failed-send diagnostics now capture a task count under the monitor instead of formatting a concurrently modified task list.
  • Added configuration initialization/default tests, explicit fault-injection checks, cached cross-ledger completion coverage and inline-mode rejection coverage for the low-parallelism ledger fallback. Callback Javadoc and both broker configuration files describe common-pool execution and the global depth property. See e006cf850823.

Validation on the pushed head: 79 focused cases and full quickCheck passed, with four JVM-visible processors and retries disabled. Fresh JVM affinity suites passed with common-pool parallelism properties zero and one, and with a nesting limit of one; parallelism one exercises inline-mode ledger-fallback rejection. On JDK 25, CompletableFuture initialization raises configured zero parallelism to two, so that run still exercises the common pool.

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.

@lhotari

lhotari commented Sep 18, 2026

Copy link
Copy Markdown
Member Author

The depth-limit continuation now uses ForkJoinPool.commonPool() when inline completion is enabled and common-pool parallelism exceeds one. See 90f1ecf78a47.

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. OpScan is a concrete recursive cached-read caller; Exclusive/Failover retains its explicit whenCompleteAsync continuation, and Shared/Key_Shared conflates read requests.

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
@lhotari
lhotari marked this pull request as ready for review September 18, 2026 16:17
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.

@merlimat merlimat left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍

@merlimat
merlimat merged commit 8b74450 into master Sep 18, 2026
48 checks passed
@lhotari lhotari added the category/performance Performance issues fix or improvements label Sep 28, 2026
@lhotari
lhotari deleted the lh-improve-inline-read-completion branch October 1, 2026 00:02
@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

category/performance Performance issues fix or improvements ready-to-test

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants