This is an automated email from the ASF dual-hosted git repository.

ejttianyu pushed a commit to branch dynamic_compaction
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dynamic_compaction by this 
push:
     new 7f9a32b  finish heavy hitter strategy
7f9a32b is described below

commit 7f9a32b3ce51cc2fc7822a0ba9cbdd8aeffa763c
Author: EJTTianyu <[email protected]>
AuthorDate: Sat Mar 13 22:34:53 2021 +0800

    finish heavy hitter strategy
---
 .../resources/conf/iotdb-engine.properties         |   4 +-
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |  40 +++
 .../db/engine/compaction/CompactionStrategy.java   |   4 +
 .../level/LevelCompactionTsFileManagement.java     |  43 ++-
 .../HitterLevelCompactionTsFileManagement.java     | 307 +++++++++++++++++++++
 .../engine/compaction/utils/CompactionUtils.java   |  95 +++++++
 .../QueryHeavyHitters.java}                        |  32 ++-
 .../QueryHitterManager.java}                       |  27 +-
 .../QueryHitterStrategy.java}                      |  20 +-
 .../engine/heavyhitter/hitter/DefaultHitter.java   |  56 ++++
 .../iotdb/tsfile/write/writer/TsFileIOWriter.java  |  29 ++
 11 files changed, 590 insertions(+), 67 deletions(-)

diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties 
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 5b482e2..fa16a25 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -298,8 +298,8 @@ default_fill_interval=-1
 ####################
 ### Merge Configurations
 ####################
-# LEVEL_COMPACTION, NO_COMPACTION
-compaction_strategy=NO_COMPACTION
+# LEVEL_COMPACTION, NO_COMPACTION, HITTER_LEVEL_COMPACTION
+compaction_strategy=HITTER_LEVEL_COMPACTION
 
 # Works when the compaction_strategy is LEVEL_COMPACTION.
 # Whether to merge unseq files into seq files or not.
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 24f4556..477b4ae 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
@@ -24,6 +24,7 @@ import java.io.File;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 import org.apache.iotdb.db.conf.directories.DirectoryManager;
+import org.apache.iotdb.db.engine.heavyhitter.QueryHitterStrategy;
 import org.apache.iotdb.db.engine.merge.selector.MergeFileStrategy;
 import org.apache.iotdb.db.engine.compaction.CompactionStrategy;
 import org.apache.iotdb.db.exception.LoadConfigurationException;
@@ -312,6 +313,21 @@ public class IoTDBConfig {
   private CompactionStrategy compactionStrategy = 
CompactionStrategy.LEVEL_COMPACTION;
 
   /**
+   * Query hitter strategy
+   */
+  private QueryHitterStrategy queryHitterStrategy = 
QueryHitterStrategy.DEFAULT_STRATEGY;
+
+  /**
+   * max query path hitter contains
+   */
+  private int maxHitterNum = 5000;
+
+  /**
+   * size ratio of the level merge
+   */
+  private int sizeRatio = 2;
+
+  /**
    * Works when the compaction_strategy is LEVEL_COMPACTION.
    * Whether to merge unseq files into seq files or not.
    */
@@ -1461,6 +1477,30 @@ public class IoTDBConfig {
     this.mergeFileStrategy = mergeFileStrategy;
   }
 
+  public QueryHitterStrategy getQueryHitterStrategy() {
+    return queryHitterStrategy;
+  }
+
+  public void setQueryHitterStrategy(
+      QueryHitterStrategy queryHitterStrategy) {
+    this.queryHitterStrategy = queryHitterStrategy;
+  }
+
+  public int getMaxHitterNum() {
+    return maxHitterNum;
+  }
+
+  public void setMaxHitterNum(int maxHitterNum) {
+    this.maxHitterNum = maxHitterNum;
+  }
+
+  public int getSizeRatio() {
+    return sizeRatio;
+  }
+
+  public void setSizeRatio(int sizeRatio) {
+    this.sizeRatio = sizeRatio;
+  }
 
   public CompactionStrategy getCompactionStrategy() {
     return compactionStrategy;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
index 96ec9f9..79a3c05 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
@@ -20,16 +20,20 @@
 package org.apache.iotdb.db.engine.compaction;
 
 import 
org.apache.iotdb.db.engine.compaction.level.LevelCompactionTsFileManagement;
+import 
org.apache.iotdb.db.engine.compaction.level.hitter.HitterLevelCompactionTsFileManagement;
 import org.apache.iotdb.db.engine.compaction.no.NoCompactionTsFileManagement;
 
 public enum CompactionStrategy {
   LEVEL_COMPACTION,
+  HITTER_LEVEL_COMPACTION,
   NO_COMPACTION;
 
   public TsFileManagement getTsFileManagement(String storageGroupName, String 
storageGroupDir) {
     switch (this) {
       case LEVEL_COMPACTION:
         return new LevelCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
+      case HITTER_LEVEL_COMPACTION:
+        return new HitterLevelCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
       case NO_COMPACTION:
       default:
         return new NoCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
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 2f279d3..296cf75 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
@@ -20,7 +20,6 @@
 package org.apache.iotdb.db.engine.compaction.level;
 
 import static org.apache.iotdb.db.conf.IoTDBConstant.FILE_NAME_SEPARATOR;
-import static 
org.apache.iotdb.db.engine.compaction.no.NoCompactionTsFileManagement.compareFileName;
 import static 
org.apache.iotdb.db.engine.compaction.utils.CompactionLogger.COMPACTION_LOG_NAME;
 import static 
org.apache.iotdb.db.engine.compaction.utils.CompactionLogger.SOURCE_NAME;
 import static 
org.apache.iotdb.db.engine.compaction.utils.CompactionLogger.TARGET_NAME;
@@ -43,11 +42,11 @@ import java.util.concurrent.ConcurrentSkipListMap;
 import java.util.concurrent.CopyOnWriteArrayList;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.cache.ChunkMetadataCache;
-import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
 import org.apache.iotdb.db.engine.compaction.TsFileManagement;
 import org.apache.iotdb.db.engine.compaction.utils.CompactionLogAnalyzer;
 import org.apache.iotdb.db.engine.compaction.utils.CompactionLogger;
 import org.apache.iotdb.db.engine.compaction.utils.CompactionUtils;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
 import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.query.control.FileReaderManager;
 import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
@@ -63,31 +62,31 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   private static final Logger logger = LoggerFactory
       .getLogger(LevelCompactionTsFileManagement.class);
 
-  private final int seqLevelNum = Math
+  protected final int seqLevelNum = Math
       .max(IoTDBDescriptor.getInstance().getConfig().getSeqLevelNum(), 1);
-  private final int seqFileNumInEachLevel = Math
+  protected final int seqFileNumInEachLevel = Math
       
.max(IoTDBDescriptor.getInstance().getConfig().getSeqFileNumInEachLevel(), 1);
-  private final int unseqLevelNum = Math
+  protected final int unseqLevelNum = Math
       .max(IoTDBDescriptor.getInstance().getConfig().getUnseqLevelNum(), 1);
-  private final int unseqFileNumInEachLevel = Math
+  protected final int unseqFileNumInEachLevel = Math
       
.max(IoTDBDescriptor.getInstance().getConfig().getUnseqFileNumInEachLevel(), 1);
 
-  private final boolean enableUnseqCompaction = 
IoTDBDescriptor.getInstance().getConfig()
+  protected final boolean enableUnseqCompaction = 
IoTDBDescriptor.getInstance().getConfig()
       .isEnableUnseqCompaction();
-  private final boolean isForceFullMerge = 
IoTDBDescriptor.getInstance().getConfig()
+  protected 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<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<>();
+  protected final Map<Long, List<SortedSet<TsFileResource>>> 
sequenceTsFileResources = new ConcurrentSkipListMap<>();
+  protected final Map<Long, List<List<TsFileResource>>> 
unSequenceTsFileResources = new ConcurrentSkipListMap<>();
+  protected final List<List<TsFileResource>> forkedSequenceTsFileResources = 
new ArrayList<>();
+  protected final List<List<TsFileResource>> forkedUnSequenceTsFileResources = 
new ArrayList<>();
 
   public LevelCompactionTsFileManagement(String storageGroupName, String 
storageGroupDir) {
     super(storageGroupName, storageGroupDir);
     clear();
   }
 
-  private void deleteLevelFilesInDisk(Collection<TsFileResource> mergeTsFiles) 
{
+  protected void deleteLevelFilesInDisk(Collection<TsFileResource> 
mergeTsFiles) {
     logger.debug("{} [compaction] merge starts to delete real file", 
storageGroupName);
     for (TsFileResource mergeTsFile : mergeTsFiles) {
       deleteLevelFile(mergeTsFile);
@@ -96,7 +95,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     }
   }
 
-  private void deleteLevelFilesInList(long timePartitionId,
+  protected void deleteLevelFilesInList(long timePartitionId,
       Collection<TsFileResource> mergeTsFiles, int level, boolean sequence) {
     logger.debug("{} [compaction] merge starts to delete file list", 
storageGroupName);
     if (sequence) {
@@ -118,7 +117,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     }
   }
 
-  private void deleteLevelFile(TsFileResource seqFile) {
+  protected void deleteLevelFile(TsFileResource seqFile) {
     seqFile.writeLock();
     try {
       ChunkMetadataCache.getInstance().remove(seqFile);
@@ -414,7 +413,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     }
   }
 
-  private void deleteAllSubLevelFiles(boolean isSeq, long timePartition) {
+  protected void deleteAllSubLevelFiles(boolean isSeq, long timePartition) {
     if (isSeq) {
       for (int level = 0; level < 
sequenceTsFileResources.get(timePartition).size();
           level++) {
@@ -452,7 +451,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     }
   }
 
-  private void forkTsFileList(
+  protected void forkTsFileList(
       List<List<TsFileResource>> forkedTsFileResources,
       List rawTsFileResources, int currMaxLevel, int currFileNumInEachLevel) {
     forkedTsFileResources.clear();
@@ -486,7 +485,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   }
 
   @SuppressWarnings("squid:S3776")
-  private void merge(List<List<TsFileResource>> mergeResources, boolean 
sequence,
+  protected void merge(List<List<TsFileResource>> mergeResources, boolean 
sequence,
       long timePartition, int currMaxLevel, int currMaxFileNumInEachLevel) {
     // wait until unseq merge has finished
     while (isUnseqMerging) {
@@ -571,13 +570,13 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
   /**
    * if level < maxLevel-1, the file need compaction else, the file can be 
merged later
    */
-  private File createNewTsFileName(File sourceFile, int level) {
+  protected File createNewTsFileName(File sourceFile, int level) {
     String path = sourceFile.getAbsolutePath();
     String prefixPath = path.substring(0, 
path.lastIndexOf(FILE_NAME_SEPARATOR) + 1);
     return new File(prefixPath + level + TSFILE_SUFFIX);
   }
 
-  private List<SortedSet<TsFileResource>> newSequenceTsFileResources(Long k) {
+  protected List<SortedSet<TsFileResource>> newSequenceTsFileResources(Long k) 
{
     List<SortedSet<TsFileResource>> newSequenceTsFileResources = new 
CopyOnWriteArrayList<>();
     for (int i = 0; i < seqLevelNum; i++) {
       newSequenceTsFileResources.add(Collections.synchronizedSortedSet(new 
TreeSet<>(
@@ -596,7 +595,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     return newSequenceTsFileResources;
   }
 
-  private List<List<TsFileResource>> newUnSequenceTsFileResources(Long k) {
+  protected List<List<TsFileResource>> newUnSequenceTsFileResources(Long k) {
     List<List<TsFileResource>> newUnSequenceTsFileResources = new 
CopyOnWriteArrayList<>();
     for (int i = 0; i < unseqLevelNum; i++) {
       newUnSequenceTsFileResources.add(new CopyOnWriteArrayList<>());
@@ -611,7 +610,7 @@ public class LevelCompactionTsFileManagement extends 
TsFileManagement {
     return Integer.parseInt(mergeLevelStr);
   }
 
-  private TsFileResource getTsFileResource(String filePath, boolean isSeq) 
throws IOException {
+  protected TsFileResource getTsFileResource(String filePath, boolean isSeq) 
throws IOException {
     if (isSeq) {
       for (List<SortedSet<TsFileResource>> tsFileResourcesWithLevel : 
sequenceTsFileResources
           .values()) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/hitter/HitterLevelCompactionTsFileManagement.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/hitter/HitterLevelCompactionTsFileManagement.java
new file mode 100644
index 0000000..5783d3c
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/level/hitter/HitterLevelCompactionTsFileManagement.java
@@ -0,0 +1,307 @@
+/*
+ * 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.compaction.level.hitter;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.engine.cache.ChunkMetadataCache;
+import 
org.apache.iotdb.db.engine.compaction.level.LevelCompactionTsFileManagement;
+import org.apache.iotdb.db.engine.compaction.utils.CompactionUtils;
+import org.apache.iotdb.db.engine.heavyhitter.QueryHitterManager;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.query.control.FileReaderManager;
+import org.apache.iotdb.tsfile.exception.write.TsFileNotCompleteException;
+import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
+import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
+import org.apache.iotdb.tsfile.read.common.Chunk;
+import org.apache.iotdb.tsfile.read.common.Path;
+import org.apache.iotdb.tsfile.write.writer.ForceAppendTsFileWriter;
+import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
+import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class HitterLevelCompactionTsFileManagement extends 
LevelCompactionTsFileManagement {
+
+  private static final Logger logger = LoggerFactory
+      .getLogger(HitterLevelCompactionTsFileManagement.class);
+  private final int sizeRatio = 
IoTDBDescriptor.getInstance().getConfig().getSizeRatio();
+  private final int firstLevelNum = Math
+      
.max(IoTDBDescriptor.getInstance().getConfig().getSeqFileNumInEachLevel(), 1);
+  private final String MERGE_SUFFIX = ".temp";
+
+  public HitterLevelCompactionTsFileManagement(String storageGroupName, String 
storageGroupDir) {
+    super(storageGroupName, storageGroupDir);
+  }
+
+  @Override
+  protected void merge(long timePartition) {
+    merge(forkedSequenceTsFileResources, timePartition);
+    if (enableUnseqCompaction && forkedUnSequenceTsFileResources.size() > 0) {
+      merge(isForceFullMerge, getTsFileList(true), 
forkedUnSequenceTsFileResources.get(0),
+          Long.MAX_VALUE);
+    }
+  }
+
+  protected void merge(List<List<TsFileResource>> mergeResources, long 
timePartition) {
+    // wait until unseq merge has finished
+    while (isUnseqMerging) {
+      try {
+        Thread.sleep(200);
+      } catch (InterruptedException e) {
+        logger.error("{} [Compaction] shutdown", storageGroupName, e);
+        Thread.currentThread().interrupt();
+        return;
+      }
+    }
+    long startTimeMillis = System.currentTimeMillis();
+    try {
+      logger.info("{} start to filter compaction condition", storageGroupName);
+      for (int i = 0; i < seqLevelNum - 1; i++) {
+        if (mergeResources.get(i).size() >= firstLevelNum * 
Math.pow(sizeRatio, i)) {
+          List<TsFileResource> toMergeTsFiles = mergeResources.get(i);
+          logger.info("{} [Hitter Compaction] merge level-{}'s {} TsFiles to 
next level",
+              storageGroupName, i, toMergeTsFiles.size());
+          for (TsFileResource toMergeTsFile : toMergeTsFiles) {
+            logger.info("{} [Hitter Compaction] start to merge TsFile {}", 
storageGroupName,
+                toMergeTsFile);
+          }
+
+          // tmp file which contains all the hitter series
+          File newLevelFile = 
createTempTsFileName(mergeResources.get(i).get(0).getTsFile());
+          TsFileResource newResource = new TsFileResource(newLevelFile);
+          // merge, read  heavy hitters time series from source files and 
write to target file
+          List<PartialPath> unmergedPaths = QueryHitterManager.getQueryHitter()
+              .getTopCompactionSeries(new PartialPath(storageGroupName));
+          CompactionUtils
+              .hitterMerge(newResource, toMergeTsFiles, storageGroupName, new 
HashSet<>(),
+                  unmergedPaths);
+          logger.info(
+              "{} [Compaction] merged level-{}'s {} TsFiles to next level, and 
start to clean up",
+              storageGroupName, i, toMergeTsFiles.size());
+          // do the clean Up
+          writeLockAllFiles(toMergeTsFiles);
+          try {
+            Set<Path> mergedPaths = new HashSet<>(unmergedPaths);
+            List<TsFileIOWriter> writers = new ArrayList<>();
+            for (TsFileResource fileResource : toMergeTsFiles) {
+              // remove cache
+              ChunkMetadataCache.getInstance().remove(fileResource);
+              
FileReaderManager.getInstance().closeFileAndRemoveReader(fileResource.getTsFilePath());
+
+              TsFileIOWriter oldFileWriter = getOldFileWriter(fileResource);
+              // filter all the chunks that have been merged
+              oldFileWriter.filterChunksHitter(mergedPaths);
+              writers.add(oldFileWriter);
+            }
+            TsFileIOWriter newFileWriter = getOldFileWriter(newResource);
+            try (TsFileSequenceReader tmpFileReader =
+                new TsFileSequenceReader(newFileWriter.getFile().getPath())) {
+              Map<String, List<ChunkMetadata>> chunkMetadataListInChunkGroups =
+                  newFileWriter.getDeviceChunkMetadataMap();
+              for (Map.Entry<String, List<ChunkMetadata>> entry : 
chunkMetadataListInChunkGroups
+                  .entrySet()) {
+                String deviceId = entry.getKey();
+                List<ChunkMetadata> chunkMetadataList = entry.getValue();
+                writeMergedChunkGroup(chunkMetadataList, deviceId, 
tmpFileReader, writers.get(0));
+
+                if (Thread.interrupted()) {
+                  Thread.currentThread().interrupt();
+                  writers.get(0).close();
+                  restoreOldFile(newResource);
+                  return;
+                }
+              }
+            }
+            Set<Long> historicalVersions = new HashSet<>();
+            for (TsFileResource tsFileResource : toMergeTsFiles) {
+              
historicalVersions.addAll(tsFileResource.getHistoricalVersions());
+            }
+            toMergeTsFiles.get(0).setHistoricalVersions(historicalVersions);
+            for (TsFileResource tsFileResource: toMergeTsFiles) {
+              // rename file
+              File oldFile = tsFileResource.getTsFile();
+              File newFile = createNewTsFileName(oldFile, i + 1);
+              FSFactoryProducer.getFSFactory().moveFile(oldFile, newFile);
+              FSFactoryProducer.getFSFactory().moveFile(
+                  FSFactoryProducer.getFSFactory().getFile(oldFile + 
TsFileResource.RESOURCE_SUFFIX),
+                  FSFactoryProducer.getFSFactory().getFile(newFile + 
TsFileResource.RESOURCE_SUFFIX));
+              tsFileResource.setFile(newFile);
+              //
+              tsFileResource.serialize();
+              tsFileResource.close();
+            }
+            for (TsFileIOWriter writer: writers){
+              writer.endFile();
+            }
+            newFileWriter.close();
+            newFileWriter.getFile().delete();
+          } finally {
+            writeUnlockAllFiles(toMergeTsFiles);
+          }
+//          System.exit(0);
+          writeLock();
+          try {
+            synchronized (sequenceTsFileResources) {
+              for (TsFileResource tsFileResource : toMergeTsFiles) {
+                // remove and add
+                
sequenceTsFileResources.get(timePartition).get(i).remove(tsFileResource);
+                sequenceTsFileResources.get(timePartition).get(i + 
1).add(tsFileResource);
+                if (mergeResources.size() > i + 1) {
+//                  mergeResources.get(i).remove(tsFileResource);
+                  mergeResources.get(i + 1).add(tsFileResource);
+                }
+              }
+            }
+          } finally {
+            writeUnlock();
+          }
+        }
+      }
+    } catch (Exception e) {
+      logger.error("Error occurred in Compaction Merge thread", e);
+    } finally {
+      // reset the merge working state to false
+      logger.info("{} [Compaction] merge end time consumption: {} ms",
+          storageGroupName, System.currentTimeMillis() - startTimeMillis);
+    }
+  }
+
+  protected File createTempTsFileName(File sourceFile) {
+    String path = sourceFile.getAbsolutePath();
+    return new File(path + MERGE_SUFFIX);
+  }
+
+  protected void writeLockAllFiles(List<TsFileResource> toMergeTsFiles) {
+    int lockCnt;
+    boolean[] locked = new boolean[toMergeTsFiles.size()];
+    while (true) {
+      lockCnt = 0;
+      for (int i = 0; i < toMergeTsFiles.size(); i++) {
+        locked[i] = toMergeTsFiles.get(i).tryWriteLock();
+        if (locked[i]) {
+          lockCnt++;
+        }
+      }
+      if (lockCnt == toMergeTsFiles.size()) {
+        break;
+      } else {
+        for (int i = 0; i < toMergeTsFiles.size(); i++) {
+          if (locked[i]) {
+            toMergeTsFiles.get(i).writeUnlock();
+          }
+        }
+      }
+    }
+  }
+
+  protected void writeUnlockAllFiles(List<TsFileResource> toMergeTsFiles) {
+    for (int i = 0; i < toMergeTsFiles.size(); i++) {
+      toMergeTsFiles.get(i).writeUnlock();
+    }
+  }
+
+  /**
+   * Open an appending writer for an old seq file so we can add new chunks to 
it.
+   */
+  private TsFileIOWriter getOldFileWriter(TsFileResource seqFile) throws 
IOException {
+    TsFileIOWriter oldFileWriter;
+    try {
+      oldFileWriter = new ForceAppendTsFileWriter(seqFile.getTsFile());
+      ((ForceAppendTsFileWriter) oldFileWriter).doTruncate();
+    } catch (TsFileNotCompleteException e) {
+      // this file may already be truncated if this merge is a system reboot 
merge
+      oldFileWriter = new RestorableTsFileIOWriter(seqFile.getTsFile());
+    }
+    return oldFileWriter;
+  }
+
+  private void writeMergedChunkGroup(List<ChunkMetadata> chunkMetadataList, 
String device,
+      TsFileSequenceReader reader, TsFileIOWriter fileWriter)
+      throws IOException {
+    fileWriter.startChunkGroup(device);
+    long maxVersion = 0;
+    for (ChunkMetadata chunkMetaData : chunkMetadataList) {
+      Chunk chunk = reader.readMemChunk(chunkMetaData);
+      fileWriter.writeChunk(chunk, chunkMetaData);
+      maxVersion =
+          chunkMetaData.getVersion() > maxVersion ? chunkMetaData.getVersion() 
: maxVersion;
+    }
+    fileWriter.writeVersion(maxVersion);
+    fileWriter.endChunkGroup();
+  }
+
+  private void restoreOldFile(TsFileResource seqFile) throws IOException {
+    RestorableTsFileIOWriter oldFileRecoverWriter = new 
RestorableTsFileIOWriter(
+        seqFile.getTsFile());
+    if (oldFileRecoverWriter.hasCrashed() && oldFileRecoverWriter.canWrite()) {
+      oldFileRecoverWriter.endFile();
+    } else {
+      oldFileRecoverWriter.close();
+    }
+  }
+
+  @Override
+  public void forkCurrentFileList(long timePartition) {
+    synchronized (sequenceTsFileResources) {
+      forkTsFileList(
+          forkedSequenceTsFileResources,
+          sequenceTsFileResources.computeIfAbsent(timePartition, 
this::newSequenceTsFileResources),
+          seqLevelNum);
+    }
+    // we have to copy all unseq file
+    synchronized (unSequenceTsFileResources) {
+      forkTsFileList(
+          forkedUnSequenceTsFileResources,
+          unSequenceTsFileResources
+              .computeIfAbsent(timePartition, 
this::newUnSequenceTsFileResources),
+          unseqLevelNum + 1);
+    }
+  }
+
+  protected void forkTsFileList(
+      List<List<TsFileResource>> forkedTsFileResources,
+      List rawTsFileResources, int currMaxLevel) {
+    forkedTsFileResources.clear();
+    for (int i = 0; i < currMaxLevel - 1; i++) {
+      List<TsFileResource> forkedLevelTsFileResources = new ArrayList<>();
+      Collection<TsFileResource> levelRawTsFileResources = 
(Collection<TsFileResource>) rawTsFileResources
+          .get(i);
+      for (TsFileResource tsFileResource : levelRawTsFileResources) {
+        if (tsFileResource.isClosed()) {
+          forkedLevelTsFileResources.add(tsFileResource);
+          if (forkedLevelTsFileResources.size() > firstLevelNum * 
Math.pow(sizeRatio, i)) {
+            break;
+          }
+        }
+      }
+      forkedTsFileResources.add(forkedLevelTsFileResources);
+    }
+  }
+}
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 9133777..9ad7926 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
@@ -43,6 +43,7 @@ import 
org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
 import org.apache.iotdb.db.metadata.PartialPath;
 import org.apache.iotdb.db.service.IoTDB;
+import org.apache.iotdb.db.utils.MergeUtils;
 import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
 import org.apache.iotdb.tsfile.read.TimeValuePair;
 import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
@@ -208,6 +209,100 @@ 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 devices the devices to be skipped(used by recover)
+   */
+  @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity 
warning
+  public static void hitterMerge(TsFileResource targetResource,
+      List<TsFileResource> tsFileResources, String storageGroup,
+      Set<String> devices, List<PartialPath> unmergedPaths) throws 
IOException, MetadataException {
+    RestorableTsFileIOWriter writer = new 
RestorableTsFileIOWriter(targetResource.getTsFile());
+    Map<String, TsFileSequenceReader> tsFileSequenceReaderMap = new 
HashMap<>();
+    Map<String, List<Modification>> modificationCache = new HashMap<>();
+    RateLimiter compactionWriteRateLimiter = 
MergeManager.getINSTANCE().getMergeWriteRateLimiter();
+
+    List<List<PartialPath>> devicePaths = 
MergeUtils.splitPathsByDevice(unmergedPaths);
+    for (List<PartialPath> pathList : devicePaths) {
+      String device = pathList.get(0).getDevice();
+      if (devices.contains(device)) {
+        continue;
+      }
+      writer.startChunkGroup(device);
+      // sort chunkMeta by measurement
+      Map<String, Map<TsFileSequenceReader, List<ChunkMetadata>>> 
measurementChunkMetadataMap = new HashMap<>();
+      for (TsFileResource levelResource : tsFileResources) {
+        TsFileSequenceReader reader = 
buildReaderFromTsFileResource(levelResource,
+            tsFileSequenceReaderMap, storageGroup);
+        if (reader == null) {
+          continue;
+        }
+        Map<String, List<ChunkMetadata>> chunkMetadataMap = new HashMap<>();
+        for (PartialPath path : pathList) {
+          chunkMetadataMap.computeIfAbsent(path.getMeasurement(), p -> new 
ArrayList<>());
+          
chunkMetadataMap.get(path.getMeasurement()).addAll(reader.getChunkMetadataList(path));
+        }
+        for (Entry<String, List<ChunkMetadata>> entry : 
chunkMetadataMap.entrySet()) {
+          for (ChunkMetadata chunkMetadata : entry.getValue()) {
+            Map<TsFileSequenceReader, List<ChunkMetadata>> 
readerChunkMetadataMap;
+            String measurementUid = chunkMetadata.getMeasurementUid();
+            if (measurementChunkMetadataMap.containsKey(measurementUid)) {
+              readerChunkMetadataMap = 
measurementChunkMetadataMap.get(measurementUid);
+            } else {
+              readerChunkMetadataMap = new LinkedHashMap<>();
+            }
+            List<ChunkMetadata> chunkMetadataList;
+            if (readerChunkMetadataMap.containsKey(reader)) {
+              chunkMetadataList = readerChunkMetadataMap.get(reader);
+            } else {
+              chunkMetadataList = new ArrayList<>();
+            }
+            chunkMetadataList.add(chunkMetadata);
+            readerChunkMetadataMap.put(reader, chunkMetadataList);
+            measurementChunkMetadataMap
+                .put(chunkMetadata.getMeasurementUid(), 
readerChunkMetadataMap);
+          }
+        }
+      }
+      long maxVersion = Long.MIN_VALUE;
+      for (Entry<String, Map<TsFileSequenceReader, List<ChunkMetadata>>> entry 
: measurementChunkMetadataMap
+          .entrySet()) {
+        Map<TsFileSequenceReader, List<ChunkMetadata>> readerChunkMetadatasMap 
= entry.getValue();
+        boolean isPageEnoughLarge = true;
+        for (List<ChunkMetadata> chunkMetadatas : 
readerChunkMetadatasMap.values()) {
+          for (ChunkMetadata chunkMetadata : chunkMetadatas) {
+            if (chunkMetadata.getNumOfPoints() < MERGE_PAGE_POINT_NUM) {
+              isPageEnoughLarge = false;
+              break;
+            }
+          }
+        }
+        if (isPageEnoughLarge) {
+          logger.debug("{} [Compaction] page enough large, use append merge", 
storageGroup);
+          // append page in chunks, so we do not have to deserialize a chunk
+          maxVersion = writeByAppendMerge(maxVersion, device, 
compactionWriteRateLimiter,
+              entry, targetResource, writer, modificationCache);
+        } else {
+          logger
+              .debug("{} [Compaction] page too small, use deserialize merge", 
storageGroup);
+          // we have to deserialize chunks to merge pages
+          maxVersion = writeByDeserializeMerge(maxVersion, device, 
compactionWriteRateLimiter,
+              entry, targetResource, writer, modificationCache);
+        }
+      }
+      writer.endChunkGroup();
+      writer.writeVersion(maxVersion);
+    }
+
+    for (TsFileSequenceReader reader : tsFileSequenceReaderMap.values()) {
+      reader.close();
+    }
+    writer.endFile();
+    targetResource.close();
+  }
+
+  /**
+   * @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)
    */
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHeavyHitters.java
similarity index 55%
copy from 
server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
copy to 
server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHeavyHitters.java
index 96ec9f9..b9c0989 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHeavyHitters.java
@@ -17,22 +17,24 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.engine.compaction;
+package org.apache.iotdb.db.engine.heavyhitter;
 
-import 
org.apache.iotdb.db.engine.compaction.level.LevelCompactionTsFileManagement;
-import org.apache.iotdb.db.engine.compaction.no.NoCompactionTsFileManagement;
+import java.util.List;
+import org.apache.iotdb.db.exception.metadata.MetadataException;
+import org.apache.iotdb.db.metadata.PartialPath;
 
-public enum CompactionStrategy {
-  LEVEL_COMPACTION,
-  NO_COMPACTION;
+/**
+ * 被用于使用选取合并收益较高的时间序列
+ */
+public interface QueryHeavyHitters {
+
+  /**
+   * 用于接收查询的时间序列
+   */
+  void acceptQuerySeries(PartialPath queryPath);
 
-  public TsFileManagement getTsFileManagement(String storageGroupName, String 
storageGroupDir) {
-    switch (this) {
-      case LEVEL_COMPACTION:
-        return new LevelCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
-      case NO_COMPACTION:
-      default:
-        return new NoCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
-    }
-  }
+  /**
+   * 用于获取合并收益最高的时间序列
+   */
+  List<PartialPath> getTopCompactionSeries(PartialPath sgName) throws 
MetadataException;
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterManager.java
similarity index 57%
copy from 
server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
copy to 
server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterManager.java
index 96ec9f9..1c9f982 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterManager.java
@@ -17,22 +17,25 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.engine.compaction;
+package org.apache.iotdb.db.engine.heavyhitter;
 
-import 
org.apache.iotdb.db.engine.compaction.level.LevelCompactionTsFileManagement;
-import org.apache.iotdb.db.engine.compaction.no.NoCompactionTsFileManagement;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.engine.heavyhitter.hitter.DefaultHitter;
 
-public enum CompactionStrategy {
-  LEVEL_COMPACTION,
-  NO_COMPACTION;
+public class QueryHitterManager {
 
-  public TsFileManagement getTsFileManagement(String storageGroupName, String 
storageGroupDir) {
-    switch (this) {
-      case LEVEL_COMPACTION:
-        return new LevelCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
-      case NO_COMPACTION:
+  private static final QueryHeavyHitters INSTANCE = loadQueryHitters();
+
+  public static QueryHeavyHitters getQueryHitter() {
+    return INSTANCE;
+  }
+
+  private static QueryHeavyHitters loadQueryHitters() {
+    switch 
(IoTDBDescriptor.getInstance().getConfig().getQueryHitterStrategy()) {
+      case DEFAULT_STRATEGY:
       default:
-        return new NoCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
+        return new 
DefaultHitter(IoTDBDescriptor.getInstance().getConfig().getMaxHitterNum());
     }
   }
+
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterStrategy.java
similarity index 55%
copy from 
server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
copy to 
server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterStrategy.java
index 96ec9f9..80a2daf 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionStrategy.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/QueryHitterStrategy.java
@@ -17,22 +17,10 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.engine.compaction;
+package org.apache.iotdb.db.engine.heavyhitter;
 
-import 
org.apache.iotdb.db.engine.compaction.level.LevelCompactionTsFileManagement;
-import org.apache.iotdb.db.engine.compaction.no.NoCompactionTsFileManagement;
+public enum  QueryHitterStrategy {
+  //用于测试的 strategy
+  DEFAULT_STRATEGY;
 
-public enum CompactionStrategy {
-  LEVEL_COMPACTION,
-  NO_COMPACTION;
-
-  public TsFileManagement getTsFileManagement(String storageGroupName, String 
storageGroupDir) {
-    switch (this) {
-      case LEVEL_COMPACTION:
-        return new LevelCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
-      case NO_COMPACTION:
-      default:
-        return new NoCompactionTsFileManagement(storageGroupName, 
storageGroupDir);
-    }
-  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/hitter/DefaultHitter.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/hitter/DefaultHitter.java
new file mode 100644
index 0000000..e85f3f6
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/heavyhitter/hitter/DefaultHitter.java
@@ -0,0 +1,56 @@
+/*
+ * 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.heavyhitter.hitter;
+
+import java.util.List;
+import org.apache.iotdb.db.engine.heavyhitter.QueryHeavyHitters;
+import org.apache.iotdb.db.exception.metadata.MetadataException;
+import org.apache.iotdb.db.metadata.MManager;
+import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.utils.MergeUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class DefaultHitter implements QueryHeavyHitters {
+
+  private static final Logger logger = 
LoggerFactory.getLogger(DefaultHitter.class);
+
+  public DefaultHitter(int maxHitterNum) {
+
+  }
+
+  @Override
+  public void acceptQuerySeries(PartialPath queryPath) {
+    // do nothing
+  }
+
+  @Override
+  public List<PartialPath> getTopCompactionSeries(PartialPath sgName) throws 
MetadataException {
+    List<PartialPath> unmergedSeries =
+        MManager.getInstance().getAllTimeseriesPath(sgName);
+    List<List<PartialPath>> devicePaths = 
MergeUtils.splitPathsByDevice(unmergedSeries);
+    if (devicePaths.size() > 0) {
+      String deviceName = devicePaths.get(0).get(0).getDevice();
+      logger.info("default hitter, top compaction device:{}", deviceName);
+      return devicePaths.get(0);
+    }
+    return null;
+  }
+}
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java
index 7ad89b2..c6dcef2 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/writer/TsFileIOWriter.java
@@ -26,6 +26,7 @@ import java.util.Iterator;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.TreeMap;
 import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
 import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
@@ -417,6 +418,34 @@ public class TsFileIOWriter {
   }
 
   /**
+   * Remove such ChunkMetadata if heavy hitters has merged
+   */
+  public void filterChunksHitter(Set<Path> mergedPaths) {
+    Iterator<ChunkGroupMetadata> chunkGroupMetaDataIterator = 
chunkGroupMetadataList.iterator();
+    while (chunkGroupMetaDataIterator.hasNext()) {
+      ChunkGroupMetadata chunkGroupMetaData = 
chunkGroupMetaDataIterator.next();
+      String deviceId = chunkGroupMetaData.getDevice();
+      int chunkNum = chunkGroupMetaData.getChunkMetadataList().size();
+      Iterator<ChunkMetadata> chunkMetaDataIterator = 
chunkGroupMetaData.getChunkMetadataList()
+          .iterator();
+      while (chunkMetaDataIterator.hasNext()) {
+        ChunkMetadata chunkMetaData = chunkMetaDataIterator.next();
+        Path path = new Path(deviceId, chunkMetaData.getMeasurementUid());
+
+        boolean chunkInValid = mergedPaths.contains(path);
+        if (chunkInValid) {
+          chunkMetaDataIterator.remove();
+          chunkNum--;
+          invalidChunkNum++;
+        }
+      }
+      if (chunkNum == 0) {
+        chunkGroupMetaDataIterator.remove();
+      }
+    }
+  }
+
+  /**
    * write MetaMarker.VERSION with version Then, cache offset-version in 
versionInfo
    */
   public void writeVersion(long version) throws IOException {

Reply via email to