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 {