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 91db5fa  Compaction not block flush (#2341)
91db5fa is described below

commit 91db5fa35171d3992a1f77c65c70295c9f1e340b
Author: zhanglingzhe0820 <[email protected]>
AuthorDate: Wed Dec 30 08:51:51 2020 +0800

    Compaction not block flush (#2341)
---
 .../compaction/CompactionMergeTaskPoolManager.java |  32 +++-
 .../db/engine/compaction/TsFileManagement.java     |  22 ++-
 .../level/LevelCompactionTsFileManagement.java     | 190 +++++++++++----------
 .../no/NoCompactionTsFileManagement.java           |   5 -
 .../engine/compaction/utils/CompactionUtils.java   |  12 +-
 .../engine/storagegroup/StorageGroupProcessor.java |  90 +++++-----
 .../java/org/apache/iotdb/db/service/IoTDB.java    |   4 +-
 .../compaction/LevelCompactionMergeTest.java       |   4 +-
 .../LevelCompactionTsFileManagementTest.java       |   1 -
 .../NoCompactionTsFileManagementTest.java          |   1 -
 .../db/integration/IoTDBLevelCompactionIT.java     |   3 +
 .../iotdb/db/integration/IoTDBRestartIT.java       |  40 ++---
 .../db/integration/IoTDBRpcCompressionIT.java      |   1 -
 .../apache/iotdb/db/utils/EnvironmentUtils.java    |   3 +
 14 files changed, 232 insertions(+), 176 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
index aa2ef5f..893c667 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionMergeTaskPoolManager.java
@@ -19,21 +19,25 @@
 
 package org.apache.iotdb.db.engine.compaction;
 
+import static 
org.apache.iotdb.db.engine.compaction.utils.CompactionLogger.COMPACTION_LOG_NAME;
+
+import java.io.File;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.TimeUnit;
 import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.db.concurrent.ThreadName;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import 
org.apache.iotdb.db.engine.compaction.TsFileManagement.CompactionMergeTask;
 import org.apache.iotdb.db.service.IService;
 import org.apache.iotdb.db.service.ServiceType;
+import org.apache.iotdb.db.utils.FilePathUtils;
+import org.apache.iotdb.db.utils.TestOnly;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 /**
- * CompactionMergeTaskPoolManager provides a ThreadPool to queue and run all 
compaction
- * tasks.
+ * CompactionMergeTaskPoolManager provides a ThreadPool to queue and run all 
compaction tasks.
  */
 public class CompactionMergeTaskPoolManager implements IService {
 
@@ -75,6 +79,26 @@ public class CompactionMergeTaskPoolManager implements 
IService {
     }
   }
 
+  @TestOnly
+  public void waitAllCompactionFinish() {
+    if (pool != null) {
+      File sgDir = FSFactoryProducer.getFSFactory().getFile(
+          
FilePathUtils.regularizePath(IoTDBDescriptor.getInstance().getConfig().getSystemDir())
+              + "storage_groups");
+      File[] subDirList = sgDir.listFiles();
+      if(subDirList!=null) {
+        for (File subDir : subDirList) {
+          while (FSFactoryProducer.getFSFactory().getFile(
+              subDir.getAbsoluteFile() + File.separator + subDir.getName() + 
COMPACTION_LOG_NAME)
+              .exists()) {
+            // wait
+          }
+        }
+      }
+      logger.info("All compaction task finish");
+    }
+  }
+
   private void waitTermination() {
     long startTime = System.currentTimeMillis();
     while (!pool.isTerminated()) {
@@ -112,7 +136,7 @@ public class CompactionMergeTaskPoolManager implements 
IService {
     return ServiceType.COMPACTION_SERVICE;
   }
 
-  public void submitTask(CompactionMergeTask compactionMergeTask)
+  public void submitTask(Runnable compactionMergeTask)
       throws RejectedExecutionException {
     if (pool != null && !pool.isTerminated()) {
       pool.submit(compactionMergeTask);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
index d10b4c1..dc4df7f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
@@ -74,11 +74,6 @@ public abstract class TsFileManagement {
   }
 
   /**
-   * get the TsFile list which has been completed hot compacted
-   */
-  public abstract List<TsFileResource> getStableTsFileList(boolean sequence);
-
-  /**
    * get the TsFile list in sequence
    */
   public abstract List<TsFileResource> getTsFileList(boolean sequence);
@@ -178,7 +173,22 @@ public abstract class TsFileManagement {
     }
   }
 
-  public void merge(boolean fullMerge, List<TsFileResource> seqMergeList,
+  public class CompactionRecoverTask implements Runnable {
+
+    private CloseCompactionMergeCallBack closeCompactionMergeCallBack;
+
+    public CompactionRecoverTask(CloseCompactionMergeCallBack 
closeCompactionMergeCallBack) {
+      this.closeCompactionMergeCallBack = closeCompactionMergeCallBack;
+    }
+
+    @Override
+    public void run() {
+      recover();
+      closeCompactionMergeCallBack.call();
+    }
+  }
+
+  public synchronized void merge(boolean fullMerge, List<TsFileResource> 
seqMergeList,
       List<TsFileResource> unSeqMergeList, long dataTTL) {
     if (isUnseqMerging) {
       if (logger.isInfoEnabled()) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/LevelCompactionTsFileManagement.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/LevelCompactionTsFileManagement.java
index c4abfbd..ed89a14 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/LevelCompactionTsFileManagement.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/LevelCompactionTsFileManagement.java
@@ -31,11 +31,13 @@ import java.io.IOException;
 import java.nio.file.Files;
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Collections;
 import java.util.HashSet;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.SortedSet;
 import java.util.TreeSet;
 import java.util.concurrent.ConcurrentSkipListMap;
 import java.util.concurrent.CopyOnWriteArrayList;
@@ -61,24 +63,21 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   private static final Logger logger = LoggerFactory
       .getLogger(LevelCompactionTsFileManagement.class);
 
-  private final int seqLevelNum = 
IoTDBDescriptor.getInstance().getConfig().getSeqLevelNum() < 1 ? 1
-      : IoTDBDescriptor.getInstance().getConfig().getSeqLevelNum();
-  private final int seqFileNumInEachLevel =
-      IoTDBDescriptor.getInstance().getConfig().getSeqFileNumInEachLevel() < 1 
? 1
-          : 
IoTDBDescriptor.getInstance().getConfig().getSeqFileNumInEachLevel();
-  private final int unseqLevelNum =
-      IoTDBDescriptor.getInstance().getConfig().getUnseqLevelNum() < 1 ? 1
-          : IoTDBDescriptor.getInstance().getConfig().getUnseqLevelNum();
-  private final int unseqFileNumInEachLevel =
-      IoTDBDescriptor.getInstance().getConfig().getUnseqFileNumInEachLevel() < 
1 ? 1
-          : 
IoTDBDescriptor.getInstance().getConfig().getUnseqFileNumInEachLevel();
+  private final int seqLevelNum = Math
+      .max(IoTDBDescriptor.getInstance().getConfig().getSeqLevelNum(), 1);
+  private final int seqFileNumInEachLevel = Math
+      
.max(IoTDBDescriptor.getInstance().getConfig().getSeqFileNumInEachLevel(), 1);
+  private final int unseqLevelNum = Math
+      .max(IoTDBDescriptor.getInstance().getConfig().getUnseqLevelNum(), 1);
+  private final int unseqFileNumInEachLevel = Math
+      
.max(IoTDBDescriptor.getInstance().getConfig().getUnseqFileNumInEachLevel(), 1);
 
   private final boolean enableUnseqCompaction = 
IoTDBDescriptor.getInstance().getConfig()
       .isEnableUnseqCompaction();
   private final boolean isForceFullMerge = 
IoTDBDescriptor.getInstance().getConfig()
       .isForceFullMerge();
   // First map is partition list; Second list is level list; Third list is 
file list in level;
-  private final Map<Long, List<TreeSet<TsFileResource>>> 
sequenceTsFileResources = new ConcurrentSkipListMap<>();
+  private final Map<Long, List<SortedSet<TsFileResource>>> 
sequenceTsFileResources = new ConcurrentSkipListMap<>();
   private final Map<Long, List<List<TsFileResource>>> 
unSequenceTsFileResources = new ConcurrentSkipListMap<>();
   private final List<List<TsFileResource>> forkedSequenceTsFileResources = new 
ArrayList<>();
   private final List<List<TsFileResource>> forkedUnSequenceTsFileResources = 
new ArrayList<>();
@@ -103,13 +102,17 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     if (sequence) {
       if (sequenceTsFileResources.containsKey(timePartitionId)) {
         if (sequenceTsFileResources.get(timePartitionId).size() > level) {
-          
sequenceTsFileResources.get(timePartitionId).get(level).removeAll(mergeTsFiles);
+          synchronized (sequenceTsFileResources) {
+            
sequenceTsFileResources.get(timePartitionId).get(level).removeAll(mergeTsFiles);
+          }
         }
       }
     } else {
       if (unSequenceTsFileResources.containsKey(timePartitionId)) {
         if (unSequenceTsFileResources.get(timePartitionId).size() > level) {
-          
unSequenceTsFileResources.get(timePartitionId).get(level).removeAll(mergeTsFiles);
+          synchronized (unSequenceTsFileResources) {
+            
unSequenceTsFileResources.get(timePartitionId).get(level).removeAll(mergeTsFiles);
+          }
         }
       }
     }
@@ -130,33 +133,23 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   }
 
   @Override
-  public List<TsFileResource> getStableTsFileList(boolean sequence) {
-    List<TsFileResource> result = new ArrayList<>();
-    if (sequence) {
-      for (List<TreeSet<TsFileResource>> sequenceTsFileList : 
sequenceTsFileResources.values()) {
-        result.addAll(sequenceTsFileList.get(seqLevelNum - 1));
-      }
-    } else {
-      for (List<List<TsFileResource>> unSequenceTsFileList : 
unSequenceTsFileResources.values()) {
-        result.addAll(unSequenceTsFileList.get(unseqLevelNum - 1));
-      }
-    }
-    return result;
-  }
-
-  @Override
   public List<TsFileResource> getTsFileList(boolean sequence) {
     List<TsFileResource> result = new ArrayList<>();
     if (sequence) {
-      for (List<TreeSet<TsFileResource>> sequenceTsFileList : 
sequenceTsFileResources.values()) {
-        for (int i = sequenceTsFileList.size() - 1; i >= 0; i--) {
-          result.addAll(sequenceTsFileList.get(i));
+      synchronized (sequenceTsFileResources) {
+        for (List<SortedSet<TsFileResource>> sequenceTsFileList : 
sequenceTsFileResources
+            .values()) {
+          for (int i = sequenceTsFileList.size() - 1; i >= 0; i--) {
+            result.addAll(sequenceTsFileList.get(i));
+          }
         }
       }
     } else {
-      for (List<List<TsFileResource>> unSequenceTsFileList : 
unSequenceTsFileResources.values()) {
-        for (int i = unSequenceTsFileList.size() - 1; i >= 0; i--) {
-          result.addAll(unSequenceTsFileList.get(i));
+      synchronized (unSequenceTsFileResources) {
+        for (List<List<TsFileResource>> unSequenceTsFileList : 
unSequenceTsFileResources.values()) {
+          for (int i = unSequenceTsFileList.size() - 1; i >= 0; i--) {
+            result.addAll(unSequenceTsFileList.get(i));
+          }
         }
       }
     }
@@ -171,14 +164,18 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   @Override
   public void remove(TsFileResource tsFileResource, boolean sequence) {
     if (sequence) {
-      for (TreeSet<TsFileResource> sequenceTsFileResource : 
sequenceTsFileResources
-          .get(tsFileResource.getTimePartition())) {
-        sequenceTsFileResource.remove(tsFileResource);
+      synchronized (sequenceTsFileResources) {
+        for (SortedSet<TsFileResource> sequenceTsFileResource : 
sequenceTsFileResources
+            .get(tsFileResource.getTimePartition())) {
+          sequenceTsFileResource.remove(tsFileResource);
+        }
       }
     } else {
-      for (List<TsFileResource> unSequenceTsFileResource : 
unSequenceTsFileResources
-          .get(tsFileResource.getTimePartition())) {
-        unSequenceTsFileResource.remove(tsFileResource);
+      synchronized (unSequenceTsFileResources) {
+        for (List<TsFileResource> unSequenceTsFileResource : 
unSequenceTsFileResources
+            .get(tsFileResource.getTimePartition())) {
+          unSequenceTsFileResource.remove(tsFileResource);
+        }
       }
     }
   }
@@ -186,17 +183,21 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   @Override
   public void removeAll(List<TsFileResource> tsFileResourceList, boolean 
sequence) {
     if (sequence) {
-      for (List<TreeSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
-          .values()) {
-        for (TreeSet<TsFileResource> levelTsFileResource : 
partitionSequenceTsFileResource) {
-          levelTsFileResource.removeAll(tsFileResourceList);
+      synchronized (sequenceTsFileResources) {
+        for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
+            .values()) {
+          for (SortedSet<TsFileResource> levelTsFileResource : 
partitionSequenceTsFileResource) {
+            levelTsFileResource.removeAll(tsFileResourceList);
+          }
         }
       }
     } else {
-      for (List<List<TsFileResource>> partitionUnSequenceTsFileResource : 
unSequenceTsFileResources
-          .values()) {
-        for (List<TsFileResource> levelTsFileResource : 
partitionUnSequenceTsFileResource) {
-          levelTsFileResource.removeAll(tsFileResourceList);
+      synchronized (unSequenceTsFileResources) {
+        for (List<List<TsFileResource>> partitionUnSequenceTsFileResource : 
unSequenceTsFileResources
+            .values()) {
+          for (List<TsFileResource> levelTsFileResource : 
partitionUnSequenceTsFileResource) {
+            levelTsFileResource.removeAll(tsFileResourceList);
+          }
         }
       }
     }
@@ -207,28 +208,33 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     long timePartitionId = tsFileResource.getTimePartition();
     int level = getMergeLevel(tsFileResource.getTsFile());
     if (sequence) {
-      if (level <= seqLevelNum - 1) {
-        // current file has too high level
-        sequenceTsFileResources
-            .computeIfAbsent(timePartitionId, 
this::newSequenceTsFileResources).get(level)
-            .add(tsFileResource);
-      } else {
-        // current file has normal level
-        sequenceTsFileResources
-            .computeIfAbsent(timePartitionId, 
this::newSequenceTsFileResources).get(seqLevelNum - 1)
-            .add(tsFileResource);
+      synchronized (sequenceTsFileResources) {
+        if (level <= seqLevelNum - 1) {
+          // current file has normal level
+          sequenceTsFileResources
+              .computeIfAbsent(timePartitionId, 
this::newSequenceTsFileResources).get(level)
+              .add(tsFileResource);
+        } else {
+          // current file has too high level
+          sequenceTsFileResources
+              .computeIfAbsent(timePartitionId, 
this::newSequenceTsFileResources)
+              .get(seqLevelNum - 1)
+              .add(tsFileResource);
+        }
       }
     } else {
-      if (level <= unseqLevelNum - 1) {
-        // current file has too high level
-        unSequenceTsFileResources
-            .computeIfAbsent(timePartitionId, 
this::newUnSequenceTsFileResources).get(level)
-            .add(tsFileResource);
-      } else {
-        // current file has normal level
-        unSequenceTsFileResources
-            .computeIfAbsent(timePartitionId, 
this::newUnSequenceTsFileResources)
-            .get(unseqLevelNum - 1).add(tsFileResource);
+      synchronized (unSequenceTsFileResources) {
+        if (level <= unseqLevelNum - 1) {
+          // current file has normal level
+          unSequenceTsFileResources
+              .computeIfAbsent(timePartitionId, 
this::newUnSequenceTsFileResources).get(level)
+              .add(tsFileResource);
+        } else {
+          // current file has too high level
+          unSequenceTsFileResources
+              .computeIfAbsent(timePartitionId, 
this::newUnSequenceTsFileResources)
+              .get(unseqLevelNum - 1).add(tsFileResource);
+        }
       }
     }
   }
@@ -243,7 +249,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   @Override
   public boolean contains(TsFileResource tsFileResource, boolean sequence) {
     if (sequence) {
-      for (TreeSet<TsFileResource> sequenceTsFileResource : 
sequenceTsFileResources
+      for (SortedSet<TsFileResource> sequenceTsFileResource : 
sequenceTsFileResources
           .computeIfAbsent(tsFileResource.getTimePartition(), 
this::newSequenceTsFileResources)) {
         if (sequenceTsFileResource.contains(tsFileResource)) {
           return true;
@@ -270,9 +276,9 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   @SuppressWarnings("squid:S3776")
   public boolean isEmpty(boolean sequence) {
     if (sequence) {
-      for (List<TreeSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
+      for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
           .values()) {
-        for (TreeSet<TsFileResource> sequenceTsFileResource : 
partitionSequenceTsFileResource) {
+        for (SortedSet<TsFileResource> sequenceTsFileResource : 
partitionSequenceTsFileResource) {
           if (!sequenceTsFileResource.isEmpty()) {
             return false;
           }
@@ -295,7 +301,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   public int size(boolean sequence) {
     int result = 0;
     if (sequence) {
-      for (List<TreeSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
+      for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource : 
sequenceTsFileResources
           .values()) {
         for (int i = seqLevelNum - 1; i >= 0; i--) {
           result += partitionSequenceTsFileResource.get(i).size();
@@ -412,7 +418,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     if (isSeq) {
       for (int level = 0; level < 
sequenceTsFileResources.get(timePartition).size();
           level++) {
-        TreeSet<TsFileResource> currLevelMergeFile = sequenceTsFileResources
+        SortedSet<TsFileResource> currLevelMergeFile = sequenceTsFileResources
             .get(timePartition).get(level);
         deleteLevelFilesInDisk(currLevelMergeFile);
         deleteLevelFilesInList(timePartition, currLevelMergeFile, level, 
isSeq);
@@ -420,7 +426,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     } else {
       for (int level = 0; level < 
unSequenceTsFileResources.get(timePartition).size();
           level++) {
-        TreeSet<TsFileResource> currLevelMergeFile = sequenceTsFileResources
+        SortedSet<TsFileResource> currLevelMergeFile = sequenceTsFileResources
             .get(timePartition).get(level);
         deleteLevelFilesInDisk(currLevelMergeFile);
         deleteLevelFilesInList(timePartition, currLevelMergeFile, level, 
isSeq);
@@ -430,16 +436,20 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
 
   @Override
   public void forkCurrentFileList(long timePartition) {
-    forkTsFileList(
-        forkedSequenceTsFileResources,
-        sequenceTsFileResources.computeIfAbsent(timePartition, 
this::newSequenceTsFileResources),
-        seqLevelNum, seqFileNumInEachLevel);
+    synchronized (sequenceTsFileResources) {
+      forkTsFileList(
+          forkedSequenceTsFileResources,
+          sequenceTsFileResources.computeIfAbsent(timePartition, 
this::newSequenceTsFileResources),
+          seqLevelNum, seqFileNumInEachLevel);
+    }
     // we have to copy all unseq file
-    forkTsFileList(
-        forkedUnSequenceTsFileResources,
-        unSequenceTsFileResources
-            .computeIfAbsent(timePartition, 
this::newUnSequenceTsFileResources),
-        unseqLevelNum + 1, unseqFileNumInEachLevel);
+    synchronized (unSequenceTsFileResources) {
+      forkTsFileList(
+          forkedUnSequenceTsFileResources,
+          unSequenceTsFileResources
+              .computeIfAbsent(timePartition, 
this::newUnSequenceTsFileResources),
+          unseqLevelNum + 1, unseqFileNumInEachLevel);
+    }
   }
 
   private void forkTsFileList(
@@ -567,10 +577,10 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     return new File(prefixPath + level + TSFILE_SUFFIX);
   }
 
-  private List<TreeSet<TsFileResource>> newSequenceTsFileResources(Long k) {
-    List<TreeSet<TsFileResource>> newSequenceTsFileResources = new 
CopyOnWriteArrayList<>();
+  private List<SortedSet<TsFileResource>> newSequenceTsFileResources(Long k) {
+    List<SortedSet<TsFileResource>> newSequenceTsFileResources = new 
CopyOnWriteArrayList<>();
     for (int i = 0; i < seqLevelNum; i++) {
-      newSequenceTsFileResources.add(new TreeSet<>(
+      newSequenceTsFileResources.add(Collections.synchronizedSortedSet(new 
TreeSet<>(
           (o1, o2) -> {
             try {
               int rangeCompare = Long
@@ -581,7 +591,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
             } catch (NumberFormatException e) {
               return compareFileName(o1.getTsFile(), o2.getTsFile());
             }
-          }));
+          })));
     }
     return newSequenceTsFileResources;
   }
@@ -603,9 +613,9 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
 
   private TsFileResource getTsFileResource(String filePath, boolean isSeq) 
throws IOException {
     if (isSeq) {
-      for (List<TreeSet<TsFileResource>> tsFileResourcesWithLevel : 
sequenceTsFileResources
+      for (List<SortedSet<TsFileResource>> tsFileResourcesWithLevel : 
sequenceTsFileResources
           .values()) {
-        for (TreeSet<TsFileResource> tsFileResources : 
tsFileResourcesWithLevel) {
+        for (SortedSet<TsFileResource> tsFileResources : 
tsFileResourcesWithLevel) {
           for (TsFileResource tsFileResource : tsFileResources) {
             if (tsFileResource.getTsFile().getAbsolutePath().equals(filePath)) 
{
               return tsFileResource;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/no/NoCompactionTsFileManagement.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/no/NoCompactionTsFileManagement.java
index 4126557..2d76372 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/no/NoCompactionTsFileManagement.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/no/NoCompactionTsFileManagement.java
@@ -51,11 +51,6 @@ public class NoCompactionTsFileManagement extends 
TsFileManagement {
   }
 
   @Override
-  public List<TsFileResource> getStableTsFileList(boolean sequence) {
-    return getTsFileList(sequence);
-  }
-
-  @Override
   public List<TsFileResource> getTsFileList(boolean sequence) {
     if (sequence) {
       return new ArrayList<>(sequenceFileTreeSet);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/utils/CompactionUtils.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/utils/CompactionUtils.java
index 928548c..017b3aa 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/utils/CompactionUtils.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/utils/CompactionUtils.java
@@ -166,7 +166,9 @@ public class CompactionUtils {
       chunkWriter = new ChunkWriterImpl(
           IoTDB.metaManager.getSeriesSchema(new PartialPath(device), 
entry.getKey()), true);
     } catch (MetadataException e) {
-      throw new IOException(e);
+      // this may caused in IT by restart
+      logger.error("{} get schema {} error,skip this sensor", device, 
entry.getKey());
+      return maxVersion;
     }
     for (TimeValuePair timeValuePair : timeValuePairMap.values()) {
       writeTVPair(timeValuePair, chunkWriter);
@@ -197,11 +199,11 @@ public class CompactionUtils {
   }
 
   /**
-   * @param targetResource the target resource to be merged to
-   * @param tsFileResources the source resource to be merged
-   * @param storageGroup the storage group name
+   * @param targetResource   the target resource to be merged to
+   * @param tsFileResources  the source resource to be merged
+   * @param storageGroup     the storage group name
    * @param compactionLogger the logger
-   * @param devices the devices to be skipped(used by recover)
+   * @param devices          the devices to be skipped(used by recover)
    */
   @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity 
warning
   public static void merge(TsFileResource targetResource,
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 955deb4..a2d6a5f 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
@@ -328,7 +328,7 @@ public class StorageGroupProcessor {
       if 
(!IoTDBDescriptor.getInstance().getConfig().isContinueMergeAfterReboot()) {
         mergingMods.delete();
       }
-      tsFileManagement.recover();
+      recoverCompaction();
       updateLatestFlushedTime();
     } catch (IOException | MetadataException e) {
       throw new StorageGroupProcessorException(e);
@@ -353,6 +353,24 @@ public class StorageGroupProcessor {
 
   }
 
+  private void recoverCompaction() {
+    if (!CompactionMergeTaskPoolManager.getInstance().isTerminated()) {
+      compactionMergeWorking = true;
+      logger.info("{} submit a compaction merge task", storageGroupName);
+      try {
+        CompactionMergeTaskPoolManager.getInstance()
+            .submitTask(
+                tsFileManagement.new 
CompactionRecoverTask(this::closeCompactionMergeCallBack));
+      } catch (RejectedExecutionException e) {
+        this.closeCompactionMergeCallBack();
+        logger.error("{} compaction submit task failed", storageGroupName);
+      }
+    } else {
+      logger.error("{} compaction pool not started ,recover failed",
+          storageGroupName);
+    }
+  }
+
   public long getMonitorSeriesValue() {
     return monitorSeriesValue;
   }
@@ -570,8 +588,6 @@ public class StorageGroupProcessor {
         continue;
       }
 
-
-
       if (i != tsFiles.size() - 1 || !writer.canWrite()) {
         // not the last file or cannot write, just close it
         tsFileResource.setClosed(true);
@@ -798,11 +814,11 @@ public class StorageGroupProcessor {
    * inserted are in the range [start, end)
    *
    * @param insertTabletPlan insert a tablet of a device
-   * @param sequence         whether is sequence
-   * @param start            start index of rows to be inserted in 
insertTabletPlan
-   * @param end              end index of rows to be inserted in 
insertTabletPlan
-   * @param results          result array
-   * @param timePartitionId  time partition id
+   * @param sequence whether is sequence
+   * @param start start index of rows to be inserted in insertTabletPlan
+   * @param end end index of rows to be inserted in insertTabletPlan
+   * @param results result array
+   * @param timePartitionId time partition id
    * @return false if any failure occurs when inserting the tablet, true 
otherwise
    */
   private boolean insertTabletToTsFileProcessor(InsertTabletPlan 
insertTabletPlan,
@@ -870,7 +886,8 @@ public class StorageGroupProcessor {
     }
   }
 
-  private void insertToTsFileProcessor(InsertRowPlan insertRowPlan, boolean 
sequence, long timePartitionId)
+  private void insertToTsFileProcessor(InsertRowPlan insertRowPlan, boolean 
sequence,
+      long timePartitionId)
       throws WriteProcessException {
     TsFileProcessor tsFileProcessor = 
getOrCreateTsFileProcessor(timePartitionId, sequence);
 
@@ -924,7 +941,7 @@ public class StorageGroupProcessor {
   public void asyncFlushMemTableInTsFileProcessor(TsFileProcessor 
tsFileProcessor) {
     writeLock();
     try {
-      if (!closingSequenceTsFileProcessor.contains(tsFileProcessor) && 
+      if (!closingSequenceTsFileProcessor.contains(tsFileProcessor) &&
           !closingUnSequenceTsFileProcessor.contains(tsFileProcessor)) {
         fileFlushPolicy.apply(this, tsFileProcessor, 
tsFileProcessor.isSequence());
       }
@@ -960,9 +977,9 @@ public class StorageGroupProcessor {
   /**
    * get processor from hashmap, flush oldest processor if necessary
    *
-   * @param timeRangeId            time partition range
+   * @param timeRangeId time partition range
    * @param tsFileProcessorTreeMap tsFileProcessorTreeMap
-   * @param sequence               whether is sequence or not
+   * @param sequence whether is sequence or not
    */
   private TsFileProcessor getOrCreateTsFileProcessorIntern(long timeRangeId,
       TreeMap<Long, TsFileProcessor> tsFileProcessorTreeMap,
@@ -1098,7 +1115,7 @@ public class StorageGroupProcessor {
   public void asyncCloseOneTsFileProcessor(boolean sequence, TsFileProcessor 
tsFileProcessor) {
     //for sequence tsfile, we update the endTimeMap only when the file is 
prepared to be closed.
     //for unsequence tsfile, we have maintained the endTimeMap when an 
insertion comes.
-    if (closingSequenceTsFileProcessor.contains(tsFileProcessor) || 
+    if (closingSequenceTsFileProcessor.contains(tsFileProcessor) ||
         closingUnSequenceTsFileProcessor.contains(tsFileProcessor)) {
       return;
     }
@@ -1284,14 +1301,6 @@ public class StorageGroupProcessor {
                 (System.currentTimeMillis() - startTime) / 1000);
           }
         }
-        while (compactionMergeWorking) {
-          closeStorageGroupCondition.wait(100);
-          if (System.currentTimeMillis() - startTime > 60_000) {
-            logger
-                .warn("{} has spent {}s to wait for closing compaction.", 
this.storageGroupName,
-                    (System.currentTimeMillis() - startTime) / 1000);
-          }
-        }
       } catch (InterruptedException e) {
         logger.error("CloseFileNodeCondition error occurs while waiting for 
closing the storage "
             + "group {}", storageGroupName, e);
@@ -1474,10 +1483,9 @@ public class StorageGroupProcessor {
    * Delete data whose timestamp <= 'timestamp' and belongs to the time series
    * deviceId.measurementId.
    *
-   * @param path      the timeseries path of the to be deleted.
+   * @param path the timeseries path of the to be deleted.
    * @param startTime the startTime of delete range.
-   * @param endTime   the endTime of delete range.
-   * @param planIndex
+   * @param endTime the endTime of delete range.
    */
   public void delete(PartialPath path, long startTime, long endTime, long 
planIndex)
       throws IOException {
@@ -1694,6 +1702,11 @@ public class StorageGroupProcessor {
     } else {
       closingUnSequenceTsFileProcessor.remove(tsFileProcessor);
     }
+    synchronized (closeStorageGroupCondition) {
+      closeStorageGroupCondition.notifyAll();
+    }
+    logger.info("signal closing storage group condition in {}", 
storageGroupName);
+
     if (!compactionMergeWorking && 
!CompactionMergeTaskPoolManager.getInstance()
         .isTerminated()) {
       compactionMergeWorking = true;
@@ -1713,10 +1726,6 @@ public class StorageGroupProcessor {
       logger.info("{} last compaction merge task is working, skip current 
merge",
           storageGroupName);
     }
-    synchronized (closeStorageGroupCondition) {
-      closeStorageGroupCondition.notifyAll();
-    }
-    logger.info("signal closing storage group condition in {}", 
storageGroupName);
   }
 
   /**
@@ -1724,9 +1733,6 @@ public class StorageGroupProcessor {
    */
   private void closeCompactionMergeCallBack() {
     this.compactionMergeWorking = false;
-    synchronized (closeStorageGroupCondition) {
-      closeStorageGroupCondition.notifyAll();
-    }
   }
 
   /**
@@ -2057,7 +2063,8 @@ public class StorageGroupProcessor {
   private void removeFullyOverlapFile(TsFileResource tsFileResource,
       Iterator<TsFileResource> iterator
       , boolean isSeq) {
-    logger.info("Removing a covered file {}, closed: {}", tsFileResource, 
tsFileResource.isClosed());
+    logger
+        .info("Removing a covered file {}, closed: {}", tsFileResource, 
tsFileResource.isClosed());
     if (!tsFileResource.isClosed()) {
       try {
         // also remove the TsFileProcessor if the overlapped file is not closed
@@ -2102,9 +2109,9 @@ public class StorageGroupProcessor {
    * returns directly; otherwise, the time stamp is the mean of the timestamps 
of the two files, the
    * version number is the version number in the tsfile with a larger 
timestamp.
    *
-   * @param tsfileName  origin tsfile name
+   * @param tsfileName origin tsfile name
    * @param insertIndex the new file will be inserted between the files 
[insertIndex, insertIndex +
-   *                    1]
+   * 1]
    * @return appropriate filename
    */
   private String getFileNameForLoadingFile(String tsfileName, int insertIndex,
@@ -2172,8 +2179,8 @@ public class StorageGroupProcessor {
   /**
    * Execute the loading process by the type.
    *
-   * @param type            load type
-   * @param tsFileResource  tsfile resource to be loaded
+   * @param type load type
+   * @param tsFileResource tsfile resource to be loaded
    * @param filePartitionId the partition id of the new file
    * @return load the file successfully
    * @UsedBy sync module, load external tsfile module.
@@ -2397,10 +2404,13 @@ public class StorageGroupProcessor {
    */
   public boolean isFileAlreadyExist(TsFileResource tsFileResource, long 
partitionNum) {
     // examine working processor first as they have the largest plan index
-    return isFileAlreadyExistInWorking(tsFileResource, partitionNum, 
getWorkSequenceTsFileProcessors()) ||
-        isFileAlreadyExistInWorking(tsFileResource, partitionNum, 
getWorkUnsequenceTsFileProcessors()) ||
-        isFileAlreadyExistInClosed(tsFileResource, partitionNum, 
getSequenceFileTreeSet()) ||
-        isFileAlreadyExistInClosed(tsFileResource, partitionNum, 
getUnSequenceFileList());
+    return
+        isFileAlreadyExistInWorking(tsFileResource, partitionNum, 
getWorkSequenceTsFileProcessors())
+            ||
+            isFileAlreadyExistInWorking(tsFileResource, partitionNum,
+                getWorkUnsequenceTsFileProcessors()) ||
+            isFileAlreadyExistInClosed(tsFileResource, partitionNum, 
getSequenceFileTreeSet()) ||
+            isFileAlreadyExistInClosed(tsFileResource, partitionNum, 
getUnSequenceFileList());
   }
 
   private boolean isFileAlreadyExistInClosed(TsFileResource tsFileResource, 
long partitionNum,
diff --git a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java 
b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
index a4bdee5..8b03141 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
@@ -107,6 +107,8 @@ public class IoTDB implements IoTDBMBean {
     registerManager.register(Measurement.INSTANCE);
     registerManager.register(TVListAllocator.getInstance());
     registerManager.register(CacheHitRatioMonitor.getInstance());
+    registerManager.register(MergeManager.getINSTANCE());
+    registerManager.register(CompactionMergeTaskPoolManager.getInstance());
     JMXService.registerMBean(getInstance(), mbeanName);
     registerManager.register(StorageEngine.getInstance());
     registerManager.register(TemporaryQueryDataFileService.getInstance());
@@ -137,8 +139,6 @@ public class IoTDB implements IoTDBMBean {
     registerManager.register(StatMonitor.getInstance());
     registerManager.register(SyncServerManager.getInstance());
     registerManager.register(UpgradeSevice.getINSTANCE());
-    registerManager.register(MergeManager.getINSTANCE());
-    registerManager.register(CompactionMergeTaskPoolManager.getInstance());
 
     logger.info("Congratulation, IoTDB is set up successfully. Now, enjoy 
yourself!");
   }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
index 6ed1198..8970a4c 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionMergeTest.java
@@ -118,7 +118,7 @@ public class LevelCompactionMergeTest extends 
LevelCompactionTest {
         deviceIds[0] + TsFileConstant.PATH_SEPARATOR + 
measurementSchemas[0].getMeasurementId());
     IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path, 
measurementSchemas[0].getType(),
         context,
-        levelCompactionTsFileManagement.getStableTsFileList(true), new 
ArrayList<>(), null, null,
+        levelCompactionTsFileManagement.getTsFileList(true), new 
ArrayList<>(), null, null,
         true);
     int count = 0;
     while (tsFilesReader.hasNextBatch()) {
@@ -128,7 +128,7 @@ public class LevelCompactionMergeTest extends 
LevelCompactionTest {
         assertEquals(batchData.getTimeByIndex(i), 
batchData.getDoubleByIndex(i), 0.001);
       }
     }
-    assertEquals(200, count);
+    assertEquals(500, count);
     
IoTDBDescriptor.getInstance().getConfig().setSeqFileNumInEachLevel(prevSeqLevelFileNum);
     IoTDBDescriptor.getInstance().getConfig().setSeqLevelNum(prevSeqLevelNum);
   }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionTsFileManagementTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionTsFileManagementTest.java
index 0d8bdab..6b51613 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionTsFileManagementTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/LevelCompactionTsFileManagementTest.java
@@ -87,7 +87,6 @@ public class LevelCompactionTsFileManagementTest extends 
LevelCompactionTest {
     levelCompactionTsFileManagement
         .remove(levelCompactionTsFileManagement.getTsFileList(false).get(0), 
false);
     assertEquals(5, 
levelCompactionTsFileManagement.getTsFileList(true).size());
-    assertEquals(5, 
levelCompactionTsFileManagement.getStableTsFileList(false).size());
     levelCompactionTsFileManagement
         .removeAll(levelCompactionTsFileManagement.getTsFileList(false), 
false);
     assertEquals(0, 
levelCompactionTsFileManagement.getTsFileList(false).size());
diff --git 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/NoCompactionTsFileManagementTest.java
 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/NoCompactionTsFileManagementTest.java
index 68f417d..f99ec45 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/NoCompactionTsFileManagementTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/NoCompactionTsFileManagementTest.java
@@ -88,7 +88,6 @@ public class NoCompactionTsFileManagementTest extends 
LevelCompactionTest {
     noCompactionTsFileManagement
         .remove(noCompactionTsFileManagement.getTsFileList(false).get(0), 
false);
     assertEquals(5, noCompactionTsFileManagement.getTsFileList(true).size());
-    assertEquals(5, 
noCompactionTsFileManagement.getStableTsFileList(false).size());
     noCompactionTsFileManagement
         .removeAll(noCompactionTsFileManagement.getTsFileList(false), false);
     assertEquals(0, noCompactionTsFileManagement.getTsFileList(false).size());
diff --git 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBLevelCompactionIT.java
 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBLevelCompactionIT.java
index 761aa26..022f2cd 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBLevelCompactionIT.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBLevelCompactionIT.java
@@ -18,6 +18,7 @@
  */
 package org.apache.iotdb.db.integration;
 
+import static 
org.apache.iotdb.db.engine.compaction.utils.CompactionLogger.COMPACTION_LOG_NAME;
 import static org.junit.Assert.assertEquals;
 
 import java.sql.Connection;
@@ -29,7 +30,9 @@ import java.util.Random;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.compaction.CompactionStrategy;
 import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.db.utils.FilePathUtils;
 import org.apache.iotdb.jdbc.Config;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
 import org.junit.After;
 import org.junit.Before;
 import org.junit.Test;
diff --git 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
index 9f643c6..4a23e84 100644
--- a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
+++ b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
@@ -28,6 +28,7 @@ import java.sql.DriverManager;
 import java.sql.ResultSet;
 import java.sql.SQLException;
 import java.sql.Statement;
+import org.apache.iotdb.db.engine.compaction.CompactionMergeTaskPoolManager;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.utils.EnvironmentUtils;
 import org.apache.iotdb.jdbc.Config;
@@ -45,7 +46,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(1,1.0)");
       statement.execute("flush");
     }
@@ -59,7 +60,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(2,1.0)");
     }
 
@@ -72,7 +73,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(3,1.0)");
 
       boolean hasResultSet = statement.execute("SELECT s1 FROM 
root.turbine.d1");
@@ -104,7 +105,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(1,1)");
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(2,2)");
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(3,3)");
@@ -119,7 +120,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("delete from root.turbine.d1.s1 where time<=1");
 
       boolean hasResultSet = statement.execute("SELECT s1 FROM 
root.turbine.d1");
@@ -143,7 +144,7 @@ public class IoTDBRestartIT {
         hasResultSet = statement.execute("SELECT s1 FROM root.turbine.d1");
         assertTrue(hasResultSet);
         exp = new String[]{
-                "3,3.0"
+            "3,3.0"
         };
         resultSet = statement.getResultSet();
         cnt = 0;
@@ -169,7 +170,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(1,1)");
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(2,2)");
     }
@@ -228,7 +229,7 @@ public class IoTDBRestartIT {
     try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(1,1)");
       statement.execute("insert into root.turbine.d1(timestamp,s1) 
values(2,2)");
     }
@@ -284,22 +285,22 @@ public class IoTDBRestartIT {
     EnvironmentUtils.envSetUp();
     Class.forName(Config.JDBC_DRIVER_NAME);
 
-    try(Connection connection = DriverManager
+    try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine1.d1(timestamp,s1,s2) 
values(1,1.1,2.2)");
       statement.execute("delete timeseries root.turbine1.d1.s1");
-      statement.execute("create timeseries root.turbine1.d1.s1 with 
datatype=INT32, encoding=RLE, compression=SNAPPY");
+      statement.execute(
+          "create timeseries root.turbine1.d1.s1 with datatype=INT32, 
encoding=RLE, compression=SNAPPY");
     }
 
     EnvironmentUtils.restartDaemon();
 
-
-    try(Connection connection = DriverManager
+    try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
 
       boolean hasResultSet = statement.execute("select * from root");
       assertTrue(hasResultSet);
@@ -319,21 +320,20 @@ public class IoTDBRestartIT {
     EnvironmentUtils.envSetUp();
     Class.forName(Config.JDBC_DRIVER_NAME);
 
-    try(Connection connection = DriverManager
+    try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
       statement.execute("insert into root.turbine1.d1(timestamp,s1,s2) 
values(1,1.1,2.2)");
       statement.execute("delete timeseries root.turbine1.d1.s1");
     }
 
     EnvironmentUtils.restartDaemon();
 
-
-    try(Connection connection = DriverManager
+    try (Connection connection = DriverManager
         .getConnection(Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root",
             "root");
-        Statement statement = connection.createStatement()){
+        Statement statement = connection.createStatement()) {
 
       boolean hasResultSet = statement.execute("select * from root");
       assertTrue(hasResultSet);
@@ -375,6 +375,8 @@ public class IoTDBRestartIT {
     }
 
     try {
+      CompactionMergeTaskPoolManager.getInstance().waitAllCompactionFinish();
+      Thread.sleep(10000);
       EnvironmentUtils.restartDaemon();
     } catch (Exception e) {
       Assert.fail();
diff --git 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRpcCompressionIT.java
 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRpcCompressionIT.java
index 615d35f..28a279f 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRpcCompressionIT.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRpcCompressionIT.java
@@ -135,7 +135,6 @@ public class IoTDBRpcCompressionIT {
         }
         assertEquals(1, cnt);
       }
-      statement.execute("merge");
       Thread.sleep(1000);
       // before merge completes
       try (ResultSet set = statement.executeQuery("SELECT * FROM root")) {
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 2280c78..5551376 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.engine.StorageEngine;
 import org.apache.iotdb.db.engine.cache.ChunkCache;
 import org.apache.iotdb.db.engine.cache.ChunkMetadataCache;
 import org.apache.iotdb.db.engine.cache.TimeSeriesMetadataCache;
+import org.apache.iotdb.db.engine.compaction.CompactionMergeTaskPoolManager;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.query.control.FileReaderManager;
@@ -79,6 +80,8 @@ public class EnvironmentUtils {
       .parseBoolean(System.getProperty("test.port.closed", "false"));
 
   public static void cleanEnv() throws IOException, StorageEngineException {
+    // wait all compaction finished
+    CompactionMergeTaskPoolManager.getInstance().waitAllCompactionFinish();
     // deregister all UDFs
     UDFRegistrationService.getInstance().deregisterAll();
 

Reply via email to