Skip to content

[fix][ml] Fix active cursor retention and deferred-update synchronization - #26511

Merged
lhotari merged 6 commits into
apache:masterfrom
lhotari:lh-fix-active-cursor-retention
Sep 9, 2026
Merged

lhotari merged 6 commits into
apache:masterfrom
lhotari:lh-fix-active-cursor-retention

Conversation

@lhotari

@lhotari lhotari commented Sep 9, 2026 •

Copy link
Copy Markdown
Member

Replacement for #26506.

Motivation

ActiveManagedCursorContainerImpl defers cursor ordering until a position query needs it. With cacheEvictionByExpectedReadCount enabled, cursor creation and deletion can continue without such queries. Removed nodes then accumulate in pendingPositionUpdates or remain in the ordered list and pendingRemovedCursors, retaining their cursor objects and associated object graphs.

This PR replaces the approach in #26506 with immediate cursor-reference release and bounded, batched node compaction. The ten-event flush in #26506 counts cursors removed before their first position flush; already-tracked removals are outside that cleanup trigger. Preserving lazy position updates is important to the efficiency of this data structure.

This also fixes a pre-existing synchronization bug: getSlowestCursorPosition() held a read lock while calling processPendingPositions(), which mutates the pending-update collection, linked list, and counters. Concurrent readers could therefore flush the same state simultaneously. The internal state-checking method had the same locking problem. This is API-correctness hardening: under the current broker configuration paths, this implementation is selected when cacheEvictionByExpectedReadCount is enabled, and the broker skips its slowest-position lookup in that mode. The production rank-query path already used the write lock.

Modifications

  • Clear the cursor reference immediately on removal. Make it volatile for the existing lock-free iterator, skip cleared references during iteration, and install the current cursor instance when reusing a node.
  • Compact removed nodes after max(64, live cursor count) removals, or immediately when the container becomes empty. Prune the pending-update list and unlink removed list nodes in linear time, repairing shared counters in one pass. Compaction preserves pending position updates and does not sort or flush them.
  • Track pending-list membership separately from pendingPosition, so cancellation and reactivation cannot add the same node repeatedly.
  • Replace the pending priority queue with a Fastutil ObjectArrayList and sort only before incremental flushes. Process descending old positions to reduce list traversal when cursors advance together. Small batches use the standard sort; larger batches skip sorting when already ordered and otherwise use in-place quicksort to avoid temporary merge buffers. Full rebuilds skip pending-list sorting. Fastutil is already a dependency.
  • Make add() cancel pending untracking when an existing cursor receives a real position, matching updateCursor(). This is a defensive consistency fix: production activation currently guards against adding an already-active cursor.
  • Fix a pre-existing, production-reachable tail-counter bug: a tail can leave or join a shared position group without changing list order. Skipping group-counter adjustments produces incorrect ranks even with ordinary position updates on master. Deduplicated reactivation exposes another trigger by avoiding a full rebuild that previously masked the bug. Reuse the existing counter when the tail remains in its own position group, avoiding an allocation without skipping shared-group adjustments.
  • Use a validated optimistic read for clean slowest-position lookups. Fall back to the write lock when pending changes need flushing or concurrent mutation invalidates the snapshot. The state-checking method also uses the write lock when flushing.
  • Add deterministic retention, compaction, reactivation, iterator, and slowest-position tests. Pin exclusive locking for pending flushes and appending after compaction removes the tail. Strengthen the test-only state checker to validate live cursor references, backward links, the tail, and the tracked node count. Extend ActiveManagedCursorContainerBenchmark with cursor-churn, read-only slowest-position, and tail-seek workloads. Add ActiveManagedCursorContainerCohortBenchmark for grouped cursor advances.

Verifying this change

  • Make sure that the change passes the CI checks.

Local validation passed:

./gradlew :managed-ledger:test \
  --tests ActiveManagedCursorContainerTest \
  --tests ActiveManagedCursorContainerRetentionTest \
  -PtestRetryCount=0 -PtestFailFast=false
./gradlew quickCheck

The two test classes now run 59 cases. The 11 retention regression cases fail on the unpatched implementation with only read-only testing accessors added. Coverage includes removal before and after position tracking, bounded retention without queries, preserving deferred updates and shared counters during compaction, repeated reactivation, replacement cursor identity, iterator behavior, and cached reads completing while another reader holds the lock. Two additional regression cases fail before the tail-counter fix and cover reactivation after a pending move, joining and leaving position groups, and subsequent compaction. Four further regression cases cover cancelling pending untracking through add() on both rebuild and incremental paths; all four fail before the defensive fix. Two further tests protect the pending-flush write lock and compaction tail repair. Both were mutation-tested: replacing the flush write lock with a read lock fails the first test; removing tail = previous from compaction fails the second. The additional structural checks run only in the test helper, not on production paths. A further 24 parameterized cases verify batches around the sorting threshold, ascending/reverse arrival order, mixed new and tracked nodes, coalesced repeated updates, shared positions, and a second flush to check pending-state reuse.

Performance measurements

The final hybrid implementation at 108223dba7b4 was measured against this PR's previous PriorityQueue implementation at c80d72a4db89, plus a fresh churn comparison against #26506. All 18 configurations completed and passed their benchmark state checks. Both sides use the same final benchmark jar and dependencies; only implementation classes are overridden for the comparison variants.

Environment: macOS arm64, Corretto 21.0.11, DEFAULT implementation, 256 MiB heap, GC profiler. Each configuration uses two forks, two 1-second warmup iterations and three 1-second measurement iterations per fork. Variants run sequentially, alternating order between parameter combinations. Tables show millions of updates/operations per second ± JMH-reported 99.9% error; percentages compare means. These are short local measurements, not production performance guarantees. Linux x86_64 validation is still needed.

Final hybrid versus the previous PriorityQueue implementation

Query ratio is updates per rank query; zero disables rank queries. randomSeekingForward01 uses one thread and randomSeekingForward10 uses ten. cohort is the new ActiveManagedCursorContainerCohortBenchmark.advanceCohort: it advances the fastest 32 cursors together in ascending old-position order, then queries rank. Both population sizes force incremental processing, and @OperationsPerInvocation(32) normalizes throughput and allocation per cursor update. This deliberately ordering-sensitive workload demonstrates the traversal benefit; it is not a measured broker traffic distribution.

Workload Cursors Query ratio PriorityQueue M ops/s Final hybrid M ops/s Mean change Allocation B/op (before → after)
randomSeekingForward01 100 1 28.80 ± 0.77 33.29 ± 1.07 +15.6% 47.83 → 47.83
randomSeekingForward01 100 10 21.61 ± 0.18 24.41 ± 1.84 +12.9% 47.14 → 47.14
randomSeekingForward01 100 100 15.20 ± 0.29 18.07 ± 0.31 +18.9% 65.35 → 65.35
randomSeekingForward01 500 100 10.65 ± 0.61 10.68 ± 0.52 +0.2% 46.37 → 46.37
randomSeekingForward10 500 10 6.16 ± 0.58 6.89 ± 0.96 +11.8% 48.37 → 48.33
cohort 100 32 29.12 ± 0.86 51.02 ± 3.14 +75.2% 47.50 → 47.50
cohort 500 32 29.11 ± 0.38 43.78 ± 4.41 +50.4% 47.50 → 47.50
cursorChurn 500 0 16.15 ± 0.40 16.80 ± 0.33 +4.0% 192.16 → 192.00

Cursor churn versus #26506

At 500 active cursors with no rank queries, cursorChurn measures 15.09 ± 0.56 M ops/s for #26506 versus 16.66 ± 0.48 M ops/s for this PR, a +10.4% change in mean throughput. #26506's implementation is pinned to bbf855c72c33 and overrides only implementation classes in the same benchmark jar.

Why this sorting strategy

FIFO experiments improved several random-seek workloads but lost roughly half their throughput when cursors advanced together. Sorting only during incremental flushes preserved descending-old-position processing, but plain ArrayList.sort introduced temporary-buffer allocations for larger batches. The final hybrid uses the standard small-array sort below 32 entries on JDK 21; for larger batches it skips sorting if already ordered, otherwise uses Fastutil's in-place quicksort. Full rebuilds skip pending-list sorting. Separate paired experiments against plain deferred sorting reduced allocation from 58.37 to 46.37 B/update for large random batches and from 80.00 to 47.50 B/update for the 500-cursor cohort. Fastutil is already a dependency.

The read-only slowest-position benchmarks remain API coverage; this implementation's slowest-position lookup is not on the current broker path.

Reproducing the current implementation's workloads
./gradlew :microbench:shadowJar

# Repeat with (cursors, query ratio): (100,1), (100,10), (100,100), (500,100).
java -jar microbench/build/libs/microbench-*-benchmarks.jar \
  'ActiveManagedCursorContainerBenchmark.randomSeekingForward01' \
  -p activeManagedCursorContainerImplType=DEFAULT \
  -p numberOfCursors=100 \
  -p getNumberOfCursorsAtSamePositionOrBeforeRatio=10 \
  -f 2 -wi 2 -i 3 -w 1s -r 1s -jvmArgs '-Xms256m -Xmx256m' -prof gc

# Ten-thread comparison: use randomSeekingForward10, 500 cursors, query ratio 10.
# Churn comparison: use cursorChurn, 500 cursors, query ratio 0.

java -jar microbench/build/libs/microbench-*-benchmarks.jar \
  'ActiveManagedCursorContainerCohortBenchmark.advanceCohort' \
  -p numberOfCursors=100,500 \
  -f 2 -wi 2 -i 3 -w 1s -r 1s -jvmArgs '-Xms256m -Xmx256m' -prof gc

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model — enforce exclusive deferred-state mutation and allow validated optimistic reads when no changes are pending.
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

…tion

Release removed cursor references immediately and compact stale nodes in
batches without flushing lazy position updates. Track queue membership to
prevent duplicate entries during reactivation.

Fix deferred-state mutation under a shared read lock. Use optimistic reads
for clean slowest-position queries and exclusive locking when flushing.

Add regression tests and cursor churn/read-only JMH workloads.

Assisted-by: OpenAI Codex
Remove the tail-only early return from incremental position updates so
joining and leaving a shared position group updates cursor ranks.

Add regression coverage for a moving cursor that is removed and
reactivated, tail group transitions, and subsequent node compaction.

Assisted-by: OpenAI Codex
Make add() clear pending removal state when an existing cursor receives
a real position, matching updateCursor(). Add regression coverage for
both rebuild and incremental flushing and both ways of untracking.

Assisted-by: OpenAI Codex
Keep the tail counter when its old and new positions are separate from the preceding group. Preserve counter updates when joining or leaving a shared position group, and add a focused tail-seek benchmark.

Assisted-by: Codex
@lhotari
lhotari marked this pull request as draft September 9, 2026 11:10
Add regression tests for exclusive pending-position flushes and appending after compaction removes the tail. Strengthen the test-only state checker to validate cursor references, backward links, the tail, and tracked node count without changing production operation paths.

Assisted-by: Codex
@lhotari
lhotari marked this pull request as ready for review September 9, 2026 11:17
@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 9, 2026
@lhotari
lhotari marked this pull request as draft September 9, 2026 12:12
Accumulate pending nodes in an ObjectArrayList and preserve descending old-position processing during incremental flushes. Use the small-array sort for short batches, skip sorting ordered larger batches, and use in-place quicksort otherwise to avoid temporary merge buffers. Full rebuilds skip pending-list sorting.

Cover batch boundaries, repeated updates, new cursors and shared position groups. Add a grouped-advance benchmark. All 59 cursor cases and quickCheck pass; final paired JMH comparisons completed without a measured throughput regression.

Assisted-by: Codex
@lhotari
lhotari marked this pull request as ready for review September 9, 2026 13:26
@lhotari
lhotari merged commit 72632e5 into apache:master Sep 9, 2026
44 checks passed
lhotari added a commit that referenced this pull request Sep 9, 2026
dao-jun pushed a commit to ascentstream/pulsar that referenced this pull request Sep 20, 2026
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants