Environment
- BiFroMQ
4.0.0
- Verified against
main @ ba2c34e0 — still present on the latest main at the time of filing
- Module:
base-kv
- File:
base-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java:71, 92-95, 131-149
Summary
Each KV space (i.e. each range) creates its own compactionExecutor, but neither doClose() nor doDestroy() shuts it down. The field is referenced in only three places in the whole class — declaration, creation, and execute — with no release point anywhere.
Three factors compound the leak:
EnvProvider.INSTANCE.newThreadFactory(name) creates non-daemon threads by default
new ThreadPoolExecutor(1, 1, 0L, MILLISECONDS, new LinkedBlockingQueue<>()) does not enable allowCoreThreadTimeOut, so the core thread never exits once a task has been submitted
ExecutorServiceMetrics.monitor(Metrics.globalRegistry, ...) registers meters that strongly reference the executor and are never removed
Steps to reproduce
RocksDBKVEngineProvider provider = new RocksDBKVEngineProvider();
// dbRootDir must exist and contain an IDENTITY file
IKVEngine<? extends ICPableKVSpace> engine = provider.createCPable("probe", conf);
engine.start();
ICPableKVSpace space = engine.createIfMissing("probeSpace");
space.open();
// reflectively obtain compactionExecutor, invoke scheduleCompact() once, then close
space.close();
Observed (before fix)
close 前 compactionExecutor.isShutdown() = false
close 后 compactionExecutor.isShutdown() = false
[ BUG ] close() 应关闭 compactionExecutor
[ BUG ] close 之后仍可向 compactionExecutor 提交任务 —— 线程池未被关闭
存活线程: kvspace-compactor-probeSpace daemon=false
=== engine.stop() 之后仍存活的所有非守护线程 ===
[非守护] kvspace-compactor-probeSpace
The strongest signal is that the probe process does not terminate after engine.stop() — it had to be killed by a timeout. Enumerating all threads confirms the surviving non-daemon thread is uniquely the compaction one, so the cause of the hang is not merely correlated with the leak.
Impact
With default settings a range triggering compaction once leaves a permanent thread behind. (RocksDBKVSpaceCompactionTrigger.recordPut() increments both keyCount and tombstoneKeyCount, so the tombstone ratio reaches N/(N+N) = 0.5 on pure puts — the totalTombstones > 200000 && ratio >= 0.3 trigger fires after 200k puts even without any deletes.)
Ranges are split and merged repeatedly, so the thread count grows without bound (~1 MB stack each), eventually producing OutOfMemoryError: unable to create native thread. The non-daemon property additionally prevents the JVM from shutting down cleanly.
Suggested fix
@Override
protected void doClose() {
logger.debug("Close key range[{}]", id);
compactionExecutor.shutdownNow();
unregisterCompactionMetrics();
if (spaceMetrics != null) {
spaceMetrics.close();
}
}
@Override
protected void doDestroy() {
compactionExecutor.shutdownNow();
unregisterCompactionMetrics();
...
}
Metric registration is released as well
ExecutorServiceMetrics.monitor(Metrics.globalRegistry, ...) registers nine meters per KV space (kvspace.executor, kvspace.executor.active, .queued, .pool.size, .pool.core, .pool.max, .completed, .idle, .queue.remaining). Those were also never removed, so the global registry accumulated one set per retired range on every split/merge.
Micrometer 1.11.12's ExecutorServiceMetrics exposes no close() and monitor(...) returns the decorated executor rather than the binder, so the fix removes the meters by name + tags, using the same Tags instance that was passed to monitor(...):
compactionMetricTags = Tags.of(tags);
compactionExecutor = ExecutorServiceMetrics.monitor(Metrics.globalRegistry, ..., "compactor", "kvspace",
compactionMetricTags);
private void unregisterCompactionMetrics() {
Metrics.globalRegistry.getMeters().stream()
.filter(m -> m.getId().getName().startsWith("kvspace.executor"))
.filter(m -> compactionMetricTags.stream()
.allMatch(t -> t.getValue().equals(m.getId().getTag(t.getKey()))))
.collect(Collectors.toList())
.forEach(Metrics.globalRegistry::remove);
}
Matching on the full tag set keeps the removal scoped to this space: AbstractKVEngine.buildKVSpace appends ("spaceId", spaceId) to the tags, so each space's meters carry a unique spaceId.
Measured on the same probe:
|
before fix |
after fix |
meters remaining after close() / stop() |
9 |
0 |
This accumulation was bounded-memory and was not the cause of the OutOfMemoryError above — that was the thread leak — but it is fixed in the same patch since it shares the lifetime.
Observed (after fix)
close 后 compactionExecutor.isShutdown() = true
[ OK ] close() 应关闭 compactionExecutor
[ OK ] close 之后提交任务被拒绝 —— 线程池已关闭
=== engine.stop() 之后仍存活的所有非守护线程 ===
Process exits with code 0 instead of being killed by a timeout, and the global registry retains no kvspace.executor meters.
Environment
4.0.0main@ba2c34e0— still present on the latestmainat the time of filingbase-kvbase-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java:71, 92-95, 131-149Summary
Each KV space (i.e. each range) creates its own
compactionExecutor, but neitherdoClose()nordoDestroy()shuts it down. The field is referenced in only three places in the whole class — declaration, creation, andexecute— with no release point anywhere.Three factors compound the leak:
EnvProvider.INSTANCE.newThreadFactory(name)creates non-daemon threads by defaultnew ThreadPoolExecutor(1, 1, 0L, MILLISECONDS, new LinkedBlockingQueue<>())does not enableallowCoreThreadTimeOut, so the core thread never exits once a task has been submittedExecutorServiceMetrics.monitor(Metrics.globalRegistry, ...)registers meters that strongly reference the executor and are never removedSteps to reproduce
Observed (before fix)
The strongest signal is that the probe process does not terminate after
engine.stop()— it had to be killed by a timeout. Enumerating all threads confirms the surviving non-daemon thread is uniquely the compaction one, so the cause of the hang is not merely correlated with the leak.Impact
With default settings a range triggering compaction once leaves a permanent thread behind. (
RocksDBKVSpaceCompactionTrigger.recordPut()increments bothkeyCountandtombstoneKeyCount, so the tombstone ratio reachesN/(N+N) = 0.5on pure puts — thetotalTombstones > 200000 && ratio >= 0.3trigger fires after 200k puts even without any deletes.)Ranges are split and merged repeatedly, so the thread count grows without bound (~1 MB stack each), eventually producing
OutOfMemoryError: unable to create native thread. The non-daemon property additionally prevents the JVM from shutting down cleanly.Suggested fix
Metric registration is released as well
ExecutorServiceMetrics.monitor(Metrics.globalRegistry, ...)registers nine meters per KV space (kvspace.executor,kvspace.executor.active,.queued,.pool.size,.pool.core,.pool.max,.completed,.idle,.queue.remaining). Those were also never removed, so the global registry accumulated one set per retired range on every split/merge.Micrometer 1.11.12's
ExecutorServiceMetricsexposes noclose()andmonitor(...)returns the decorated executor rather than the binder, so the fix removes the meters by name + tags, using the sameTagsinstance that was passed tomonitor(...):Matching on the full tag set keeps the removal scoped to this space:
AbstractKVEngine.buildKVSpaceappends("spaceId", spaceId)to the tags, so each space's meters carry a uniquespaceId.Measured on the same probe:
close()/stop()This accumulation was bounded-memory and was not the cause of the
OutOfMemoryErrorabove — that was the thread leak — but it is fixed in the same patch since it shares the lifetime.Observed (after fix)
Process exits with code 0 instead of being killed by a timeout, and the global registry retains no
kvspace.executormeters.