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();
}
}