Skip to content

[fix][broker] Fix AvgShedder assignment cache keying with stable bundle names - #26246

Merged
lhotari merged 12 commits into
apache:masterfrom
void-ptr974:fix/avg-shedder-stable-bundle-key
Sep 11, 2026
Merged

lhotari merged 12 commits into
apache:masterfrom
void-ptr974:fix/avg-shedder-stable-bundle-key

Conversation

@void-ptr974

@void-ptr974 void-ptr974 commented Jul 26, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

AvgShedder must preserve each bundle's planned destination between load-shedding planning and broker assignment. It previously keyed that state by BundleData, whose load statistics participate in equals and hashCode and are updated in place. A refreshed key can become unreachable, while different bundles with equal statistics can overwrite each other, causing assignment to ignore the planned destination.

Modifications

  • Key pending destinations by canonical bundle name in a ConcurrentHashMap.
  • Scope the pending plan to one shedding attempt and clear it after processing, including planning failures.
  • Pass the stable bundle name through a backward-compatible selectBrokerForBundle(...) strategy method.
  • Use the planned destination while it is valid. If it is unavailable, use normal placement; if it is overloaded, select through the existing strategy from healthy candidates and retain the plan when none are available.
  • Keep the four-argument selector as the normal uncached fallback and return Optional.empty() for empty candidates.

The existing directed-unload affinity entry carries the selected destination to the next ownership lookup after the temporary AvgShedder plan is cleared.

Compatibility

New strategy callbacks have default implementations, and the historical selector remains available. Existing strategy implementations continue to work without changes.

The change is limited to the Modular load manager. It does not change configuration defaults, metadata or wire formats, protocols, REST APIs, or the Simple and Extensible load managers.

Verification

  • ./gradlew --no-daemon quickCheck
  • AvgShedderTest: 3 tests passed
  • ModularLoadManagerStrategyTest: 9 tests passed
  • ModularLoadManagerImplTest: 17 tests passed
  • Focused broker tests ran with -PtestRetryCount=0

Coverage includes mutable and equal BundleData, stable-name propagation, attempt cleanup on success and failure, unavailable and overloaded planned destinations, bundle-aware overload retry, affinity handoff, and legacy strategy fallback.

GitHub CI has not yet been triggered for this PR.

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

  • Dependencies
  • Public API — adds backward-compatible default methods
  • Schema or metadata formats
  • Configuration defaults
  • Threading model
  • Binary protocol
  • REST endpoints
  • Admin CLI
  • Metrics
  • Deployment

@void-ptr974
void-ptr974 marked this pull request as draft July 26, 2026 07:55
@void-ptr974
void-ptr974 force-pushed the fix/avg-shedder-stable-bundle-key branch from 3c2419a to acd9d15 Compare July 26, 2026 09:08
@void-ptr974 void-ptr974 changed the title [fix][broker] Use stable bundle names in AvgShedder [fix][broker] Fix AvgShedder assignment cache keying with stable bundle names Jul 26, 2026
@void-ptr974
void-ptr974 marked this pull request as ready for review July 26, 2026 09:58

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for digging into this — the diagnosis is correct and I was able to confirm it independently: BundleData is annotated @EqualsAndHashCode (pulsar-common/.../BundleData.java:26) over fields that BundleData.update(NamespaceBundleStats) mutates in place on every load-report refresh, so a HashMap<BundleData, String> key silently becomes unreachable, and two bundles with identical stats collide on one key. Re-keying by bundle name is the right fix.

I also verified the key spaces line up: ModularLoadManagerImpl:902 uses serviceUnit.toString(), which is the same key space as loadData.getBundleData() and getBundleDataForLoadShedding(), so the write in selectBundleForUnloading and the read in the placement path agree. The new selectBrokerForBundle default method is source- and binary-compatible, and LeastLongTermMessageRate, LeastResourceUsageWithWeight and RoundRobinBrokerSelector are correctly unaffected.

Two things I'd like to discuss before this goes in — mainly (1).

1. Planned destinations are now permanent — nothing ever removes an entry from bundleBrokerMap

AvgShedder.java:290-310

selectBrokerWithBundleName returns the cached broker and never removes it. The only writes are the put in selectBundleForUnloading (:196) and the fallback put (:305); there is no removal anywhere, and onActiveBrokersChange (:206) just delegates to the no-op default.

Before this PR the bug provided accidental expiry: the key's hash drifted on the next load-report cycle, the entry became unreachable, and the next assignment re-randomized. After the fix the mapping is sticky for the lifetime of the process, which I think is a bigger behaviour change than "retain a planned destination across load-data refreshes" suggests.

Two consequences:

Manual unload stops moving bundles. A bundle is shed to broker2; the entry is written and consumed correctly (this is the intended behaviour, and it now works). An operator later runs pulsar-admin namespaces unload to move it off broker2. The next lookup hits selectBrokerForBundle, broker2 is still a candidate, so the bundle is assigned straight back to broker2. From the operator's point of view the unload is a no-op, until AvgShedder itself happens to shed the bundle again.

The overload-escape path becomes dead code for AvgShedder. ModularLoadManagerImpl:981-987 re-runs placement over the widened candidate set when the selected broker is above loadBalancerBrokerOverloadedThresholdPercentage. With a sticky entry, the second call returns the identical broker.

I don't have a strong opinion on which way to resolve it — options I can see are dropping the entry once an assignment has consumed it, bounding validity via loadData.getRecentlyUnloadedBundles(), or clearing consumed entries at the start of findBundlesForUnloading. Whichever is chosen needs to keep the entry valid across the retry at ModularLoadManagerImpl:981 within a single assignment.

2. bundleBrokerMap grows without bound, while loadData.getBundleData() is pruned

AvgShedder.java:52

ModularLoadManagerImpl actively removes bundle entries when a bundle goes inactive (:600-604) and when it is split (:796). bundleBrokerMap has no equivalent, so on a long-lived leader with bundle splits, namespace deletion or topic churn, every bundle name ever assigned is retained forever.

This is pre-existing rather than a regression (and it actually grew faster before, since every hash drift orphaned a fresh entry), so I'd be fine with it as a follow-up. But since the PR is already rewriting this field, pruning against loadData.getBundleData().keySet() is only a few lines.

3. bundleBrokerMap is a plain HashMap mutated from two threads under different monitors

AvgShedder.java:52, :196, :305

doLoadShedding() is synchronized on the ModularLoadManagerImpl instance (:637) and reaches bundleBrokerMap.put through findBundlesForUnloading. The assignment path synchronizes on a different monitor — synchronized (brokerCandidateCache) (:902) — and also does get/put. Nothing orders those two, so concurrent put during a resize can corrupt the table or spin.

Again pre-existing and not introduced here, but LeastResourceUsageWithWeight.selectBroker is synchronized for exactly this reason, and switching the declaration to ConcurrentHashMap is a one-word change while that line is already being touched.

4. The identity-scan compatibility path is fragile, and its only callers are this PR's tests

AvgShedder.java:312-323

findBundleNameByIdentity recovers the key by scanning every entry of loadData.getBundleData() for reference equality (entry.getValue() == bundleToAssign). That is O(number of bundles in the cluster) per call, and it only works when the caller passes the exact instance stored in loadData — a contract that appears nowhere in the ModularLoadManagerStrategy javadoc. A caller that passes a defensive copy silently falls into the uncached random branch with no signal.

Since ModularLoadManagerImpl now always uses the name-aware path, the scan exists only for third-party callers and for the tests added here. Leaving the 4-arg selectBroker as a plain uncached fallback (optionally @Deprecated) would be simpler and equally correct; the assertions in AvgShedderTest and testSheddingMultiplePairs would move to selectBrokerForBundle.

5. Undocumented behaviour change: empty candidates no longer throw

AvgShedder.java:278, :292-294

Previously empty candidates reached getExpectedBroker, hit % 0, and the catch (Throwable) fallback threw ArithmeticException again (:335, :343), which propagated out of selectBrokerForAssignment. Both paths now short-circuit to Optional.empty(), which ModularLoadManagerImpl:970 handles as "No brokers available".

That's an improvement, but it's a real behaviour change that isn't called out beyond "keep a valid uncached fallback" — worth an explicit line under Modifications.

6. Tests

The fixture cleanup is genuinely nice: dropping the setNumSamples(i) hacks and adding assertEquals(loadData.getBundleData().get("bundle1-0"), loadData.getBundleData().get("bundle3-0")) in testSheddingMultiplePairs turns the equal-but-distinct BundleData collision into an explicit precondition instead of something the old test worked around. I checked that those are two distinct instances, so the identity scan resolves them unambiguously.

A few gaps:

  • Nothing tests the actual integration point of the fix — that ModularLoadManagerImpl passes the correct bundle name through. ModularLoadManagerImpl.selectBroker(ServiceUnitId) is already @VisibleForTesting, so a spy strategy asserting the received name would lock the wiring in.
  • Nothing covers the sticky-forever behaviour from (1), which is the riskiest part of the change.
  • assertNotEquals(plannedBundleData.hashCode(), originalHashCode) couples the test to Lombok's generated hash including topics. Comparing the objects (or a pre-mutation copy) expresses the same intent without depending on hash internals.
  • The new BundleData() → new BundleData(1, 1) and TimeAverageMessageData() → TimeAverageMessageData(1) fixture change in testHitHighThreshold is necessary — with maxSamples == 0, update() can't record the sample — but it's unexplained. A one-line comment would save the next reader the detour.

Nit

AvgShedder.java:291: the BundleData bundleToAssign continuation line is indented 2 columns past the opening paren. Purely cosmetic, checkstyle won't flag it.


Intent and implementation match; the only place the description understates things is (1), where "retain a planned destination across load-data refreshes" is doing some heavy lifting for "retain it indefinitely".

I reviewed by reading only — I did not run the build or tests, so I'm relying on your local run plus CI for that. I did confirm that every import in the rewritten ModularLoadManagerStrategyTest is still used (Field and Map are still needed by the untouched LeastResourceUsageWithWeight tests), so there shouldn't be an unused-import checkstyle failure.

Assisted-by: Claude Opus 5 (Claude Code), with a second independent pass from Codex gpt-5.6-sol (codex review) that reported no actionable findings. All findings above come from the Claude pass; the code references were checked against the PR head.

@void-ptr974
void-ptr974 marked this pull request as draft July 27, 2026 02:16
@void-ptr974

Copy link
Copy Markdown
Contributor Author

Thanks for the detailed review — addressed in follow-up commits.

  • Changed AvgShedder to store pending destinations as bundle name -> broker.
  • Added selectBrokerForBundle(...) and wired ModularLoadManagerImpl to pass the bundle name to the placement strategy.
  • Added onUnloadAttemptCompleted(...); ModularLoadManagerImpl calls it in finally after processing a shedding plan.
  • Replaced the shared map with ConcurrentHashMap. The existing 4-argument selector remains as an uncached compatibility fallback.
  • Added tests for mutable/equal BundleData, unavailable planned destinations, cleanup, and the three-broker unload flow.

@void-ptr974
void-ptr974 marked this pull request as ready for review July 27, 2026 13:13
Preserve the original pending destination across temporary fallback selection, exclude an overloaded broker from placement retry, deprecate the nameless AvgShedder selector, and replace the live load-data E2E with deterministic coverage.

Assisted-by: OpenAI Codex
Clear pending destinations after every shedding attempt, including planning failures. Preserve legacy overload retry and subclass selection behavior.

Assisted-by: Codex
Remove misleading initialization and shutdown logs from the normal no-plan fallback, and describe the attempt lifecycle directly.

Assisted-by: Codex

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the follow-up work — the three commits since my last pass (d1329a0f, bc30c41c, db349c57) close every point I had open, and I re-checked the core of the fix rather than just the replies.

What I verified at db349c57:

  • The new key is stable for exactly as long as it needs to be. pendingBundleToBroker is keyed by the canonical bundle name that ModularLoadManagerImpl.selectBrokerForAssignment already derives from serviceUnit.toString(), and its lifetime is now one shedding attempt: cleared at AvgShedder.java:66 when planning starts, and cleared again from the finally at ModularLoadManagerImpl.java:706-708. So it does not need to survive a restart, a bundle split or a range change — a name that no longer exists simply misses and falls through to normal placement. That is the right scope for this map.
  • The unbounded-growth risk is gone, and it was real before. On master the only operations on bundleBrokerMap are put (AvgShedder.java:196, :286) and get — there is no removal anywhere, and because the BundleData keys mutate in place the stale entries also become unreachable. The attempt-scoped clear() replaces a genuine leak.
  • There is a test that fails without the change. AvgShedderTest.testSheddingMultiplePairs now drops the timeAverageMessageData.setNumSamples(i) lines whose comment on master reads "as AvgShedder map BundleData to broker, the hashCode of different BundleData should be different so we need to set some different fields", and asserts instead that bundle1-0 and bundle3-0 are equals but notSame (AvgShedderTest.java:272-273) while still resolving to broker2 and broker4 respectively. Under the old Map<BundleData, String> those two puts collide and one destination overwrites the other. That test papering-over is good evidence that the bug was already known.
  • Cache access is thread-safe as written. Writers run under the synchronized doLoadShedding; readers reach it from selectBrokerForAssignment, which holds a different monitor (brokerCandidateCache), so reads genuinely race with planning. Since the fallback write-back was removed in 7c31e428, all that is left is single-key get/put/clear on a ConcurrentHashMap, so the worst case is a reader seeing a half-built plan and falling back to normal placement. That is acceptable. The sibling brokerScoreMap / brokerHitCountFor* maps are still plain HashMaps, but they are touched only from the shedding thread, unchanged from master.
  • The lifecycle answer on the affinity map checks out. The directed unload goes doLoadShedding -> admin unloadNamespaceBundle(ns, range, dest) -> NamespacesBase.setNamespaceBundleAffinityAsync -> ModularLoadManagerImpl.setNamespaceBundleAffinity (NamespacesBase.java:1507), and ModularLoadManagerWrapper.getLeastLoaded:67 consumes it one-shot. So clearing the plan when the attempt ends really does not strand the destination.

What I think is still open: I agree with @Denovo1998's newest point about the overload-retry path, and I left one comment there with the AvgShedder-specific consequence. There is also a small @Deprecated signalling inconsistency worth a one-line fix.

One process note: CI has never run on this PR — db349c57 and every earlier commit report zero check runs, and the PR carries no labels, so it is sitting behind the ready-to-test gate rather than failing. Someone with write access needs to trigger it before this can be judged green; please do not read the empty status as a passing run.

@void-ptr974

Copy link
Copy Markdown
Contributor Author

Thanks again for the thoughtful review. I’ve updated the implementation to address these concerns while keeping the existing strategy contract compatible. Please take another look when you have a chance.

@Denovo1998 Denovo1998 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.

LGTM!

@lhotari lhotari added the triage/lhotari/important lhotari's triaging label for important issues or PRs label Sep 10, 2026

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

The follow-up changes address both remaining comments: overload retry now stays bundle-aware, healthy alternatives exclude the shedding source, and the four-argument fallback is no longer deprecated. I found no further issues in the full change at 1191348c234. The directed-unload affinity carries the destination independently of the temporary plan.

This assessment is based on code and test inspection; I did not run the tests locally. There are currently no CI results on the reviewed head, so CI validation remains outstanding.

@lhotari lhotari added this to the 5.0.0-M2 milestone Sep 11, 2026
@lhotari
lhotari merged commit 04f51b6 into apache:master Sep 11, 2026
43 checks passed
lhotari pushed a commit that referenced this pull request Sep 12, 2026
lhotari pushed a commit that referenced this pull request Sep 12, 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