Environment
- BiFroMQ
4.0.0 / main @ eef5e3af
- Module:
base-kv/base-kv-store-client
- File:
base-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java (refreshMutPipelines)
Summary
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(). 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.
Impact
RangeLeaderBalancer performs leadership transfers routinely (it is registered by default in dist-worker / inbox-store deployments). Each transfer leaks one pipeline + one gRPC bidi stream + one server-side pipeline: monotonically growing resource usage on long-running brokers.
Suggested 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();
}
}
}
}
Test
New test class 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
A 6-node BifroMQ 4.0.0 cluster (standalone.sh, systemd, mTLS clientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5, $oshare subscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.
The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale lib/ replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.
Post-deploy checks (all measured after the rollout):
| Check |
Result |
| Cluster mesh (inter-node 8898/8899 connections, per node) |
408–417 established on every node |
| Live client connections (3 bridge replicas + 10 backend services over mTLS) |
13/13 reconnected and stable |
| Cross-node delivery end-to-end |
publish on one node → bridge (connected elsewhere) consumed it → forwarded (counter 0 → 1) |
| Client-visible protocol smoke |
connect / subscribe / publish with QoS 1 all CONNACK 0 / PUBACK 0 |
Baselines captured per node for leak tracking (JVM threads 216–229, RSS 2.8–5.1 GB, gRPC inter-node connections 408–417); the leak manifests as monotonic growth across further leadership transfers, which is being watched.
Environment
4.0.0/main @ eef5e3afbase-kv/base-kv-store-clientbase-kv/base-kv-store-client/src/main/java/org/apache/bifromq/basekv/client/BaseKVStoreClient.java(refreshMutPipelines)Summary
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(). The droppedManagedMutationPipelinekeeps its RX subscription and its dedicated gRPC bidi stream (plus a server-sideMutatePipeline), and it is unreachable fromBaseKVStoreClient.close()once dropped.Impact
RangeLeaderBalancerperforms leadership transfers routinely (it is registered by default in dist-worker / inbox-store deployments). Each transfer leaks one pipeline + one gRPC bidi stream + one server-side pipeline: monotonically growing resource usage on long-running brokers.Suggested 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):
Test
New test class 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
A 6-node BifroMQ 4.0.0 cluster (
standalone.sh, systemd, mTLSclientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5,$osharesubscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale
lib/replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.Post-deploy checks (all measured after the rollout):
Baselines captured per node for leak tracking (JVM threads 216–229, RSS 2.8–5.1 GB, gRPC inter-node connections 408–417); the leak manifests as monotonic growth across further leadership transfers, which is being watched.