This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch fix_initCompactionSchedule in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit eb8af1537166167b996626b54f653f250a6a09f5 Author: HTHou <[email protected]> AuthorDate: Tue Oct 31 19:09:48 2023 +0800 Fix Compaction Schedule didn't start --- .../iotdb/db/storageengine/StorageEngine.java | 22 +++++++++++++--------- .../db/storageengine/dataregion/DataRegion.java | 2 +- 2 files changed, 14 insertions(+), 10 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java index dcfbc918f74..58b9c28d984 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java @@ -193,7 +193,7 @@ public class StorageEngine implements IService { isAllSgReady.set(allSgReady); } - public void recover() throws StartupException { + public void asyncRecover() throws StartupException { setAllSgReady(false); cachedThreadPool = IoTDBThreadPoolFactory.newCachedThreadPool(ThreadName.STORAGE_ENGINE_CACHED_POOL.getName()); @@ -218,11 +218,20 @@ public class StorageEngine implements IService { checkResults(futures, "StorageEngine failed to recover."); setAllSgReady(true); ttlMapForRecover.clear(); + initCompactionSchedule(); }, ThreadName.STORAGE_ENGINE_RECOVER_TRIGGER.getName()); recoverEndTrigger.start(); } + private void initCompactionSchedule() { + for (DataRegion dataRegion : dataRegionMap.values()) { + if (dataRegion != null) { + dataRegion.initCompactionSchedule(); + } + } + } + private void asyncRecover(List<Future<Void>> futures) { Map<String, List<DataRegionId>> localDataRegionInfo = getLocalDataRegionInfo(); localDataRegionInfo.values().forEach(list -> recoverDataRegionNum += list.size()); @@ -295,12 +304,7 @@ public class StorageEngine implements IService { throw new StorageEngineFailureException(e); } - recover(); - for (DataRegion dataRegion : dataRegionMap.values()) { - if (dataRegion != null) { - dataRegion.initCompaction(); - } - } + asyncRecover(); ttlCheckThread = IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.TTL_CHECK.getName()); @@ -654,8 +658,8 @@ public class StorageEngine implements IService { region.deleteFolder(systemDir); if (config.isClusterMode() && config - .getDataRegionConsensusProtocolClass() - .equals(ConsensusFactory.IOT_CONSENSUS)) { + .getDataRegionConsensusProtocolClass() + .equals(ConsensusFactory.IOT_CONSENSUS)) { // delete wal WALManager.getInstance() .deleteWALNode( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java index 42780fa177f..5dd46accb01 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java @@ -600,7 +600,7 @@ public class DataRegion implements IDataRegionForQuery { lastFlushTimeMap.setMultiDeviceGlobalFlushedTime(endTimeMap); } - public void initCompaction() { + public void initCompactionSchedule() { if (!config.isEnableSeqSpaceCompaction() && !config.isEnableUnseqSpaceCompaction() && !config.isEnableCrossSpaceCompaction()
