This is an automated email from the ASF dual-hosted git repository.
qiaojialin 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 d984fe8 [IOTDB-1496] Timed flush memtable (#3610)
d984fe8 is described below
commit d984fe834d98c270f3f35a1499703410935774ca
Author: Alan Choo <[email protected]>
AuthorDate: Tue Jul 27 17:57:22 2021 +0800
[IOTDB-1496] Timed flush memtable (#3610)
---
.../resources/conf/iotdb-engine.properties | 15 +++++
.../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 36 +++++++++++
.../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 28 +++++++++
.../org/apache/iotdb/db/engine/StorageEngine.java | 43 ++++++++++++++
.../iotdb/db/engine/memtable/AbstractMemTable.java | 7 +++
.../apache/iotdb/db/engine/memtable/IMemTable.java | 2 +
.../engine/storagegroup/StorageGroupProcessor.java | 22 +++++++
.../db/engine/storagegroup/TsFileProcessor.java | 5 ++
.../virtualSg/VirtualStorageGroupManager.java | 9 +++
.../apache/iotdb/db/rescon/MemTableManager.java | 4 ++
.../storagegroup/StorageGroupProcessorTest.java | 69 ++++++++++++++++++----
.../apache/iotdb/db/utils/EnvironmentUtils.java | 4 ++
12 files changed, 233 insertions(+), 11 deletions(-)
diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 9eb19dd..63e79d4 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -227,6 +227,21 @@ timestamp_precision=ms
# Datatype: long
# memtable_size_threshold=1073741824
+# Whether to timed flush unsequence tsfiles' memtables.
+# Datatype: boolean
+# enable_timed_flush_unseq_memtable=false
+
+# When a memTable's created time is older than current time minus this, the
memtable is flushed to disk.
+# Only check unsequence tsfiles' memtables.
+# The default flush interval is 12 * 60 * 60 * 1000. (unit: ms)
+# Datatype: long
+# unseq_memtable_flush_interval_in_ms=43200000
+
+# The interval to check whether the memtable needs flushing.
+# The default flush interval is 1 * 60 * 60 * 1000. (unit: ms)
+# Datatype: long
+# unseq_memtable_flush_check_interval_in_ms=3600000
+
# When the average point number of timeseries in memtable exceeds this, the
memtable is flushed to disk. The default threshold is 10000.
# Datatype: int
# avg_series_point_number_threshold=10000
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index 4b8d6b8..3979903 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -295,6 +295,18 @@ public class IoTDBConfig {
/** When a memTable's size (in byte) exceeds this, the memtable is flushed
to disk. Unit: byte */
private long memtableSizeThreshold = 1024 * 1024 * 1024L;
+ /** Whether to timed flush unsequence tsfiles' memtables. */
+ private boolean enableTimedFlushUnseqMemtable = false;
+
+ /**
+ * When a memTable's created time is older than current time minus this, the
memtable is flushed
+ * to disk.(only check unsequence tsfiles' memtables) Unit: ms
+ */
+ private long unseqMemtableFlushInterval = 12 * 60 * 60 * 1000L;
+
+ /** The interval to check whether the memtable needs flushing. Unit: ms */
+ private long unseqMemtableFlushCheckInterval = 60 * 60 * 1000L;
+
/** When average series point number reaches this, flush the memtable to
disk */
private int avgSeriesPointNumberThreshold = 10000;
@@ -1500,6 +1512,30 @@ public class IoTDBConfig {
this.memtableSizeThreshold = memtableSizeThreshold;
}
+ public boolean isEnableTimedFlushUnseqMemtable() {
+ return enableTimedFlushUnseqMemtable;
+ }
+
+ public void setEnableTimedFlushUnseqMemtable(boolean
enableTimedFlushUnseqMemtable) {
+ this.enableTimedFlushUnseqMemtable = enableTimedFlushUnseqMemtable;
+ }
+
+ public long getUnseqMemtableFlushInterval() {
+ return unseqMemtableFlushInterval;
+ }
+
+ public void setUnseqMemtableFlushInterval(long unseqMemtableFlushInterval) {
+ this.unseqMemtableFlushInterval = unseqMemtableFlushInterval;
+ }
+
+ public long getUnseqMemtableFlushCheckInterval() {
+ return unseqMemtableFlushCheckInterval;
+ }
+
+ public void setUnseqMemtableFlushCheckInterval(long
unseqMemtableFlushCheckInterval) {
+ this.unseqMemtableFlushCheckInterval = unseqMemtableFlushCheckInterval;
+ }
+
public int getAvgSeriesPointNumberThreshold() {
return avgSeriesPointNumberThreshold;
}
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index a3f271a..1b0e4eb 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -271,6 +271,34 @@ public class IoTDBDescriptor {
conf.setMemtableSizeThreshold(memTableSizeThreshold);
}
+ conf.setEnableTimedFlushUnseqMemtable(
+ Boolean.parseBoolean(
+ properties.getProperty(
+ "enable_timed_flush_unseq_memtable",
+ Boolean.toString(conf.isEnableTimedFlushUnseqMemtable()))));
+
+ long unseqMemTableFlushInterval =
+ Long.parseLong(
+ properties
+ .getProperty(
+ "unseq_memtable_flush_interval_in_ms",
+ Long.toString(conf.getUnseqMemtableFlushInterval()))
+ .trim());
+ if (unseqMemTableFlushInterval > 0) {
+ conf.setUnseqMemtableFlushInterval(unseqMemTableFlushInterval);
+ }
+
+ long unseqMemTableFlushCheckInterval =
+ Long.parseLong(
+ properties
+ .getProperty(
+ "unseq_memtable_flush_check_interval_in_ms",
+ Long.toString(conf.getUnseqMemtableFlushCheckInterval()))
+ .trim());
+ if (unseqMemTableFlushCheckInterval > 0) {
+
conf.setUnseqMemtableFlushCheckInterval(unseqMemTableFlushCheckInterval);
+ }
+
conf.setAvgSeriesPointNumberThreshold(
Integer.parseInt(
properties.getProperty(
diff --git a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
index 8250d5f..2150984 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/StorageEngine.java
@@ -101,6 +101,7 @@ public class StorageEngine implements IService {
private static final IoTDBConfig config =
IoTDBDescriptor.getInstance().getConfig();
private static final long TTL_CHECK_INTERVAL = 60 * 1000L;
+
/**
* Time range for dividing storage group, the time unit is the same with
IoTDB's
* TimestampPrecision
@@ -123,6 +124,7 @@ public class StorageEngine implements IService {
private AtomicBoolean isAllSgReady = new AtomicBoolean(false);
private ScheduledExecutorService ttlCheckThread;
+ private ScheduledExecutorService memtableTimedFlushCheckThread;
private TsFileFlushPolicy fileFlushPolicy = new DirectFlushPolicy();
private ExecutorService recoveryThreadPool;
// add customized listeners here for flush and close events
@@ -273,6 +275,17 @@ public class StorageEngine implements IService {
ttlCheckThread = Executors.newSingleThreadScheduledExecutor();
ttlCheckThread.scheduleAtFixedRate(
this::checkTTL, TTL_CHECK_INTERVAL, TTL_CHECK_INTERVAL,
TimeUnit.MILLISECONDS);
+ logger.info("start ttl check thread successfully.");
+
+ if (config.isEnableTimedFlushUnseqMemtable()) {
+ memtableTimedFlushCheckThread =
Executors.newSingleThreadScheduledExecutor();
+ memtableTimedFlushCheckThread.scheduleAtFixedRate(
+ this::timedFlushMemTable,
+ config.getUnseqMemtableFlushCheckInterval(),
+ config.getUnseqMemtableFlushCheckInterval(),
+ TimeUnit.MILLISECONDS);
+ logger.info("start memtable timed flush check thread successfully.");
+ }
}
private void checkTTL() {
@@ -287,6 +300,16 @@ public class StorageEngine implements IService {
}
}
+ private void timedFlushMemTable() {
+ try {
+ for (VirtualStorageGroupManager processor : processorMap.values()) {
+ processor.timedFlushMemTable();
+ }
+ } catch (Exception e) {
+ logger.error("An error occurred when checking memtable flush interval",
e);
+ }
+ }
+
@Override
public void stop() {
syncCloseAllProcessor();
@@ -301,6 +324,17 @@ public class StorageEngine implements IService {
"StorageEngine failed to stop because of " + "ttlCheckThread.", e);
}
}
+ if (memtableTimedFlushCheckThread != null) {
+ memtableTimedFlushCheckThread.shutdownNow();
+ try {
+ memtableTimedFlushCheckThread.awaitTermination(60, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ logger.warn("Memtable flush interval check thread still doesn't exit
after 60s");
+ Thread.currentThread().interrupt();
+ throw new StorageEngineFailureException(
+ "StorageEngine failed to stop because of
memtableFlushCheckThread.", e);
+ }
+ }
recoveryThreadPool.shutdownNow();
for (PartialPath storageGroup :
IoTDB.metaManager.getAllStorageGroupPaths()) {
this.releaseWalDirectByteBufferPoolInOneStorageGroup(storageGroup);
@@ -324,6 +358,15 @@ public class StorageEngine implements IService {
Thread.currentThread().interrupt();
}
}
+ if (memtableTimedFlushCheckThread != null) {
+ memtableTimedFlushCheckThread.shutdownNow();
+ try {
+ memtableTimedFlushCheckThread.awaitTermination(30, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ logger.warn("Memtable flush interval check thread still doesn't exit
after 30s");
+ Thread.currentThread().interrupt();
+ }
+ }
recoveryThreadPool.shutdownNow();
this.reset();
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
index 0acb9e1..569d2f4 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/AbstractMemTable.java
@@ -73,6 +73,8 @@ public abstract class AbstractMemTable implements IMemTable {
private long minPlanIndex = Long.MAX_VALUE;
+ private long createdTime = System.currentTimeMillis();
+
public AbstractMemTable() {
this.memTableMap = new HashMap<>();
}
@@ -436,4 +438,9 @@ public abstract class AbstractMemTable implements IMemTable
{
maxPlanIndex = Math.max(index, maxPlanIndex);
minPlanIndex = Math.min(index, minPlanIndex);
}
+
+ @Override
+ public long getCreatedTime() {
+ return createdTime;
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
index e5c9e32..ff865ea 100644
--- a/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
+++ b/server/src/main/java/org/apache/iotdb/db/engine/memtable/IMemTable.java
@@ -140,4 +140,6 @@ public interface IMemTable {
long getMaxPlanIndex();
long getMinPlanIndex();
+
+ long getCreatedTime();
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
index 7062f54..d124341 100755
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
@@ -1552,6 +1552,28 @@ public class StorageGroupProcessor {
}
}
+ public void timedFlushMemTable() {
+ writeLock("timedFlushMemTable");
+ try {
+ // only check unsequence tsfiles' memtables
+ List<TsFileProcessor> tsFileProcessors =
+ new ArrayList<>(workUnsequenceTsFileProcessors.values());
+ long timestampBaseline = System.currentTimeMillis() -
config.getUnseqMemtableFlushInterval();
+ for (TsFileProcessor tsFileProcessor : tsFileProcessors) {
+ if (tsFileProcessor.getWorkMemTableCreatedTime() < timestampBaseline) {
+ logger.info(
+ "Exceed flush interval, so flush work memtable of time partition
{} in storage group {}[{}]",
+ tsFileProcessor.getTimeRangeId(),
+ logicalStorageGroupName,
+ virtualStorageGroupId);
+ fileFlushPolicy.apply(this, tsFileProcessor,
tsFileProcessor.isSequence());
+ }
+ }
+ } finally {
+ writeUnlock();
+ }
+ }
+
/** This method will be blocked until all tsfile processors are closed. */
public void syncCloseAllWorkingTsFileProcessors() {
synchronized (closeStorageGroupCondition) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index ff19c89..b59c45c 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -1299,6 +1299,11 @@ public class TsFileProcessor {
return workMemTable != null ? workMemTable.getTVListsRamCost() : 0;
}
+ /** Return Long.MAX_VALUE if workMemTable is null */
+ public long getWorkMemTableCreatedTime() {
+ return workMemTable != null ? workMemTable.getCreatedTime() :
Long.MAX_VALUE;
+ }
+
public boolean isSequence() {
return sequence;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/virtualSg/VirtualStorageGroupManager.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/virtualSg/VirtualStorageGroupManager.java
index e760183..1f65389 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/virtualSg/VirtualStorageGroupManager.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/virtualSg/VirtualStorageGroupManager.java
@@ -105,6 +105,15 @@ public class VirtualStorageGroupManager {
}
}
+ /** push check memtable flush interval down to all sg */
+ public void timedFlushMemTable() {
+ for (StorageGroupProcessor storageGroupProcessor :
virtualStorageGroupProcessor) {
+ if (storageGroupProcessor != null) {
+ storageGroupProcessor.timedFlushMemTable();
+ }
+ }
+ }
+
/**
* get processor from device id
*
diff --git
a/server/src/main/java/org/apache/iotdb/db/rescon/MemTableManager.java
b/server/src/main/java/org/apache/iotdb/db/rescon/MemTableManager.java
index 856ffc8..461b64b 100644
--- a/server/src/main/java/org/apache/iotdb/db/rescon/MemTableManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/rescon/MemTableManager.java
@@ -102,6 +102,10 @@ public class MemTableManager {
notifyAll();
}
+ public synchronized void close() {
+ currentMemtableNumber = 0;
+ }
+
private static class InstanceHolder {
private static final MemTableManager INSTANCE = new MemTableManager();
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
index 3d7e2b6..0f0eca9 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.constant.TestConstant;
import org.apache.iotdb.db.engine.MetadataManagerHelper;
import org.apache.iotdb.db.engine.compaction.CompactionStrategy;
+import org.apache.iotdb.db.engine.flush.FlushManager;
import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy;
import org.apache.iotdb.db.engine.merge.manage.MergeManager;
import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
@@ -38,6 +39,7 @@ import org.apache.iotdb.db.metadata.mnode.MeasurementMNode;
import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
import org.apache.iotdb.db.qp.physical.crud.InsertTabletPlan;
import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.db.rescon.MemTableManager;
import org.apache.iotdb.db.utils.EnvironmentUtils;
import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
@@ -52,6 +54,8 @@ import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
@@ -59,6 +63,8 @@ import java.util.Collections;
import java.util.List;
public class StorageGroupProcessorTest {
+ private static final IoTDBConfig config =
IoTDBDescriptor.getInstance().getConfig();
+ private static Logger logger =
LoggerFactory.getLogger(StorageGroupProcessorTest.class);
private String storageGroup = "root.vehicle.d0";
private String systemDir = TestConstant.OUTPUT_DATA_DIR.concat("info");
@@ -67,11 +73,13 @@ public class StorageGroupProcessorTest {
private StorageGroupProcessor processor;
private QueryContext context = EnvironmentUtils.TEST_QUERY_CONTEXT;
+ private boolean prevEnableTimedFlushMemtable = false;
+
@Before
public void setUp() throws Exception {
- IoTDBDescriptor.getInstance()
- .getConfig()
- .setCompactionStrategy(CompactionStrategy.NO_COMPACTION);
+ prevEnableTimedFlushMemtable = config.isEnableTimedFlushUnseqMemtable();
+ config.setEnableTimedFlushUnseqMemtable(true);
+ config.setCompactionStrategy(CompactionStrategy.NO_COMPACTION);
MetadataManagerHelper.initMetadata();
EnvironmentUtils.envSetUp();
processor = new DummySGP(systemDir, storageGroup);
@@ -85,9 +93,8 @@ public class StorageGroupProcessorTest {
EnvironmentUtils.cleanDir(TestConstant.OUTPUT_DATA_DIR);
MergeManager.getINSTANCE().stop();
EnvironmentUtils.cleanEnv();
- IoTDBDescriptor.getInstance()
- .getConfig()
- .setCompactionStrategy(CompactionStrategy.LEVEL_COMPACTION);
+ config.setEnableTimedFlushUnseqMemtable(prevEnableTimedFlushMemtable);
+ config.setCompactionStrategy(CompactionStrategy.LEVEL_COMPACTION);
}
private void insertToStorageGroupProcessor(TSRecord record)
@@ -305,7 +312,6 @@ public class StorageGroupProcessorTest {
public void testEnableDiscardOutOfOrderDataForInsertRowPlan()
throws WriteProcessException, QueryProcessException,
IllegalPathException, IOException,
TriggerExecutionException {
- IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
boolean defaultValue = config.isEnableDiscardOutOfOrderData();
config.setEnableDiscardOutOfOrderData(true);
@@ -347,7 +353,6 @@ public class StorageGroupProcessorTest {
@Test
public void testEnableDiscardOutOfOrderDataForInsertTablet1()
throws QueryProcessException, IllegalPathException, IOException,
TriggerExecutionException {
- IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
boolean defaultEnableDiscard = config.isEnableDiscardOutOfOrderData();
long defaultTimePartition = config.getPartitionInterval();
boolean defaultEnablePartition = config.isEnablePartition();
@@ -429,7 +434,6 @@ public class StorageGroupProcessorTest {
@Test
public void testEnableDiscardOutOfOrderDataForInsertTablet2()
throws QueryProcessException, IllegalPathException, IOException,
TriggerExecutionException {
- IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
boolean defaultEnableDiscard = config.isEnableDiscardOutOfOrderData();
long defaultTimePartition = config.getPartitionInterval();
boolean defaultEnablePartition = config.isEnablePartition();
@@ -511,7 +515,6 @@ public class StorageGroupProcessorTest {
@Test
public void testEnableDiscardOutOfOrderDataForInsertTablet3()
throws QueryProcessException, IllegalPathException, IOException,
TriggerExecutionException {
- IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
boolean defaultEnableDiscard = config.isEnableDiscardOutOfOrderData();
long defaultTimePartition = config.getPartitionInterval();
boolean defaultEnablePartition = config.isEnablePartition();
@@ -610,7 +613,7 @@ public class StorageGroupProcessorTest {
}
processor.syncCloseAllWorkingTsFileProcessors();
-
processor.merge(IoTDBDescriptor.getInstance().getConfig().isForceFullMerge());
+ processor.merge(config.isForceFullMerge());
while (processor.getTsFileManagement().isUnseqMerging) {
// wait
}
@@ -626,6 +629,50 @@ public class StorageGroupProcessorTest {
}
}
+ @Test
+ public void testCheckMemTableFlushInterval()
+ throws IllegalPathException, InterruptedException, WriteProcessException,
+ TriggerExecutionException {
+ // create one seq memtable & close
+ TSRecord record = new TSRecord(10000, deviceId);
+ record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId,
String.valueOf(1000)));
+ processor.insert(new InsertRowPlan(record));
+ Assert.assertEquals(1,
MemTableManager.getInstance().getCurrentMemtableNumber());
+ processor.syncCloseAllWorkingTsFileProcessors();
+ Assert.assertEquals(0,
MemTableManager.getInstance().getCurrentMemtableNumber());
+
+ // create one unseq memtable
+ record = new TSRecord(1, deviceId);
+ record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId,
String.valueOf(1000)));
+ processor.insert(new InsertRowPlan(record));
+ Assert.assertEquals(1,
MemTableManager.getInstance().getCurrentMemtableNumber());
+
+ // check memtable's flush interval & flush the unsequence memtable
+ long preFLushInterval = config.getUnseqMemtableFlushInterval();
+ config.setUnseqMemtableFlushInterval(5);
+
+ Thread.sleep(500);
+
+ processor.timedFlushMemTable();
+
+ FlushManager flushManager = FlushManager.getInstance();
+ int waitCnt = 0;
+ while (flushManager.getNumberOfPendingTasks() != 0
+ || FlushManager.getInstance().getNumberOfPendingSubTasks() != 0
+ || flushManager.getNumberOfWorkingTasks() != 0
+ || FlushManager.getInstance().getNumberOfWorkingSubTasks() != 0) {
+ Thread.sleep(500);
+ ++waitCnt;
+ if (waitCnt % 10 == 0) {
+ logger.info("already wait {} s", waitCnt / 2);
+ }
+ }
+
+ Assert.assertEquals(0,
MemTableManager.getInstance().getCurrentMemtableNumber());
+
+ config.setUnseqMemtableFlushInterval(preFLushInterval);
+ }
+
class DummySGP extends StorageGroupProcessor {
DummySGP(String systemInfoDir, String storageGroupName) throws
StorageGroupProcessorException {
diff --git
a/server/src/test/java/org/apache/iotdb/db/utils/EnvironmentUtils.java
b/server/src/test/java/org/apache/iotdb/db/utils/EnvironmentUtils.java
index cc03ea3..f9c679e 100644
--- a/server/src/test/java/org/apache/iotdb/db/utils/EnvironmentUtils.java
+++ b/server/src/test/java/org/apache/iotdb/db/utils/EnvironmentUtils.java
@@ -39,6 +39,7 @@ import org.apache.iotdb.db.query.control.FileReaderManager;
import org.apache.iotdb.db.query.control.QueryResourceManager;
import org.apache.iotdb.db.query.control.TracingManager;
import org.apache.iotdb.db.query.udf.service.UDFRegistrationService;
+import org.apache.iotdb.db.rescon.MemTableManager;
import org.apache.iotdb.db.rescon.PrimitiveArrayManager;
import org.apache.iotdb.db.rescon.SystemInfo;
import org.apache.iotdb.db.service.IoTDB;
@@ -153,6 +154,9 @@ public class EnvironmentUtils {
// clear system info
SystemInfo.getInstance().close();
+ // clear memtable manager info
+ MemTableManager.getInstance().close();
+
// delete all directory
cleanAllDir();
config.setTsFileSizeThreshold(oldTsFileThreshold);