Repository navigation
[fix][broker] Fix ownership-generation races in OwnershipCache removeOwnership and lock-expiry cleanup - #26197
Conversation
1f7e054 to
8f4d0b6
Compare
|
I ran an AI-assisted review of this PR (Claude as the local reviewer plus OpenAI Codex 1. In 2. The "race-free" The comment on 3. (optional) In the common uncontended case ( 4. (optional) Straggler-topic semantics change A topic installed in 5. (optional) Barrier latency notes, for awareness only
Also: nice catch on the |
…Ownership and lock-expiry cleanup
8f4d0b6 to
724a40a
Compare
|
@lhotari
|
|
please merge origin/master (don't rebase) to the PR branch and resolve conflicts |
# Conflicts: # pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java # pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java
merged @lhotari |
|
Hello, we recently went through a production issue where two brokers seemed to hold ownership of the same topic. |
@alexandrebrg Which broker version did it happen with? |
|
The broker version was Thank you |
@alexandrebrg In that case, you might also be affected by the issues fixed in #25910 and #25913. There's also #26000, which fixes an issue related to leader broker loss that could make other issues more likely, such as the bug that the current PR under discussion will fix. |
|
Good to know, we'll try to proceed to an upgrade then |
lhotari
left a comment
There was a problem hiding this comment.
I read this at b08c6f94 in a worktree and exercised it rather than only reading it: OwnershipCacheTest passes 13/13 (7.7s) at the head commit, and I mutation-tested the production change to find out which parts the new tests actually pin.
The invariant is clear and, as far as I can tell, right. An OwnedBundle instance is now one ownership generation, because it carries the ResourceLock that created it; removeOwnership(OwnedBundle) only takes effect if locallyAcquiredLocks still maps the bundle to that lock; serialize() orders acquire against both release overloads per bundle, so an acquire can no longer be reported successful for a generation that is mid-release; and invalidateLocalOwnerCache(bundle, expected) and cleanUnloadedTopicFromCache(bundle, snapshot) extend the same identity rule to the owned-bundle cache and the broker topic cache. The generation-blind holes in the motivation are real, and this addresses them at the right level rather than papering over the symptom.
Four of the five parts of the change I mutated are each killed by a specific test, which is a good result for a change of this kind:
| mutation of the production change | killed by |
|---|---|
removeOwnership(OwnedBundle) → generation-blind locallyAcquiredLocks.remove(bundle) |
testStaleUnloadDoesNotReleaseReacquiredOwnership |
serialize(...) → return operation.get() (no barrier at all) |
testRemoveOwnershipWithAcquisitionInFlight, testTryAcquiringOwnershipWaitsForInFlightOwnedBundleRelease |
lock-expiry listener → drop locallyAcquiredLocks.remove(namespaceBundle, rl) |
testExpiredLockIsRemovedFromLocallyAcquiredLocks |
| cache loader → never throw when expiry won the publication race | testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership |
ownedBundle == expectedOwnedBundle in invalidateLocalOwnerCache → if (ex == null) |
nothing — still 13/13 green |
Three comments below. The first is the one place the generation rule was not applied, and it happens to be the most destructive of the three paths; the second is a defect in the new barrier itself; the third is that last table row.
Two of them turn on the same interleaving, which is worth stating once: the lock-expiry listener is registered with a plain thenRun, so when expiredFuture completes concurrently with registration the listener body can run on the completing thread, racing the loader rather than strictly before or after it. Under that ordering the loader fails the publication check, Caffeine drops the failed entry, a lookup installs a newer generation, and the still-running listener then operates on it. That is the same window testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership already builds a controllable ResourceLock for, so it should be cheap to drive deterministically.
Two things that are not defects but change how a reader should scope the fix:
- The ownership-generation half only affects the legacy load balancer.
ExtensibleLoadManagerImpl.tryAcquiringOwnershipdoes not go throughOwnershipCacheat all — it routes toassign(...)onServiceUnitStateChannel(ExtensibleLoadManagerImpl.java:578-592) — andNamespaceService.unloadNamespaceBundle/removeOwnedServiceUnitAsyncbranch away fromOwnershipCachewhen the extension is enabled. What theServiceUnitStateChannelImplchange contributes here is the topic-cache snapshot, which does apply to both paths. Saying which of the three motivating defects is fixed on which path would save the next reviewer the trip. cleanUnloadedTopicFromCache(NamespaceBundle)was removed rather than deprecated. Every in-repo caller is updated, so this is only about out-of-tree callers of a broker-internal class — keeping a deprecated one-arg overload that captures the snapshot itself would cost nothing.
A few things I checked that turned out fine, in case it saves re-arguing them: closing a pre-taken snapshot instead of iterating the live topics map does not meaningfully widen the straggler window, because the old topics.forEach ran after updateBundleState, which completes synchronously (ownedBundlesCache uses a direct executor), so the two observation points are effectively the same instant; dropping updateBundleState(bundle, false) from handleUnloadRequest loses nothing, since the CAS above it already flipped isActive and the removed call resolved through the cache and could deactivate a newer generation; locallyAcquiredLocks.remove(bundle, lock) really is identity-based, because ResourceLockImpl overrides hashCode() but not equals(); and the release path does not deadlock against the new barrier, because ResourceLockImpl.release() completes expiredFuture before its own result future, so a re-entrant serialize() call only queues behind the in-flight barrier. The ex.getMessage() → e.getMessage() correction at NamespaceService.java:1102 is a real fix too — the old message reported the wrong exception.
One meta point, since it affects how much the green run is worth: a clean local run is weak evidence for a change like this, so I have based the assessment on the invariant and on what survives mutation instead.
b9ae6cc to
0adab8c
Compare
…the ownership barrier without recursion
0adab8c to
2166893
Compare
lhotari
left a comment
There was a problem hiding this comment.
The follow-up in 2166893 addresses my three earlier comments: the delayed expiry listener checks its ownership generation, the normal queue drain hands off through the executor, and the new delayed-listener test exercises both generation guards. I have one nonblocking test-maintenance comment below.
…sting setter instead of reflection
…Ownership and lock-expiry cleanup (#26197) Co-authored-by: maxlisongsong <[email protected]> Co-authored-by: Lari Hotari <[email protected]> (cherry picked from commit 37d64fe)
…Ownership and lock-expiry cleanup (#26197) Co-authored-by: maxlisongsong <[email protected]> Co-authored-by: Lari Hotari <[email protected]> (cherry picked from commit 37d64fe)
…Ownership and lock-expiry cleanup (apache#26197) Co-authored-by: maxlisongsong <[email protected]> Co-authored-by: Lari Hotari <[email protected]>
Motivation
OwnershipCache.removeOwnership(bundle)removes whicheverResourceLockhappens to be inlocallyAcquiredLocksat that moment, without checking that it belongs to the ownership"generation" the caller was working on. Three defects follow from this:
A stale unload chain can destroy a re-acquired ownership. When a bundle's resource lock
expires, the lock-expiry callback in the cache loader fires
unloadNamespaceBundle(...)asynchronously (dropping its future) and invalidates the local owner cache immediately. A
concurrent lookup can then re-acquire the bundle (new lock, new
OwnedBundle) while the staleunload chain is still running. When that chain reaches its final
removeOwnership(bundle), itremoves and releases the newly acquired lock, deleting the new owner's znode while the local
cache still claims active ownership. This opens a transient dual-serving window (fenced-ledger
errors on persistent topics, silent message loss on non-persistent topics) and causes ownership
flapping. The same stale chain also deactivates the newer
OwnedBundlethrough the cache lookupin
updateBundleState(bundle, false), and itscleanUnloadedTopicFromCache(bundle)step canevict the newer generation's topics from the broker topic cache.
removeOwnershipreports success while an acquisition is in flight, leaving zombieownership. The lock is only put into
locallyAcquiredLocksafteracquireLockcompletes, soa
removeOwnershipcall racing with an in-flight acquisition finds the map empty, returns acompleted future ("released"), and the lock installed afterwards is never cleaned up: the broker
silently keeps the lock and the cache entry for a bundle the caller believes released (e.g. a
deleted bundle/namespace), polluting load reports until restart.
The lock-expiry callback never removes the expired lock from
locallyAcquiredLocks(asymmetric add/remove bookkeeping), so a dead lock handle lingers in the map and drifts from
the owned-bundle cache. The callback also drops the future returned by
unloadNamespaceBundle(...), so failures of the triggered unload are silent.Which load manager each part applies to
ExtensibleLoadManagerImpl.tryAcquiringOwnershipdoes not go throughOwnershipCache(it routes toServiceUnitStateChannel.assign), andNamespaceService.unloadNamespaceBundle/removeOwnedServiceUnitAsyncbranch away fromOwnershipCachewhen the extension is enabled. So:OwnershipCache/OwnedBundle): the ownership-generationbinding,
removeOwnership(OwnedBundle), the per-bundle acquire/release barrier, and all of thelock-expiry listener changes. Defects 1 (lock/cache half), 2 and 3 above are legacy-only.
BrokerService): the shared topic-future snapshot and the identity-safecleanUnloadedTopicFromCache(bundle, snapshot), used by bothOwnedBundle.handleUnloadRequestand
ServiceUnitStateChannelImpl.closeServiceUnit. This is the topic-cache half of defect 1.Modifications
Ownership generation (
OwnershipCache,OwnedBundle)OwnedBundleto theResourceLockacquired for it, so anOwnedBundleinstancerepresents a single ownership generation. Add
removeOwnership(OwnedBundle)which releases onlythat generation via
ConcurrentHashMap.remove(key, value); if the bundle has since beenre-acquired, the newer ownership is left untouched.
OwnedBundle.handleUnloadRequestnow uses itand no longer calls
updateBundleState(its ownisActiveflag is already flipped via CASearlier in the method; the cache lookup could deactivate a newer generation).
serialize(...)) that orderstryAcquiringOwnershipagainst bothremoveOwnership(NamespaceBundle)andremoveOwnership(OwnedBundle)for the same bundle, so arelease can neither run in the gap between an acquire's cache publication and its lock install,
nor report ownership fully released while a fresh acquisition lands moments later. The barrier is
handed to the next queued operation on
PulsarService.getExecutor()rather than inline, so thedrain depth does not depend on the queue length (a nested inline drain overflowed the stack at
roughly 1,800 queued operations and left the bundle permanently unacquirable on that broker). If
the executor rejects the hand-off (broker shutting down), the barrier is released inline instead so
the queued releases still run.
locallyAcquiredLockswith a two-arg remove, triggersunloadNamespaceBundleonly if the cachestill holds the listener's own
OwnedBundle(a late listener must not close a re-acquiredgeneration's topics or release its lock), invalidates the owned-bundle cache entry only if it
still holds this generation, delivers
onNamespaceBundleUnloadto the ownership listeners onlywhen that invalidation actually happened (a generation that was never published produced no owned
event, and a newer generation in the cache is still owned; listeners such as the topic-policies
service keep per-namespace counts on these events, so an unmatched unload event is not harmless),
and logs failures of the triggered unload instead of dropping the future.
OwnedBundle, the load fails (IllegalStateException) instead of publishing an already-deadownership; Caffeine drops a failed load, so no zombie cache entry remains.
removeOwnership(OwnedBundle)releases nothing for anOwnedBundlethat is not bound to a lock(not created by the cache) instead of falling back to the generation-blind release; the two public
lock-less
OwnedBundleconstructors, which have no callers, are deprecated.OwnedBundle.handleUnloadRequestguards the topic snapshot and the start of the unload againstsynchronous exceptions and feeds them into the existing non-fatal path, so a failure there (for
example
RejectedExecutionExceptionduring shutdown) still releases the ownership instead ofleaving the bundle inactive with its lock registered forever. The cleanup step keeps using the
snapshot even when
unloadServiceUnitthrew, since the topic closes may already have been startedfrom it.
Topic cache (
BrokerService,ServiceUnitStateChannelImpl)getTopicFuturesInBundle(bundle)captures one snapshot of the bundle's topic futures at unloadstart; the new
unloadServiceUnit(..., snapshot)overload closes exactly that set, andcleanUnloadedTopicFromCache(bundle, snapshot)removes a topic only if its cached future isidentical to the captured one. Both unload call sites thread the same snapshot through both
steps. A topic that lands in the cache after the snapshot is deliberately left alone (see the
comments at both call sites). The one-arg
cleanUnloadedTopicFromCache(NamespaceBundle)is keptas a deprecated overload that captures the snapshot itself.
Split path (
NamespaceService)splitAndOwnBundleOnceAndRetrycomposes on the future returned byremoveOwnership(bundle)instead of dropping it (a release failure is logged, not turned into a split failure), and the
adjacent
exceptionallyhandler now reports its own exception instead of the enclosing, alwaysnull one. The wait is bounded by
metadataStoreOperationTimeoutSeconds: the release is queuedbehind any in-flight acquire of the same bundle, which has no timeout of its own, so on timeout
the split completes and the release keeps running in the background.
Known follow-ups not addressed here
Raised during review and left out to keep this change scoped; neither changes the invariant above.
tryAcquiringOwnershipcalls for the same bundle used to share one Caffeine load; they now runsequentially and each issues its own metadata write. Bounded by the number of distinct lookup-key
variants in flight (a handful).
isDone()check, the
OwnedBundleis published and then dropped as soon as the listener's deferredinvalidation runs on publication. In that window a lookup can be answered with this broker for an
ownership that is already gone (clients re-lookup and converge), and since the deferred
invalidation and the acquire's
onNamespaceBundleOwnedare both dependents of the same cachefuture,
onNamespaceBundleUnloadcan be delivered beforeonNamespaceBundleOwnedfor that onegeneration.
Verifying this change
This change added tests and can be verified as follows:
OwnershipCacheTest.testStaleUnloadDoesNotReleaseReacquiredOwnership: simulates theexpiry-callback interleaving (cache invalidated, bundle re-acquired, then the stale unload chain
completes) and asserts the re-acquired lock, znode, and active state survive.
OwnershipCacheTest.testTryAcquiringOwnershipWaitsForInFlightOwnedBundleRelease: races anacquire against an in-flight
removeOwnership(OwnedBundle)and asserts the acquire is queuedbehind it and returns a fresh generation.
OwnershipCacheTest.testRemoveOwnershipWithAcquisitionInFlight: gatesacquireLockbehind acontrollable future to call
removeOwnershipdeterministically inside the in-flight window, thenasserts no lock/cache entry/znode remains once both settle.
OwnershipCacheTest.testExpiredLockIsRemovedFromLocallyAcquiredLocks: releases the lock withoutgoing through
removeOwnershipand assertslocallyAcquiredLocksis cleaned up together with thecache and exactly one unload event is delivered.
OwnershipCacheTest.testExpiryBeforePublicationDoesNotLeaveActiveZombieOwnership: analready-expired lock handed to an in-flight load makes the acquisition fail and leaves no trace,
and neither an owned nor an unload event is delivered for it.
OwnershipCacheTest.testDelayedExpiryListenerOfOldGenerationLeavesReacquiredOwnershipUntouched:completes the old generation's expiry future only after a newer generation was acquired, with the
listener's unload routed through the real by-name lookup; asserts the newer generation stays
cached, active, its lock is never released, and no unload event is delivered for it. Pins both
the listener's unload guard and the identity check in
invalidateLocalOwnerCache.OwnershipCacheTest.testSerializedOperationsDrainDeepBacklogWithoutWedgingBundle: queues 5,000acquires behind a gated first one and asserts all complete and the barrier map is empty.
OwnershipCacheTest.testUnloadRequestStillReleasesOwnershipWhenTopicUnloadThrowsSynchronously:unloadServiceUnitthrows synchronously;handleUnloadRequestmust return a future, stillrelease the ownership, and still run the cleanup with the original snapshot.
OwnershipCacheTest.testSerializedOperationsStillDrainWhenExecutorRejectsTasks: with a shut-downexecutor, a release queued behind an acquire still completes and the barrier map ends up empty.
OwnershipCacheTest.testRemoveOwnershipOfDetachedOwnedBundleDoesNotReleaseCurrentOwnership:releasing through an
OwnedBundlenot created by the cache leaves the current lock and znode alone.NamespaceServiceTest.testSplitCompletesWhenReleasingTheOldBundleStalls: with the old bundle'srelease never completing, the split still completes once the metadata operation timeout elapses.
BrokerServiceTest.testCleanUnloadedTopicFromCacheIsGenerationSafe,...IsGenerationSafeForBookkeeping,...RemovesMatchingSnapshot: against a real embeddedbroker, a stale cleanup with an old snapshot leaves a reloaded topic and its bookkeeping intact,
while a matching snapshot is removed.
OwnershipCacheTest(18),TopicPoliciesTest,MetadataStoreTopicPoliciesTest,BrokerServiceTest,NamespaceUnloadingTest,NamespaceOwnershipListenerTestand the split /unload cases of
NamespaceServiceTestpass locally.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
Threading model:
OwnershipCacheoperations for the same bundle are now serialized per bundle, andeach queued operation starts on
PulsarService.getExecutor()once the previous one settles.