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