This is an automated email from the ASF dual-hosted git repository.
ejttianyu pushed a commit to branch proceeding_vldb
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/proceeding_vldb by this push:
new 85ea06e finish tired compaction strategy
85ea06e is described below
commit 85ea06e4d7b71da741f4e4e98f9046455fef0ebe
Author: EJTTianyu <[email protected]>
AuthorDate: Tue Feb 9 12:01:44 2021 +0800
finish tired compaction strategy
---
.../resources/conf/iotdb-engine.properties | 3 +
.../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 13 +
.../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 3 +
.../db/engine/compaction/TsFileManagement.java | 14 +-
.../level/LevelCompactionTsFileManagement.java | 7 +-
.../tired/TiredCompactionTsFileManagement.java | 420 +++++++++++++++++++--
.../engine/storagegroup/TiredCompactionTest.java | 97 +++++
.../engine/storagegroup/TiredCompactionTest1.java | 98 +++++
8 files changed, 611 insertions(+), 44 deletions(-)
diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 1cb7f23..de43030 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -316,6 +316,9 @@ size_ratio=2
# The max num of level.
level_num=4
+# tired compaction size,用于 size tired 合并,指定一次合并的文件数量
+tired_file_num=4
+
# During a merge, if a chunk with less number of points than this parameter,
the chunk will be
# merged with its succeeding chunks even if it is not overflowed, until the
merged chunks reach
# this threshold and the new chunk will be flushed.
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index e3115af..1a08b9a 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -349,6 +349,14 @@ public class IoTDBConfig {
this.levelNum = levelNum;
}
+ public int getTiredFileNum() {
+ return tiredFileNum;
+ }
+
+ public void setTiredFileNum(int tiredFileNum) {
+ this.tiredFileNum = tiredFileNum;
+ }
+
/**
* 第一层的数据文件最大量
*/
@@ -385,6 +393,11 @@ public class IoTDBConfig {
private int seqLevelNum = 3;
/**
+ * tired compaction size,用于 size tired 合并,指定一次合并的文件数量
+ */
+ private int tiredFileNum = 4;
+
+ /**
* Works when compaction_strategy is LEVEL_COMPACTION.
* The max ujseq file num of each level.
* When the num of files in one level exceeds this,
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 13c796a..3c6e703 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -335,6 +335,9 @@ public class IoTDBDescriptor {
.getProperty("seq_level_num",
Integer.toString(conf.getSeqLevelNum()))));
+ conf.setTiredFileNum(Integer.parseInt(properties.
+ getProperty("tired_file_num",
Integer.toString(conf.getTiredFileNum()))));
+
conf.setSeqFileNumInEachLevel(Integer.parseInt(properties
.getProperty("seq_file_num_in_each_level",
Integer.toString(conf.getSeqFileNumInEachLevel()))));
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 278051c..ed84cae 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
@@ -40,6 +40,7 @@ import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
+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;
@@ -55,6 +56,7 @@ import
org.apache.iotdb.db.engine.modification.ModificationFile;
import
org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor.CloseCompactionMergeCallBack;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
import org.apache.iotdb.db.exception.MergeException;
+import org.apache.iotdb.db.metadata.PartialPath;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -193,7 +195,7 @@ public abstract class TsFileManagement {
return newUnSequenceTsFileResources;
}
- private void forkTsFileList(
+ protected void forkTsFileList(
List<List<TsFileResource>> forkedTsFileResources,
List rawTsFileResources, int currMaxLevel) {
forkedTsFileResources.clear();
@@ -493,4 +495,14 @@ public abstract class TsFileManagement {
return cmp;
}
}
+
+ protected boolean isCompactionWorking() {
+ try {
+ return StorageEngine.getInstance().getProcessor(new
PartialPath(storageGroupName))
+ .isCompactionMergeWorking();
+ } catch (Exception e) {
+ //TODO do nothing
+ }
+ return false;
+ }
}
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 8ccc109..a4ac9b1 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
@@ -505,7 +505,7 @@ public class LevelCompactionTsFileManagement extends
TsFileManagement {
long seqEndTime = seqFile.getEndTime(deviceId);
if (!(unseqEndTime < seqStartTime || unseqStartTime > seqEndTime)) {
long level = (long) getMergeLevel(seqFile.getTsFile()) + 1;
- if(files.containsKey(level) && files.get(level).contains(seqFile)) {
+ if (files.containsKey(level) && files.get(level).contains(seqFile)) {
continue;
} else {
files.computeIfAbsent((long) getMergeLevel(seqFile.getTsFile()) +
1,
@@ -662,7 +662,7 @@ public class LevelCompactionTsFileManagement extends
TsFileManagement {
sequenceTsFileResources.get(timePartition).get((int) (res.getKey() -
1))
.removeAll(res.getValue());
}
- for (TsFileResource deleteRes : res.getValue()){
+ for (TsFileResource deleteRes : res.getValue()) {
deleteRes.delete();
}
}
@@ -705,6 +705,9 @@ public class LevelCompactionTsFileManagement extends
TsFileManagement {
@Override
protected void merge(long timePartition) {
+ if (isCompactionWorking()) {
+ return;
+ }
handleSpecificCase(timePartition);
if (processUnseq()) {
Map<Long, Map<Long, List<TsFileResource>>> selectFiles =
selectMergeFile(timePartition);
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/tired/TiredCompactionTsFileManagement.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/tired/TiredCompactionTsFileManagement.java
index e41b6d2..f2f688b 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/tired/TiredCompactionTsFileManagement.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/tired/TiredCompactionTsFileManagement.java
@@ -19,46 +19,125 @@
package org.apache.iotdb.db.engine.compaction.tired;
+import static org.apache.iotdb.db.conf.IoTDBConstant.FILE_NAME_SEPARATOR;
+import static org.apache.iotdb.db.utils.MergeUtils.writeBatchPoint;
+import static
org.apache.iotdb.tsfile.common.constant.TsFileConstant.TSFILE_SUFFIX;
+
+import java.io.File;
+import java.io.IOException;
import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
+import java.util.Map.Entry;
+import java.util.Set;
+import java.util.SortedSet;
import java.util.TreeSet;
+import java.util.concurrent.CopyOnWriteArrayList;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.compaction.TsFileManagement;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.metadata.MManager;
+import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.metadata.mnode.MNode;
+import org.apache.iotdb.db.metadata.mnode.MeasurementMNode;
+import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.db.query.reader.series.SeriesRawDataBatchReader;
+import org.apache.iotdb.db.utils.MergeUtils;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
+import org.apache.iotdb.tsfile.read.common.BatchData;
+import org.apache.iotdb.tsfile.read.reader.IBatchReader;
+import org.apache.iotdb.tsfile.utils.Pair;
+import org.apache.iotdb.tsfile.write.chunk.ChunkWriterImpl;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
+import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class TiredCompactionTsFileManagement extends TsFileManagement {
+ public static final String MERGE_SUFFIX = ".merge";
+
private static final Logger logger = LoggerFactory.getLogger(
TiredCompactionTsFileManagement.class);
- // includes sealed and unsealed sequence TsFiles
- private TreeSet<TsFileResource> sequenceFileTreeSet = new TreeSet<>(
- (o1, o2) -> {
- try {
- int rangeCompare =
Long.compare(Long.parseLong(o1.getTsFile().getParentFile().getName()),
- Long.parseLong(o2.getTsFile().getParentFile().getName()));
- return rangeCompare == 0 ? compareFileName(o1.getTsFile(),
o2.getTsFile()) : rangeCompare;
- } catch (NumberFormatException e) {
- return compareFileName(o1.getTsFile(), o2.getTsFile());
- }
- });
-
- // includes sealed and unsealed unSequence TsFiles
- private List<TsFileResource> unSequenceFileList = new ArrayList<>();
+ private final int unseqLevelNum = Math
+ .max(IoTDBDescriptor.getInstance().getConfig().getUnseqLevelNum(), 1);
public TiredCompactionTsFileManagement(String storageGroupName, String
storageGroupDir) {
super(storageGroupName, storageGroupDir);
}
@Override
+ public void forkCurrentFileList(long timePartition) throws IOException {
+ synchronized (sequenceTsFileResources) {
+ forkTsFileList(
+ forkedSequenceTsFileResources,
+ sequenceTsFileResources.computeIfAbsent(timePartition,
this::newSequenceTsFileResources),
+ config.getLevelNum());
+ }
+ synchronized (unSequenceTsFileResources) {
+ forkTsFileList(
+ forkedUnSequenceTsFileResources,
+ unSequenceTsFileResources
+ .computeIfAbsent(timePartition,
this::newUnSequenceTsFileResources),
+ config.getLevelNum());
+ }
+ }
+
+ @Override
+ protected List<SortedSet<TsFileResource>> newSequenceTsFileResources(Long k)
{
+ List<SortedSet<TsFileResource>> newSequenceTsFileResources = new
CopyOnWriteArrayList<>();
+ for (int i = 0; i < config.getLevelNum(); i++) {
+ newSequenceTsFileResources.add(Collections.synchronizedSortedSet(new
TreeSet<>(
+ (o1, o2) -> {
+ try {
+ int rangeCompare = Long
+
.compare(Long.parseLong(o1.getTsFile().getParentFile().getName()),
+
Long.parseLong(o2.getTsFile().getParentFile().getName()));
+ return rangeCompare == 0 ? compareFileName(o1.getTsFile(),
o2.getTsFile())
+ : rangeCompare;
+ } catch (NumberFormatException e) {
+ return compareFileName(o1.getTsFile(), o2.getTsFile());
+ }
+ })));
+ }
+ return newSequenceTsFileResources;
+ }
+
+ @Override
+ protected List<List<TsFileResource>> newUnSequenceTsFileResources(Long k) {
+ List<List<TsFileResource>> newUnSequenceTsFileResources = new
CopyOnWriteArrayList<>();
+ for (int i = 0; i < config.getLevelNum(); i++) {
+ newUnSequenceTsFileResources.add(new CopyOnWriteArrayList<>());
+ }
+ return newUnSequenceTsFileResources;
+ }
+
+ @Override
public List<TsFileResource> getTsFileList(boolean sequence) {
+ List<TsFileResource> result = new ArrayList<>();
if (sequence) {
- return new ArrayList<>(sequenceFileTreeSet);
+ synchronized (sequenceTsFileResources) {
+ for (List<SortedSet<TsFileResource>> sequenceTsFileList :
sequenceTsFileResources
+ .values()) {
+ for (int i = sequenceTsFileList.size() - 1; i >= 0; i--) {
+ result.addAll(sequenceTsFileList.get(i));
+ }
+ }
+ }
} else {
- return unSequenceFileList;
+ synchronized (unSequenceTsFileResources) {
+ for (List<List<TsFileResource>> unSequenceTsFileList :
unSequenceTsFileResources.values()) {
+ for (int i = unSequenceTsFileList.size() - 1; i >= 0; i--) {
+ result.addAll(unSequenceTsFileList.get(i));
+ }
+ }
+ }
}
+ return result;
}
@Override
@@ -69,70 +148,165 @@ public class TiredCompactionTsFileManagement extends
TsFileManagement {
@Override
public void remove(TsFileResource tsFileResource, boolean sequence) {
if (sequence) {
- sequenceFileTreeSet.remove(tsFileResource);
+ synchronized (sequenceTsFileResources) {
+ for (SortedSet<TsFileResource> sequenceTsFileResource :
sequenceTsFileResources
+ .get(tsFileResource.getTimePartition())) {
+ sequenceTsFileResource.remove(tsFileResource);
+ }
+ }
} else {
- unSequenceFileList.remove(tsFileResource);
+ synchronized (unSequenceTsFileResources) {
+ for (List<TsFileResource> unSequenceTsFileResource :
unSequenceTsFileResources
+ .get(tsFileResource.getTimePartition())) {
+ unSequenceTsFileResource.remove(tsFileResource);
+ }
+ }
}
}
@Override
public void removeAll(List<TsFileResource> tsFileResourceList, boolean
sequence) {
if (sequence) {
- sequenceFileTreeSet.removeAll(tsFileResourceList);
+ synchronized (sequenceTsFileResources) {
+ for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource :
sequenceTsFileResources
+ .values()) {
+ for (SortedSet<TsFileResource> levelTsFileResource :
partitionSequenceTsFileResource) {
+ levelTsFileResource.removeAll(tsFileResourceList);
+ }
+ }
+ }
} else {
- unSequenceFileList.removeAll(tsFileResourceList);
+ synchronized (unSequenceTsFileResources) {
+ for (List<List<TsFileResource>> partitionUnSequenceTsFileResource :
unSequenceTsFileResources
+ .values()) {
+ for (List<TsFileResource> levelTsFileResource :
partitionUnSequenceTsFileResource) {
+ levelTsFileResource.removeAll(tsFileResourceList);
+ }
+ }
+ }
}
}
+ public static int getMergeLevel(File file) {
+ String mergeLevelStr = file.getPath()
+ .substring(file.getPath().lastIndexOf(FILE_NAME_SEPARATOR) + 1)
+ .replaceAll(TSFILE_SUFFIX, "");
+ return Integer.parseInt(mergeLevelStr);
+ }
+
@Override
public void add(TsFileResource tsFileResource, boolean sequence) {
+ long timePartitionId = tsFileResource.getTimePartition();
+ int level = getMergeLevel(tsFileResource.getTsFile());
if (sequence) {
- sequenceFileTreeSet.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 {
- unSequenceFileList.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);
+ }
+ }
}
}
@Override
public void addAll(List<TsFileResource> tsFileResourceList, boolean
sequence) {
- if (sequence) {
- sequenceFileTreeSet.addAll(tsFileResourceList);
- } else {
- unSequenceFileList.addAll(tsFileResourceList);
+ for (TsFileResource tsFileResource : tsFileResourceList) {
+ add(tsFileResource, sequence);
}
}
@Override
public boolean contains(TsFileResource tsFileResource, boolean sequence) {
if (sequence) {
- return sequenceFileTreeSet.contains(tsFileResource);
+ for (SortedSet<TsFileResource> sequenceTsFileResource :
sequenceTsFileResources
+ .computeIfAbsent(tsFileResource.getTimePartition(),
this::newSequenceTsFileResources)) {
+ if (sequenceTsFileResource.contains(tsFileResource)) {
+ return true;
+ }
+ }
} else {
- return unSequenceFileList.contains(tsFileResource);
+ for (List<TsFileResource> unSequenceTsFileResource :
unSequenceTsFileResources
+ .computeIfAbsent(tsFileResource.getTimePartition(),
this::newUnSequenceTsFileResources)) {
+ if (unSequenceTsFileResource.contains(tsFileResource)) {
+ return true;
+ }
+ }
}
+ return false;
}
@Override
public void clear() {
- sequenceFileTreeSet.clear();
- unSequenceFileList.clear();
+ sequenceTsFileResources.clear();
+ unSequenceTsFileResources.clear();
}
@Override
+ @SuppressWarnings("squid:S3776")
public boolean isEmpty(boolean sequence) {
if (sequence) {
- return sequenceFileTreeSet.isEmpty();
+ for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource :
sequenceTsFileResources
+ .values()) {
+ for (SortedSet<TsFileResource> sequenceTsFileResource :
partitionSequenceTsFileResource) {
+ if (!sequenceTsFileResource.isEmpty()) {
+ return false;
+ }
+ }
+ }
} else {
- return unSequenceFileList.isEmpty();
+ for (List<List<TsFileResource>> partitionUnSequenceTsFileResource :
unSequenceTsFileResources
+ .values()) {
+ for (List<TsFileResource> unSequenceTsFileResource :
partitionUnSequenceTsFileResource) {
+ if (!unSequenceTsFileResource.isEmpty()) {
+ return false;
+ }
+ }
+ }
}
+ return true;
}
@Override
public int size(boolean sequence) {
+ int result = 0;
if (sequence) {
- return sequenceFileTreeSet.size();
+ for (List<SortedSet<TsFileResource>> partitionSequenceTsFileResource :
sequenceTsFileResources
+ .values()) {
+ for (int i = seqLevelNum - 1; i >= 0; i--) {
+ result += partitionSequenceTsFileResource.get(i).size();
+ }
+ }
} else {
- return unSequenceFileList.size();
+ for (List<List<TsFileResource>> partitionUnSequenceTsFileResource :
unSequenceTsFileResources
+ .values()) {
+ for (int i = unseqLevelNum - 1; i >= 0; i--) {
+ result += partitionUnSequenceTsFileResource.get(i).size();
+ }
+ }
}
+ return result;
}
@Override
@@ -141,22 +315,186 @@ public class TiredCompactionTsFileManagement extends
TsFileManagement {
}
@Override
- public void forkCurrentFileList(long timePartition) {
- logger.info("{} do not need fork", storageGroupName);
- }
-
- @Override
protected Map<Long, Map<Long, List<TsFileResource>>> selectMergeFile(long
timePartition) {
- return null;
+ Map<Long, Map<Long, List<TsFileResource>>> mergedFilesRet = new
HashMap<>();
+ Map<Long, List<TsFileResource>> selectFiles = new HashMap<>();
+ boolean isMergedFileFound = false;
+ // judge seq and unseq independently
+ for (int i = 0; i < config.getLevelNum() - 1; i++) {
+ if (config.getTiredFileNum() <=
forkedSequenceTsFileResources.get(i).size()) {
+ isMergedFileFound = true;
+ List<TsFileResource> mergedRes = forkedSequenceTsFileResources.get(i)
+ .subList(0, config.getTiredFileNum());
+ selectFiles.put((long) i, mergedRes);
+ break;
+ }
+ }
+ for (int i = 0; i < config.getLevelNum() - 1 && !isMergedFileFound; i++) {
+ if (config.getTiredFileNum() <=
forkedUnSequenceTsFileResources.get(i).size()) {
+ List<TsFileResource> mergedRes = forkedUnSequenceTsFileResources.get(i)
+ .subList(0, config.getTiredFileNum());
+ selectFiles.put((long) i, mergedRes);
+ break;
+ }
+ }
+ mergedFilesRet.put(timePartition, selectFiles);
+ return mergedFilesRet;
}
@Override
protected void merge(long timePartition) {
- logger.info("{} no merge logic", storageGroupName);
+ if (isCompactionWorking()) {
+ return;
+ }
+ Map<Long, Map<Long, List<TsFileResource>>> selectFiles =
selectMergeFile(timePartition);
+ mergeFiles(selectFiles, timePartition);
}
@Override
protected void mergeFiles(Map<Long, Map<Long, List<TsFileResource>>>
resources,
long timePartition) {
+ Map<Long, List<TsFileResource>> mergeResources =
resources.get(timePartition);
+ List<String> fileNames = new ArrayList<>();
+ long mergedLevel = 0;
+ String parentPath = "";
+ List<TsFileResource> seqFiles = new ArrayList<>();
+ List<TsFileResource> unseqFiles = new ArrayList<>();
+ // 获取最大的 level
+ for (Entry<Long, List<TsFileResource>> resource :
mergeResources.entrySet()) {
+ mergedLevel = resource.getKey();
+ if (parentPath.equals("")) {
+ parentPath = resource.getValue().get(0).getTsFile().getParent();
+ }
+ if (parentPath.contains("unsequence")) {
+ unseqFiles.addAll(resource.getValue());
+ } else {
+ seqFiles.addAll(resource.getValue());
+ }
+ for (TsFileResource res : resource.getValue()) {
+ fileNames.add(res.getTsFile().getName());
+ }
+ }
+ if (seqFiles.isEmpty() && unseqFiles.isEmpty()){
+ return;
+ }
+ Collections.sort(fileNames);
+ try {
+ // get historical versions
+ Set<Long> historicalVersions = new HashSet<>();
+ for (TsFileResource tsFileResource : seqFiles) {
+ historicalVersions.addAll(tsFileResource.getHistoricalVersions());
+ }
+ for (TsFileResource tsFileResource : unseqFiles) {
+ historicalVersions.addAll(tsFileResource.getHistoricalVersions());
+ }
+
+ Set<PartialPath> devices = MManager.getInstance()
+ .getDevices(new PartialPath(storageGroupName));
+ Map<PartialPath, ChunkWriterImpl> chunkWriterCacheMap = new HashMap<>();
+ for (PartialPath device : devices) {
+ MNode deviceNode = MManager.getInstance().getNodeByPath(device);
+ for (Entry<String, MNode> entry : deviceNode.getChildren().entrySet())
{
+ MeasurementSchema measurementSchema = ((MeasurementMNode)
entry.getValue()).getSchema();
+ chunkWriterCacheMap
+ .put(new PartialPath(device.toString(), entry.getKey()),
+ new ChunkWriterImpl(measurementSchema));
+ }
+ }
+ List<PartialPath> unmergedSeries =
+ MManager.getInstance().getAllTimeseriesPath(new
PartialPath(storageGroupName));
+ Pair<RestorableTsFileIOWriter, TsFileResource> newTsFilePair =
createNewFileWriter(
+ MERGE_SUFFIX, parentPath, fileNames, mergedLevel + 1);
+ RestorableTsFileIOWriter newFileWriter = newTsFilePair.left;
+ TsFileResource newResource = newTsFilePair.right;
+
+ List<List<PartialPath>> devicePaths =
MergeUtils.splitPathsByDevice(unmergedSeries);
+ for (List<PartialPath> pathList : devicePaths) {
+ String device = pathList.get(0).getDevice();
+ newFileWriter.startChunkGroup(device);
+
+ for (PartialPath path : pathList) {
+ long currMinTime = Long.MAX_VALUE;
+ long currMaxTime = Long.MIN_VALUE;
+ ChunkWriterImpl chunkWriter = chunkWriterCacheMap.get(path);
+ newFileWriter.addSchema(path, chunkWriter.getMeasurementSchema());
+ QueryContext context = new QueryContext();
+ IBatchReader tsFilesReader = new SeriesRawDataBatchReader(path,
+ chunkWriter.getMeasurementSchema().getType(),
+ context, seqFiles, unseqFiles, null, null, true);
+ while (tsFilesReader.hasNextBatch()) {
+ BatchData batchData = tsFilesReader.nextBatch();
+ currMinTime = Math.min(currMinTime, batchData.getTimeByIndex(0));
+ for (int i = 0; i < batchData.length(); i++) {
+ writeBatchPoint(batchData, i, chunkWriter);
+ }
+ if (!tsFilesReader.hasNextBatch()) {
+ currMaxTime =
Math.max(batchData.getTimeByIndex(batchData.length() - 1), currMaxTime);
+ }
+ }
+ synchronized (newFileWriter) {
+ chunkWriter.writeToFileWriter(newFileWriter);
+ }
+ newResource.updateStartTime(path.getDevice(), currMinTime);
+ newResource.updateEndTime(path.getDevice(), currMaxTime);
+ tsFilesReader.close();
+ }
+ newFileWriter.writeVersion(0L);
+ newFileWriter.endChunkGroup();
+ }
+ newResource.setHistoricalVersions(historicalVersions);
+ newResource.serialize();
+ newFileWriter.endFile();
+
+ cleanUp(resources, newResource, mergedLevel, timePartition, parentPath);
+ } catch (Exception e) {
+ //TODO do nothing
+ }
+ }
+
+ protected Pair<RestorableTsFileIOWriter, TsFileResource> createNewFileWriter
+ (String mergeSuffix, String seqDir, List<String> fileNames, long level)
throws IOException {
+ String fileName = fileNames.get(0);
+ String mergeLevelStr = fileName
+ .substring(0, fileName.lastIndexOf(FILE_NAME_SEPARATOR) + 1)
+ + level + TSFILE_SUFFIX + mergeSuffix;
+ // use the minimum version as the version of the new file
+ File newFile = FSFactoryProducer.getFSFactory().getFile(seqDir,
mergeLevelStr);
+ return new Pair<>(new RestorableTsFileIOWriter(newFile), new
TsFileResource(newFile));
+ }
+
+ private void cleanUp(Map<Long, Map<Long, List<TsFileResource>>> resources,
+ TsFileResource newTsResource, long level, long timePartition, String
parentPath) {
+ writeLock();
+ try {
+ Map<Long, List<TsFileResource>> cleanRes = resources.get(timePartition);
+ for (Entry<Long, List<TsFileResource>> res : cleanRes.entrySet()) {
+ if (parentPath.contains("unsequence")) {
+ unSequenceTsFileResources.get(timePartition).get((int) (long)
(res.getKey()))
+ .removeAll(res.getValue());
+ } else {
+ sequenceTsFileResources.get(timePartition).get((int) (long)
(res.getKey()))
+ .removeAll(res.getValue());
+ }
+ for (TsFileResource deleteRes : res.getValue()) {
+ deleteRes.delete();
+ }
+ }
+ File oldFile = newTsResource.getTsFile();
+ File newLevelFile = new
File(oldFile.getAbsolutePath().replace(MERGE_SUFFIX, ""));
+ FSFactoryProducer.getFSFactory().moveFile(oldFile, newLevelFile);
+ FSFactoryProducer.getFSFactory().moveFile(
+ FSFactoryProducer.getFSFactory().getFile(oldFile + RESOURCE_SUFFIX),
+ FSFactoryProducer.getFSFactory().getFile(newLevelFile +
RESOURCE_SUFFIX));
+ newTsResource.setFile(newLevelFile);
+ if (parentPath.contains("unsequence")) {
+ unSequenceTsFileResources.get(timePartition).get((int) (level +
1)).add(newTsResource);
+ } else {
+ sequenceTsFileResources.get(timePartition).get((int) (level +
1)).add(newTsResource);
+ }
+ } catch (Exception e) {
+ //TODO do nothing
+ } finally {
+ writeUnlock();
+ }
}
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest.java
new file mode 100644
index 0000000..5656c38
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest.java
@@ -0,0 +1,97 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.db.engine.storagegroup;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.constant.TestConstant;
+import org.apache.iotdb.db.engine.MetadataManagerHelper;
+import org.apache.iotdb.db.engine.compaction.CompactionStrategy;
+import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy;
+import org.apache.iotdb.db.engine.merge.manage.MergeManager;
+import org.apache.iotdb.db.exception.StorageGroupProcessorException;
+import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
+import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.write.record.TSRecord;
+import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+public class TiredCompactionTest {
+
+ private String storageGroup = "root.vehicle.d0";
+ private String systemDir = TestConstant.OUTPUT_DATA_DIR.concat("info");
+ private String deviceId = "root.vehicle.d0";
+ private String measurementId = "s0";
+ private StorageGroupProcessor processor;
+ private QueryContext context = EnvironmentUtils.TEST_QUERY_CONTEXT;
+
+ @Before
+ public void setUp() throws Exception {
+
IoTDBDescriptor.getInstance().getConfig().setCompactionStrategy(CompactionStrategy.TIRED_COMPACTION);
+ IoTDBDescriptor.getInstance().getConfig().setOverlapSplit(false);
+ MetadataManagerHelper.initMetadata();
+ EnvironmentUtils.envSetUp();
+ processor = new DummySGP(systemDir, storageGroup);
+ MergeManager.getINSTANCE().start();
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ processor.syncDeleteDataFiles();
+ EnvironmentUtils.cleanEnv();
+ EnvironmentUtils.cleanDir(TestConstant.OUTPUT_DATA_DIR);
+ MergeManager.getINSTANCE().stop();
+ EnvironmentUtils.cleanEnv();
+ IoTDBDescriptor.getInstance().getConfig().setOverlapSplit(true);
+ }
+
+ @Test
+ public void testNoSplitSG() throws Exception {
+ insertAndSyncClose(1, 10);
+ insertAndSyncClose(11, 20);
+ insertAndSyncClose(21, 30);
+ insertAndSyncClose(31, 40);
+ }
+
+ private void insertAndSyncClose(int start, int end) throws Exception {
+ TSRecord record;
+ for (int j = start; j <= end; j++) {
+ record = new TSRecord(j, deviceId);
+ record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId,
String.valueOf(j)));
+ processor.insert(new InsertRowPlan(record));
+ }
+
+ processor.syncCloseAllWorkingTsFileProcessors();
+
+ while (processor.isCompactionMergeWorking()) {
+ Thread.sleep(1000);
+ }
+ }
+
+ class DummySGP extends NoSplitStorageGroupProcessor {
+
+ DummySGP(String systemInfoDir, String storageGroupName) throws
StorageGroupProcessorException {
+ super(systemInfoDir, storageGroupName, new
TsFileFlushPolicy.DirectFlushPolicy());
+ }
+
+ }
+}
\ No newline at end of file
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest1.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest1.java
new file mode 100644
index 0000000..dddcd81
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TiredCompactionTest1.java
@@ -0,0 +1,98 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.db.engine.storagegroup;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.constant.TestConstant;
+import org.apache.iotdb.db.engine.MetadataManagerHelper;
+import org.apache.iotdb.db.engine.compaction.CompactionStrategy;
+import org.apache.iotdb.db.engine.flush.TsFileFlushPolicy;
+import org.apache.iotdb.db.engine.merge.manage.MergeManager;
+import org.apache.iotdb.db.exception.StorageGroupProcessorException;
+import org.apache.iotdb.db.qp.physical.crud.InsertRowPlan;
+import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.write.record.TSRecord;
+import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+public class TiredCompactionTest1 {
+
+ private String storageGroup = "root.vehicle.d0";
+ private String systemDir = TestConstant.OUTPUT_DATA_DIR.concat("info");
+ private String deviceId = "root.vehicle.d0";
+ private String measurementId = "s0";
+ private StorageGroupProcessor processor;
+ private QueryContext context = EnvironmentUtils.TEST_QUERY_CONTEXT;
+
+ @Before
+ public void setUp() throws Exception {
+
IoTDBDescriptor.getInstance().getConfig().setCompactionStrategy(CompactionStrategy.TIRED_COMPACTION);
+ IoTDBDescriptor.getInstance().getConfig().setOverlapSplit(false);
+ MetadataManagerHelper.initMetadata();
+ EnvironmentUtils.envSetUp();
+ processor = new DummySGP(systemDir, storageGroup);
+ MergeManager.getINSTANCE().start();
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ processor.syncDeleteDataFiles();
+ EnvironmentUtils.cleanEnv();
+ EnvironmentUtils.cleanDir(TestConstant.OUTPUT_DATA_DIR);
+ MergeManager.getINSTANCE().stop();
+ EnvironmentUtils.cleanEnv();
+ IoTDBDescriptor.getInstance().getConfig().setOverlapSplit(true);
+ }
+
+ @Test
+ public void testSplitSG() throws Exception {
+ insertAndSyncClose(1, 10);
+ insertAndSyncClose(11, 20);
+
+ insertAndSyncClose(21, 30);
+ insertAndSyncClose(31, 40);
+ }
+
+ private void insertAndSyncClose(int start, int end) throws Exception {
+ TSRecord record;
+ for (int j = start; j <= end; j++) {
+ record = new TSRecord(j, deviceId);
+ record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId,
String.valueOf(j)));
+ processor.insert(new InsertRowPlan(record));
+ }
+
+ processor.syncCloseAllWorkingTsFileProcessors();
+
+ while (processor.isCompactionMergeWorking()) {
+ Thread.sleep(1000);
+ }
+ }
+
+ class DummySGP extends StorageGroupProcessor {
+
+ DummySGP(String systemInfoDir, String storageGroupName) throws
StorageGroupProcessorException {
+ super(systemInfoDir, storageGroupName, new
TsFileFlushPolicy.DirectFlushPolicy());
+ }
+
+ }
+}
\ No newline at end of file