From a39deb367aac8cb7a957d9ef947e2c3cf83a40ad Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Thu, 24 Sep 2026 09:25:59 +0300 Subject: [PATCH] Shut down KV space compaction executors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../localengine/rocksdb/RocksDBKVSpace.java | 20 +++++++++++++++- .../rocksdb/RocksDBCPableKVEngineTest.java | 24 +++++++++++++++++++ 2 files changed, 43 insertions(+), 1 deletion(-) diff --git a/base-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java b/base-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java index adae463e9..16f821b2a 100644 --- a/base-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java +++ b/base-kv/base-kv-local-engine-rocksdb/src/main/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBKVSpace.java @@ -36,6 +36,7 @@ import com.google.protobuf.ByteString; import com.google.protobuf.Struct; import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.Meter; import io.micrometer.core.instrument.Metrics; import io.micrometer.core.instrument.Tags; import io.micrometer.core.instrument.Timer; @@ -43,6 +44,7 @@ import java.io.File; import java.io.IOException; import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.LinkedBlockingQueue; @@ -69,6 +71,7 @@ abstract class RocksDBKVSpace extends AbstractKVSpace protected final IWriteStatsRecorder writeStats; private final File keySpaceDBDir; private final ExecutorService compactionExecutor; + private final List compactionMeterIds; private final AtomicBoolean compacting; private final ISyncContext.IRefresher metadataRefresher; private SpaceMetrics spaceMetrics; @@ -89,10 +92,17 @@ public RocksDBKVSpace(String id, syncContext = new SyncContext(); metadataRefresher = syncContext.refresher(); compacting = new AtomicBoolean(false); + Tags compactionMetricTags = this.tags; compactionExecutor = ExecutorServiceMetrics.monitor(Metrics.globalRegistry, new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>(), EnvProvider.INSTANCE.newThreadFactory("kvspace-compactor-" + id)), - "compactor", "kvspace", Tags.of(tags)); + "compactor", "kvspace", compactionMetricTags); + compactionMeterIds = Metrics.globalRegistry.getMeters().stream() + .map(Meter::getId) + .filter(meterId -> meterId.getName().startsWith("kvspace.executor")) + .filter(meterId -> compactionMetricTags.stream() + .allMatch(tag -> tag.getValue().equals(meterId.getTag(tag.getKey())))) + .toList(); if (boolVal(conf, MANUAL_COMPACTION)) { int minKeys = (int) numVal(conf, COMPACT_MIN_TOMBSTONE_KEYS); int minRanges = (int) numVal(conf, COMPACT_MIN_TOMBSTONE_RANGES); @@ -130,6 +140,8 @@ protected void publishMetadata(Map metadataUpdates) { @Override protected void doClose() { logger.debug("Close key range[{}]", id); + compactionExecutor.shutdownNow(); + unregisterCompactionMetrics(); if (spaceMetrics != null) { spaceMetrics.close(); } @@ -137,6 +149,8 @@ protected void doClose() { @Override protected void doDestroy() { + compactionExecutor.shutdownNow(); + unregisterCompactionMetrics(); // Destroy the whole space root directory, including pointer file and all generations. try { if (keySpaceDBDir.exists()) { @@ -147,6 +161,10 @@ protected void doDestroy() { } } + private void unregisterCompactionMetrics() { + compactionMeterIds.forEach(Metrics.globalRegistry::remove); + } + protected File spaceRootDir() { return keySpaceDBDir; } diff --git a/base-kv/base-kv-local-engine-rocksdb/src/test/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBCPableKVEngineTest.java b/base-kv/base-kv-local-engine-rocksdb/src/test/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBCPableKVEngineTest.java index df0974467..dce5518b7 100644 --- a/base-kv/base-kv-local-engine-rocksdb/src/test/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBCPableKVEngineTest.java +++ b/base-kv/base-kv-local-engine-rocksdb/src/test/java/org/apache/bifromq/basekv/localengine/rocksdb/RocksDBCPableKVEngineTest.java @@ -21,14 +21,18 @@ import static org.apache.bifromq.basekv.localengine.rocksdb.RocksDBDefaultConfigs.DB_CHECKPOINT_ROOT_DIR; import static org.apache.bifromq.basekv.localengine.rocksdb.RocksDBDefaultConfigs.DB_ROOT_DIR; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import com.google.protobuf.Struct; import com.google.protobuf.Value; +import io.micrometer.core.instrument.Metrics; import io.reactivex.rxjava3.disposables.Disposable; import java.io.File; +import java.lang.reflect.Field; import java.nio.file.Paths; +import java.util.concurrent.ExecutorService; import lombok.SneakyThrows; import org.apache.bifromq.basekv.localengine.ICPableKVSpace; import org.apache.bifromq.basekv.localengine.IKVEngine; @@ -71,4 +75,24 @@ public void removeCheckpointFileWhenDestroy() { assertTrue(engine.spaces().isEmpty()); assertFalse(engine.spaces().containsKey(rangeId)); } + + @SneakyThrows + @Test + public void shutdownCompactionExecutorWhenSpaceClosed() { + String rangeId = "compaction_range"; + IKVSpace range = engine.createIfMissing(rangeId); + Field executorField = RocksDBKVSpace.class.getDeclaredField("compactionExecutor"); + executorField.setAccessible(true); + ExecutorService executor = (ExecutorService) executorField.get(range); + + assertFalse(executor.isShutdown()); + range.close(); + assertTrue(executor.isShutdown()); + assertEquals( + Metrics.globalRegistry.getMeters().stream() + .filter(meter -> meter.getId().getName().startsWith("kvspace.executor")) + .filter(meter -> rangeId.equals(meter.getId().getTag("spaceId"))) + .count(), + 0); + } }