Skip to content

[BUG] Per-KV-space RocksDB compaction executor is never shut down: non-daemon threads accumulate and prevent JVM exit #290

Description

@ImDanXie

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:

  1. EnvProvider.INSTANCE.newThreadFactory(name) creates non-daemon threads by default
  2. new ThreadPoolExecutor(1, 1, 0L, MILLISECONDS, new LinkedBlockingQueue<>()) does not enable allowCoreThreadTimeOut, so the core thread never exits once a task has been submitted
  3. 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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions