Conversation
refreshMutPipelines() rebuilds the per-store/per-range mutation pipeline map and assigns it, but only closed pipelines of stores that vanished entirely. A pipeline whose range is no longer led by that store (the regular outcome of a RangeLeaderBalancer leadership transfer or a merge) was simply dropped: its RX subscription, its dedicated gRPC bidi stream and the server-side MutatePipeline leaked on every topology event - unreachable from close() - giving monotonically growing resource usage on long-running brokers. Replace the store-level-only clear loop with a general pass that closes every pipeline instance absent from (or replaced in) the new map, extracted as a package-private static helper for direct unit testing. Tests: range-leadership-move / store-disappear / instance-replaced cases (new BaseKVStoreClientPipelineRefreshTest). Fixes apache#302
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
BaseKVStoreClient.refreshMutPipelines()rebuilds the per-store/per-range mutation pipeline map on every route refresh, reusing pipelines whose range is still led by that store and dropping the rest. The only cleanup loop closes pipelines of stores that disappeared entirely; a pipeline whose range is no longer led by that store — the routine outcome of aRangeLeaderBalancerleadership transfer or a merge — is silently dropped from the map withoutclose().Impact
The dropped
ManagedMutationPipelinekeeps its RX subscription and its dedicated gRPC bidi stream (plus a server-sideMutatePipeline), and it is unreachable fromBaseKVStoreClient.close()once dropped.RangeLeaderBalanceris registered by default in dist-worker / inbox-store deployments, so every leadership transfer leaks one pipeline + one gRPC bidi stream + one server-side pipeline: monotonically growing resource usage on long-running brokers.Fix
Close every pipeline instance absent from (or replaced in) the new map — one pass covers all three forms (range moved away, store disappeared, instance replaced):
The helper is a package-private static for direct unit testing; the previous store-level-only clear loop is subsumed (a vanished store's ranges are all absent from the new map).
Test evidence
New
BaseKVStoreClientPipelineRefreshTestwith three cases: range-leadership-move / store-disappear / instance-replaced (mocks verify the dropped pipelines are closed and the kept ones are not).Production validation
Deployed in a rolling 6-node cluster (wholesale
lib/replacement, one node at a time, zero rollbacks); cluster mesh, all 13 live client connections and cross-node delivery verified post-deploy. Per-node baselines captured for leak tracking (JVM threads 216–229, RSS 2.8–5.1 GB, gRPC inter-node connections 408–417) — the leak would show as monotonic growth across further leadership transfers.Fixes #302