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

Reply via email to