This is an automated email from the ASF dual-hosted git repository.

zyk pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new e5b6676ccc [IOTDB-5459] Memory control cross SchemaRegion (#8966)
e5b6676ccc is described below

commit e5b6676cccb0e64cc6571ffd2e2d4645f3e593d8
Author: Chen YZ <[email protected]>
AuthorDate: Sun Feb 5 21:53:17 2023 +0800

    [IOTDB-5459] Memory control cross SchemaRegion (#8966)
---
 .../iotdb/commons/concurrent/ThreadName.java       |   5 +-
 .../db/metadata/mtree/store/CachedMTreeStore.java  |  88 ++-------
 .../mtree/store/disk/MTreeFlushTaskManager.java    |  71 -------
 .../mtree/store/disk/MTreeReleaseTaskManager.java  |  73 --------
 .../mtree/store/disk/cache/CacheManager.java       |   8 +-
 .../mtree/store/disk/cache/CacheMemoryManager.java | 207 +++++++++++++++++++++
 .../db/metadata/rescon/SchemaResourceManager.java  |  10 +-
 7 files changed, 236 insertions(+), 226 deletions(-)

diff --git 
a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
 
b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
index 9cdf23ae5a..588d473841 100644
--- 
a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
+++ 
b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
@@ -61,8 +61,9 @@ public enum ThreadName {
   
ASYNC_DATANODE_HEARTBEAT_CLIENT_POOL("AsyncDataNodeHeartbeatServiceClientPool"),
   ASYNC_CONFIGNODE_CLIENT_POOL("AsyncConfigNodeIServiceClientPool"),
   
ASYNC_DATANODE_MPP_DATA_EXCHANGE_CLIENT_POOL("AsyncDataNodeMPPDataExchangeServiceClientPool"),
-
-  
ASYNC_DATANODE_IOT_CONSENSUS_CLIENT_POOL("AsyncDataNodeMPPDataExchangeServiceClientPool");
+  
ASYNC_DATANODE_IOT_CONSENSUS_CLIENT_POOL("AsyncDataNodeMPPDataExchangeServiceClientPool"),
+  SCHEMA_REGION_RELEASE_POOL("SchemaRegion-Release-Task"),
+  SCHEMA_REGION_FLUSH_POOL("SchemaRegion-Flush-Task");
 
   private final String name;
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/CachedMTreeStore.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/CachedMTreeStore.java
index 45e13b2dca..05429f4b07 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/CachedMTreeStore.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/CachedMTreeStore.java
@@ -31,10 +31,8 @@ import 
org.apache.iotdb.db.metadata.mnode.iterator.AbstractTraverserIterator;
 import org.apache.iotdb.db.metadata.mnode.iterator.CachedTraverserIterator;
 import org.apache.iotdb.db.metadata.mnode.iterator.IMNodeIterator;
 import org.apache.iotdb.db.metadata.mtree.store.disk.ICachedMNodeContainer;
-import org.apache.iotdb.db.metadata.mtree.store.disk.MTreeFlushTaskManager;
-import org.apache.iotdb.db.metadata.mtree.store.disk.MTreeReleaseTaskManager;
+import org.apache.iotdb.db.metadata.mtree.store.disk.cache.CacheMemoryManager;
 import org.apache.iotdb.db.metadata.mtree.store.disk.cache.ICacheManager;
-import org.apache.iotdb.db.metadata.mtree.store.disk.cache.LRUCacheManager;
 import org.apache.iotdb.db.metadata.mtree.store.disk.memcontrol.IMemManager;
 import 
org.apache.iotdb.db.metadata.mtree.store.disk.memcontrol.MemManagerHolder;
 import org.apache.iotdb.db.metadata.mtree.store.disk.schemafile.ISchemaFile;
@@ -59,20 +57,13 @@ public class CachedMTreeStore implements IMTreeStore {
 
   private final IMemManager memManager = 
MemManagerHolder.getMemManagerInstance();
 
-  private final ICacheManager cacheManager = new LRUCacheManager();
+  private final ICacheManager cacheManager =
+      CacheMemoryManager.getInstance().createLRUCacheManager(this);
 
   private ISchemaFile file;
 
   private IMNode root;
 
-  private final MTreeFlushTaskManager flushTaskManager = 
MTreeFlushTaskManager.getInstance();
-  private int flushCount = 0;
-  private volatile boolean hasFlushTask;
-
-  private final MTreeReleaseTaskManager releaseTaskManager = 
MTreeReleaseTaskManager.getInstance();
-  private volatile boolean hasReleaseTask;
-  private int releaseCount = 0;
-
   private final StampedWriterPreferredLock lock = new 
StampedWriterPreferredLock();
 
   public CachedMTreeStore(PartialPath storageGroup, int schemaRegionId)
@@ -80,9 +71,7 @@ public class CachedMTreeStore implements IMTreeStore {
     file = SchemaFile.initSchemaFile(storageGroup.getFullPath(), 
schemaRegionId);
     root = file.init();
     cacheManager.initRootStatus(root);
-
-    hasFlushTask = false;
-    hasReleaseTask = false;
+    ensureMemoryStatus();
   }
 
   @Override
@@ -311,7 +300,7 @@ public class CachedMTreeStore implements IMTreeStore {
   }
 
   @Override
-  public IEntityMNode setToEntity(IMNode node) throws MetadataException {
+  public IEntityMNode setToEntity(IMNode node) {
     IEntityMNode result = MNodeUtils.setToEntity(node);
     if (result != node) {
       memManager.updatePinnedSize(IMNodeSizeEstimator.getEntityNodeBaseSize());
@@ -321,7 +310,7 @@ public class CachedMTreeStore implements IMTreeStore {
   }
 
   @Override
-  public IMNode setToInternal(IEntityMNode entityMNode) throws 
MetadataException {
+  public IMNode setToInternal(IEntityMNode entityMNode) {
     IMNode result = MNodeUtils.setToInternal(entityMNode);
     if (result != entityMNode) {
       
memManager.updatePinnedSize(-IMNodeSizeEstimator.getEntityNodeBaseSize());
@@ -445,9 +434,6 @@ public class CachedMTreeStore implements IMTreeStore {
         }
       }
       file = null;
-
-      hasFlushTask = false;
-      hasReleaseTask = false;
     } finally {
       lock.unlockWrite();
     }
@@ -455,8 +441,14 @@ public class CachedMTreeStore implements IMTreeStore {
 
   @Override
   public boolean createSnapshot(File snapshotDir) {
-    flushVolatileNodes();
-    return file.createSnapshot(snapshotDir);
+    lock.writeLock();
+    try {
+      flushVolatileNodes();
+      ensureMemoryStatus();
+      return file.createSnapshot(snapshotDir);
+    } finally {
+      lock.unlockWrite();
+    }
   }
 
   public static CachedMTreeStore loadFromSnapshot(
@@ -470,49 +462,21 @@ public class CachedMTreeStore implements IMTreeStore {
     file = SchemaFile.loadSnapshot(snapshotDir, storageGroup, schemaRegionId);
     root = file.init();
     cacheManager.initRootStatus(root);
-
-    hasFlushTask = false;
-    hasReleaseTask = false;
   }
 
   private void ensureMemoryStatus() {
-    if (memManager.isExceedFlushThreshold() && !hasReleaseTask) {
-      registerReleaseTask();
-    }
+    CacheMemoryManager.getInstance().ensureMemoryStatus();
   }
 
-  private synchronized void registerReleaseTask() {
-    if (hasReleaseTask) {
-      return;
-    }
-    hasReleaseTask = true;
-    releaseTaskManager.submit(this::tryExecuteMemoryRelease);
-  }
-
-  /**
-   * Execute cache eviction until the memory status is under safe mode or no 
node could be evicted.
-   * If the memory status is still full, which means the nodes in memory are 
all volatile nodes, new
-   * added or updated, fire flush task.
-   */
-  private void tryExecuteMemoryRelease() {
-    lock.threadReadLock();
-    try {
-      executeMemoryRelease();
-      releaseCount++;
-      hasReleaseTask = false;
-    } finally {
-      lock.threadReadUnlock();
-    }
-    if (memManager.isExceedFlushThreshold() && !hasFlushTask) {
-      registerFlushTask();
-    }
+  public StampedWriterPreferredLock getLock() {
+    return lock;
   }
 
   /**
    * Keep fetching evictable nodes from cacheManager until the memory status 
is under safe mode or
    * no node could be evicted. Update the memory status after evicting each 
node.
    */
-  private void executeMemoryRelease() {
+  public void executeMemoryRelease() {
     while (memManager.isExceedReleaseThreshold() && !memManager.isEmpty()) {
       if (!cacheManager.evict()) {
         break;
@@ -520,17 +484,8 @@ public class CachedMTreeStore implements IMTreeStore {
     }
   }
 
-  private synchronized void registerFlushTask() {
-    if (hasFlushTask) {
-      return;
-    }
-    hasFlushTask = true;
-    flushTaskManager.submit(this::flushVolatileNodes);
-  }
-
   /** Sync all volatile nodes to schemaFile and execute memory release after 
flush. */
-  private void flushVolatileNodes() {
-    lock.writeLock();
+  public void flushVolatileNodes() {
     try {
       IStorageGroupMNode updatedStorageGroupMNode = 
cacheManager.collectUpdatedStorageGroupMNodes();
       if (updatedStorageGroupMNode != null) {
@@ -557,15 +512,10 @@ public class CachedMTreeStore implements IMTreeStore {
         }
         cacheManager.updateCacheStatusAfterPersist(volatileNode);
       }
-      executeMemoryRelease();
-      hasFlushTask = false;
-      flushCount++;
     } catch (Throwable e) {
       logger.error(
           "Error occurred during MTree flush, current SchemaRegion is {}", 
root.getFullPath(), e);
       e.printStackTrace();
-    } finally {
-      lock.unlockWrite();
     }
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeFlushTaskManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeFlushTaskManager.java
deleted file mode 100644
index 6a4a61de10..0000000000
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeFlushTaskManager.java
+++ /dev/null
@@ -1,71 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied.  See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-package org.apache.iotdb.db.metadata.mtree.store.disk;
-
-import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.concurrent.ExecutorService;
-
-public class MTreeFlushTaskManager {
-
-  private static final Logger logger = 
LoggerFactory.getLogger(MTreeFlushTaskManager.class);
-  private static final String MTREE_FLUSH_THREAD_POOL_NAME = 
"MTree-flush-task";
-
-  private volatile ExecutorService flushTaskExecutor;
-
-  private MTreeFlushTaskManager() {}
-
-  private static class MTreeFlushTaskManagerHolder {
-    private static final MTreeFlushTaskManager INSTANCE = new 
MTreeFlushTaskManager();
-
-    private MTreeFlushTaskManagerHolder() {}
-  }
-
-  public static MTreeFlushTaskManager getInstance() {
-    return MTreeFlushTaskManagerHolder.INSTANCE;
-  }
-
-  public void init() {
-    flushTaskExecutor = 
IoTDBThreadPoolFactory.newCachedThreadPool(MTREE_FLUSH_THREAD_POOL_NAME);
-  }
-
-  public void clear() {
-    if (flushTaskExecutor != null) {
-      flushTaskExecutor.shutdown();
-      while (!flushTaskExecutor.isTerminated()) ;
-      flushTaskExecutor = null;
-    }
-  }
-
-  public void submit(Runnable task) {
-    flushTaskExecutor.submit(
-        () -> {
-          try {
-            task.run();
-          } catch (Throwable throwable) {
-            logger.error("Something wrong happened during MTree flush.", 
throwable);
-            throwable.printStackTrace();
-            throw throwable;
-          }
-        });
-  }
-}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeReleaseTaskManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeReleaseTaskManager.java
deleted file mode 100644
index 7ba8cceabe..0000000000
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/MTreeReleaseTaskManager.java
+++ /dev/null
@@ -1,73 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied.  See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-package org.apache.iotdb.db.metadata.mtree.store.disk;
-
-import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.concurrent.ExecutorService;
-
-public class MTreeReleaseTaskManager {
-
-  private static final Logger logger = 
LoggerFactory.getLogger(MTreeReleaseTaskManager.class);
-  private static final String MTREE_RELEASE_THREAD_POOL_NAME = 
"MTree-release-task";
-
-  private volatile ExecutorService releaseTaskExecutor;
-
-  private MTreeReleaseTaskManager() {}
-
-  private static class MTreeReleaseTaskManagerHolder {
-    private static final MTreeReleaseTaskManager INSTANCE = new 
MTreeReleaseTaskManager();
-
-    private MTreeReleaseTaskManagerHolder() {}
-  }
-
-  public static MTreeReleaseTaskManager getInstance() {
-    return MTreeReleaseTaskManager.MTreeReleaseTaskManagerHolder.INSTANCE;
-  }
-
-  public void init() {
-    releaseTaskExecutor =
-        
IoTDBThreadPoolFactory.newCachedThreadPool(MTREE_RELEASE_THREAD_POOL_NAME);
-  }
-
-  public void clear() {
-    if (releaseTaskExecutor != null) {
-      releaseTaskExecutor.shutdown();
-      while (!releaseTaskExecutor.isTerminated()) ;
-      releaseTaskExecutor = null;
-    }
-  }
-
-  public void submit(Runnable task) {
-    releaseTaskExecutor.submit(
-        () -> {
-          try {
-            task.run();
-          } catch (Throwable throwable) {
-            logger.error("Something wrong happened during MTree release.", 
throwable);
-            throwable.printStackTrace();
-            throw throwable;
-          }
-        });
-  }
-}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheManager.java
index 985bd0ac01..4e9fb1819f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheManager.java
@@ -173,11 +173,11 @@ public abstract class CacheManager implements 
ICacheManager {
   public void updateCacheStatusAfterUpdate(IMNode node) {
     CacheEntry cacheEntry = getCacheEntry(node);
     if (!cacheEntry.isVolatile()) {
-      synchronized (cacheEntry) {
-        // the status change affects the subTre collect in nodeBuffer
-        cacheEntry.setVolatile(true);
-      }
       if (!node.isStorageGroup()) {
+        synchronized (cacheEntry) {
+          // the status change affects the subTre collect in nodeBuffer
+          cacheEntry.setVolatile(true);
+        }
         // if node is StorageGroup, getBelongedContainer is null
         getBelongedContainer(node).updateMNode(node.getName());
         // MNode update operation like node replace may reset the mapping 
between cacheEntry and
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheMemoryManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheMemoryManager.java
new file mode 100644
index 0000000000..031826ac1c
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/store/disk/cache/CacheMemoryManager.java
@@ -0,0 +1,207 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.db.metadata.mtree.store.disk.cache;
+
+import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
+import org.apache.iotdb.commons.concurrent.ThreadName;
+import org.apache.iotdb.db.metadata.mtree.store.CachedMTreeStore;
+import org.apache.iotdb.db.metadata.mtree.store.disk.memcontrol.IMemManager;
+import 
org.apache.iotdb.db.metadata.mtree.store.disk.memcontrol.MemManagerHolder;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+
+/**
+ * CacheMemoryManager is used to register the CachedMTreeStore and create the 
CacheManager.
+ * CacheMemoryManager provides the {@link 
CacheMemoryManager#ensureMemoryStatus} interface, which
+ * starts asynchronous threads to free and flush the disk when memory usage 
exceeds a threshold.
+ */
+public class CacheMemoryManager {
+
+  private static final Logger logger = 
LoggerFactory.getLogger(CacheMemoryManager.class);
+
+  private final List<CachedMTreeStore> storeList = new ArrayList<>();
+
+  private final IMemManager memManager = 
MemManagerHolder.getMemManagerInstance();
+
+  private static final int CONCURRENT_NUM = 10;
+
+  private ExecutorService flushTaskExecutor;
+  private ExecutorService releaseTaskExecutor;
+
+  private volatile boolean hasFlushTask;
+  private int flushCount = 0;
+
+  private volatile boolean hasReleaseTask;
+  private int releaseCount = 0;
+
+  public synchronized ICacheManager createLRUCacheManager(CachedMTreeStore 
store) {
+    synchronized (storeList) {
+      ICacheManager cacheManager = new LRUCacheManager();
+      storeList.add(store);
+      return cacheManager;
+    }
+  }
+
+  public void init() {
+    flushTaskExecutor =
+        IoTDBThreadPoolFactory.newFixedThreadPool(
+            CONCURRENT_NUM, ThreadName.SCHEMA_REGION_FLUSH_POOL.getName());
+    releaseTaskExecutor =
+        IoTDBThreadPoolFactory.newFixedThreadPool(
+            CONCURRENT_NUM, ThreadName.SCHEMA_REGION_RELEASE_POOL.getName());
+  }
+
+  public void ensureMemoryStatus() {
+    if (memManager.isExceedReleaseThreshold() && !hasReleaseTask) {
+      registerReleaseTask();
+    }
+  }
+
+  private synchronized void registerReleaseTask() {
+    if (hasReleaseTask) {
+      return;
+    }
+    hasReleaseTask = true;
+    releaseTaskExecutor.submit(
+        () -> {
+          try {
+            tryExecuteMemoryRelease();
+          } catch (Throwable throwable) {
+            logger.error("Something wrong happened during MTree release.", 
throwable);
+            throwable.printStackTrace();
+            throw throwable;
+          }
+        });
+  }
+
+  /**
+   * Execute cache eviction until the memory status is under safe mode or no 
node could be evicted.
+   * If the memory status is still full, which means the nodes in memory are 
all volatile nodes, new
+   * added or updated, fire flush task.
+   */
+  private void tryExecuteMemoryRelease() {
+    synchronized (storeList) {
+      CompletableFuture.allOf(
+              storeList.stream()
+                  .map(
+                      store ->
+                          CompletableFuture.runAsync(
+                              () -> {
+                                store.getLock().threadReadLock();
+                                try {
+                                  store.executeMemoryRelease();
+                                } finally {
+                                  store.getLock().threadReadUnlock();
+                                }
+                              },
+                              releaseTaskExecutor))
+                  .toArray(CompletableFuture[]::new))
+          .join();
+      releaseCount++;
+      hasReleaseTask = false;
+      if (memManager.isExceedFlushThreshold() && !hasFlushTask) {
+        registerFlushTask();
+      }
+    }
+  }
+
+  private synchronized void registerFlushTask() {
+    if (hasFlushTask) {
+      return;
+    }
+    hasFlushTask = true;
+    flushTaskExecutor.submit(
+        () -> {
+          try {
+            tryFlushVolatileNodes();
+          } catch (Throwable throwable) {
+            logger.error("Something wrong happened during MTree flush.", 
throwable);
+            throwable.printStackTrace();
+            throw throwable;
+          }
+        });
+  }
+
+  /** Sync all volatile nodes to schemaFile and execute memory release after 
flush. */
+  private void tryFlushVolatileNodes() {
+    synchronized (storeList) {
+      CompletableFuture.allOf(
+              storeList.stream()
+                  .map(
+                      store ->
+                          CompletableFuture.runAsync(
+                              () -> {
+                                store.getLock().writeLock();
+                                try {
+                                  store.flushVolatileNodes();
+                                  store.executeMemoryRelease();
+                                } finally {
+                                  store.getLock().unlockWrite();
+                                }
+                              },
+                              flushTaskExecutor))
+                  .toArray(CompletableFuture[]::new))
+          .join();
+      hasFlushTask = false;
+      flushCount++;
+    }
+  }
+
+  public void clear() {
+    if (releaseTaskExecutor != null) {
+      while (true) {
+        if (!hasReleaseTask) break;
+      }
+      releaseTaskExecutor.shutdown();
+      while (true) {
+        if (releaseTaskExecutor.isTerminated()) break;
+      }
+      releaseTaskExecutor = null;
+    }
+    // the release task may submit flush task, thus must be shut down and 
clear first
+    if (flushTaskExecutor != null) {
+      while (true) {
+        if (!hasFlushTask) break;
+      }
+      flushTaskExecutor.shutdown();
+      while (true) {
+        if (flushTaskExecutor.isTerminated()) break;
+      }
+      flushTaskExecutor = null;
+    }
+  }
+
+  private CacheMemoryManager() {}
+
+  private static class GlobalCacheManagerHolder {
+    private static final CacheMemoryManager INSTANCE = new 
CacheMemoryManager();
+
+    private GlobalCacheManagerHolder() {}
+  }
+
+  public static CacheMemoryManager getInstance() {
+    return CacheMemoryManager.GlobalCacheManagerHolder.INSTANCE;
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/rescon/SchemaResourceManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/rescon/SchemaResourceManager.java
index d46cb7addf..771a010767 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/rescon/SchemaResourceManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/rescon/SchemaResourceManager.java
@@ -21,8 +21,7 @@ package org.apache.iotdb.db.metadata.rescon;
 
 import org.apache.iotdb.commons.service.metric.MetricService;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import org.apache.iotdb.db.metadata.mtree.store.disk.MTreeFlushTaskManager;
-import org.apache.iotdb.db.metadata.mtree.store.disk.MTreeReleaseTaskManager;
+import org.apache.iotdb.db.metadata.mtree.store.disk.cache.CacheMemoryManager;
 import 
org.apache.iotdb.db.metadata.mtree.store.disk.memcontrol.MemManagerHolder;
 import org.apache.iotdb.db.metadata.schemaregion.SchemaEngineMode;
 
@@ -58,14 +57,11 @@ public class SchemaResourceManager {
   private static void initSchemaFileModeResource() {
     MemManagerHolder.initMemManagerInstance();
     MemManagerHolder.getMemManagerInstance().init();
-    MTreeFlushTaskManager.getInstance().init();
-    MTreeReleaseTaskManager.getInstance().init();
+    CacheMemoryManager.getInstance().init();
   }
 
   private static void clearSchemaFileModeResource() {
     MemManagerHolder.getMemManagerInstance().clear();
-    // the release task may submit flush task, thus must be shut down and 
clear first
-    MTreeReleaseTaskManager.getInstance().clear();
-    MTreeFlushTaskManager.getInstance().clear();
+    CacheMemoryManager.getInstance().clear();
   }
 }

Reply via email to