Skip to content

fix(kv-client): close mutation pipelines dropped by a route refresh — fixes #302 - #309

Open
ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/kv-client-pipeline-leak
Open

ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/kv-client-pipeline-leak

Conversation

@ImDanXie

Copy link
Copy Markdown
Contributor

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 a RangeLeaderBalancer leadership transfer or a merge — is silently dropped from the map without close().

Impact

The dropped ManagedMutationPipeline keeps its RX subscription and its dedicated gRPC bidi stream (plus a server-side MutatePipeline), and it is unreachable from BaseKVStoreClient.close() once dropped. RangeLeaderBalancer is 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):

mutPplns = nextMutPplns;
closeDroppedMutationPipelines(currentMutPplns, nextMutPplns);
static void closeDroppedMutationPipelines(Map<String, Map<KVRangeId, IMutationPipeline>> current,
                                          Map<String, Map<KVRangeId, IMutationPipeline>> next) {
    Map<String, Map<KVRangeId, IMutationPipeline>> safeNext = next == null ? emptyMap() : next;
    for (Map.Entry<String, Map<KVRangeId, IMutationPipeline>> byStore : current.entrySet()) {
        Map<KVRangeId, IMutationPipeline> nextRanges = safeNext.getOrDefault(byStore.getKey(), emptyMap());
        for (Map.Entry<KVRangeId, IMutationPipeline> byRange : byStore.getValue().entrySet()) {
            if (nextRanges.get(byRange.getKey()) != byRange.getValue()) {
                byRange.getValue().close();
            }
        }
    }
}

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 BaseKVStoreClientPipelineRefreshTest with 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

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
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[BUG] BaseKVStoreClient leaks a mutation pipeline and its gRPC stream on every range-leader change

1 participant