This is an automated email from the ASF dual-hosted git repository. marklau99 pushed a commit to branch fix-compaction-tmp-error in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit fdd4685c29362923bf8692c49356e6d02ed80609 Author: Liu Xuxin <[email protected]> AuthorDate: Fri Apr 28 16:11:16 2023 +0800 refactor compaction metrics by getting temoporal file info by compaction summary --- .../iotdb/db/engine/TsFileMetricManager.java | 63 ++++++++++++++-------- .../performer/impl/FastCompactionPerformer.java | 14 +---- .../impl/ReadChunkCompactionPerformer.java | 11 +--- .../impl/ReadPointCompactionPerformer.java | 20 +------ .../execute/task/AbstractCompactionTask.java | 4 ++ .../execute/task/CompactionTaskSummary.java | 18 +++++++ .../readchunk/AlignedSeriesCompactionExecutor.java | 5 -- 7 files changed, 69 insertions(+), 66 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/engine/TsFileMetricManager.java b/server/src/main/java/org/apache/iotdb/db/engine/TsFileMetricManager.java index 13a205a0ad..a3cf6c8336 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/TsFileMetricManager.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/TsFileMetricManager.java @@ -19,8 +19,13 @@ package org.apache.iotdb.db.engine; +import org.apache.iotdb.db.engine.compaction.execute.task.AbstractCompactionTask; +import org.apache.iotdb.db.engine.compaction.execute.task.CompactionTaskSummary; +import org.apache.iotdb.db.engine.compaction.execute.task.InnerSpaceCompactionTask; +import org.apache.iotdb.db.engine.compaction.schedule.CompactionTaskManager; import org.apache.iotdb.db.service.metrics.FileMetrics; +import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; @@ -35,6 +40,8 @@ public class TsFileMetricManager { private final AtomicInteger modFileNum = new AtomicInteger(0); private final AtomicLong modFileSize = new AtomicLong(0); + private long lastUpdateTime = 0; + private static final long UPDATE_INTERVAL = 10_000L; // compaction temporal files private final AtomicLong innerSeqCompactionTempFileSize = new AtomicLong(0); @@ -102,41 +109,55 @@ public class TsFileMetricManager { modFileSize.addAndGet(-size); } - public void addCompactionTempFileSize(boolean innerSpace, boolean seq, long delta) { - if (innerSpace) { - long unused = - seq - ? innerSeqCompactionTempFileSize.addAndGet(delta) - : innerUnseqCompactionTempFileSize.addAndGet(delta); - } else { - crossCompactionTempFileSize.addAndGet(delta); - } + public long getInnerCompactionTempFileSize(boolean seq) { + updateCompactionTempSize(); + return seq ? innerSeqCompactionTempFileSize.get() : innerUnseqCompactionTempFileSize.get(); } - public void addCompactionTempFileNum(boolean innerSpace, boolean seq, int delta) { - if (innerSpace) { - long unused = - seq - ? innerSeqCompactionTempFileNum.addAndGet(delta) - : innerUnseqCompactionTempFileNum.addAndGet(delta); - } else { - crossCompactionTempFileNum.addAndGet(delta); + private synchronized void updateCompactionTempSize() { + if (System.currentTimeMillis() - lastUpdateTime <= UPDATE_INTERVAL) { + return; + } + lastUpdateTime = System.currentTimeMillis(); + + innerSeqCompactionTempFileSize.set(0); + innerSeqCompactionTempFileNum.set(0); + innerUnseqCompactionTempFileSize.set(0); + innerUnseqCompactionTempFileNum.set(0); + crossCompactionTempFileSize.set(0); + crossCompactionTempFileNum.set(0); + + List<AbstractCompactionTask> runningTasks = + CompactionTaskManager.getInstance().getRunningCompactionTaskList(); + for (AbstractCompactionTask task : runningTasks) { + CompactionTaskSummary summary = task.getSummary(); + if (task instanceof InnerSpaceCompactionTask) { + if (task.isInnerSeqTask()) { + innerSeqCompactionTempFileSize.addAndGet(summary.getTemporalFileSize()); + innerSeqCompactionTempFileNum.addAndGet(1); + } else { + innerUnseqCompactionTempFileSize.addAndGet(summary.getTemporalFileSize()); + innerUnseqCompactionTempFileNum.addAndGet(1); + } + } else { + crossCompactionTempFileSize.addAndGet(summary.getTemporalFileSize()); + crossCompactionTempFileNum.addAndGet(summary.getTemporalFileNum()); + } } - } - - public long getInnerCompactionTempFileSize(boolean seq) { - return seq ? innerSeqCompactionTempFileSize.get() : innerUnseqCompactionTempFileSize.get(); } public long getCrossCompactionTempFileSize() { + updateCompactionTempSize(); return crossCompactionTempFileSize.get(); } public long getInnerCompactionTempFileNum(boolean seq) { + updateCompactionTempSize(); return seq ? innerSeqCompactionTempFileNum.get() : innerUnseqCompactionTempFileNum.get(); } public long getCrossCompactionTempFileNum() { + updateCompactionTempSize(); return crossCompactionTempFileNum.get(); } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/FastCompactionPerformer.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/FastCompactionPerformer.java index d029122f0d..272ef7f18e 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/FastCompactionPerformer.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/FastCompactionPerformer.java @@ -22,7 +22,6 @@ import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.engine.TsFileMetricManager; import org.apache.iotdb.db.engine.compaction.execute.performer.ICrossCompactionPerformer; import org.apache.iotdb.db.engine.compaction.execute.performer.ISeqCompactionPerformer; import org.apache.iotdb.db.engine.compaction.execute.performer.IUnseqCompactionPerformer; @@ -105,8 +104,7 @@ public class FastCompactionPerformer @Override public void perform() throws IOException, MetadataException, StorageEngineException, InterruptedException { - TsFileMetricManager.getInstance() - .addCompactionTempFileNum(!isCrossCompaction, !seqFiles.isEmpty(), targetFiles.size()); + this.subTaskSummary.setTemporalFileNum(targetFiles.size()); try (MultiTsFileDeviceIterator deviceIterator = new MultiTsFileDeviceIterator(seqFiles, unseqFiles, readerCacheMap); AbstractCompactionWriter compactionWriter = @@ -139,11 +137,7 @@ public class FastCompactionPerformer // check whether to flush chunk metadata or not compactionWriter.checkAndMayFlushChunkMetadata(); // Add temp file metrics - long currentTempFileSize = compactionWriter.getWriterSize(); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize( - !isCrossCompaction, !seqFiles.isEmpty(), currentTempFileSize - tempFileSize); - tempFileSize = currentTempFileSize; + subTaskSummary.setTemporalFileSize(compactionWriter.getWriterSize()); sortedSourceFiles.clear(); } compactionWriter.endFile(); @@ -156,10 +150,6 @@ public class FastCompactionPerformer sortedSourceFiles = null; readerCacheMap = null; modificationCache = null; - TsFileMetricManager.getInstance() - .addCompactionTempFileNum(!isCrossCompaction, !seqFiles.isEmpty(), -targetFiles.size()); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(!isCrossCompaction, !seqFiles.isEmpty(), -tempFileSize); } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadChunkCompactionPerformer.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadChunkCompactionPerformer.java index e97f2e5fe3..da0c9c91dd 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadChunkCompactionPerformer.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadChunkCompactionPerformer.java @@ -22,7 +22,6 @@ import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.engine.TsFileMetricManager; import org.apache.iotdb.db.engine.compaction.execute.performer.ISeqCompactionPerformer; import org.apache.iotdb.db.engine.compaction.execute.task.CompactionTaskSummary; import org.apache.iotdb.db.engine.compaction.execute.utils.MultiTsFileDeviceIterator; @@ -50,7 +49,6 @@ public class ReadChunkCompactionPerformer implements ISeqCompactionPerformer { private TsFileResource targetResource; private List<TsFileResource> seqFiles; private CompactionTaskSummary summary; - private long tempFileSize = 0L; public ReadChunkCompactionPerformer(List<TsFileResource> sourceFiles, TsFileResource targetFile) { this.seqFiles = sourceFiles; @@ -72,7 +70,6 @@ public class ReadChunkCompactionPerformer implements ISeqCompactionPerformer { ((double) SystemInfo.getInstance().getMemorySizeForCompaction() / IoTDBDescriptor.getInstance().getConfig().getCompactionThreadCount() * IoTDBDescriptor.getInstance().getConfig().getChunkMetadataSizeProportion()); - TsFileMetricManager.getInstance().addCompactionTempFileNum(true, true, 1); try (MultiTsFileDeviceIterator deviceIterator = new MultiTsFileDeviceIterator(seqFiles); TsFileIOWriter writer = new TsFileIOWriter(targetResource.getTsFile(), true, sizeForFileWriter)) { @@ -87,19 +84,13 @@ public class ReadChunkCompactionPerformer implements ISeqCompactionPerformer { compactNotAlignedSeries(device, targetResource, writer, deviceIterator); } // update temporal file metrics - long newTempFileSize = writer.getPos(); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(true, true, newTempFileSize - tempFileSize); - tempFileSize = newTempFileSize; + summary.setTemporalFileSize(writer.getPos()); } for (TsFileResource tsFileResource : seqFiles) { targetResource.updatePlanIndexes(tsFileResource); } writer.endFile(); - } finally { - TsFileMetricManager.getInstance().addCompactionTempFileSize(true, true, -tempFileSize); - TsFileMetricManager.getInstance().addCompactionTempFileNum(true, true, -1); } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadPointCompactionPerformer.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadPointCompactionPerformer.java index 8384916b92..a0a38caa8a 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadPointCompactionPerformer.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/performer/impl/ReadPointCompactionPerformer.java @@ -25,7 +25,6 @@ import org.apache.iotdb.commons.path.AlignedPath; import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.engine.TsFileMetricManager; import org.apache.iotdb.db.engine.compaction.execute.performer.ICrossCompactionPerformer; import org.apache.iotdb.db.engine.compaction.execute.performer.IUnseqCompactionPerformer; import org.apache.iotdb.db.engine.compaction.execute.task.CompactionTaskSummary; @@ -103,8 +102,7 @@ public class ReadPointCompactionPerformer QueryResourceManager.getInstance() .getQueryFileManager() .addUsedFilesForQuery(queryId, queryDataSource); - TsFileMetricManager.getInstance() - .addCompactionTempFileNum(seqFiles.isEmpty(), false, targetFiles.size()); + summary.setTemporalFileNum(targetFiles.size()); try (AbstractCompactionWriter compactionWriter = getCompactionWriter(seqFiles, unseqFiles, targetFiles)) { // Do not close device iterator, because tsfile reader is managed by FileReaderManager. @@ -124,6 +122,7 @@ public class ReadPointCompactionPerformer compactNonAlignedSeries( device, deviceIterator, compactionWriter, fragmentInstanceContext, queryDataSource); } + summary.setTemporalFileSize(compactionWriter.getWriterSize()); } compactionWriter.endFile(); @@ -131,10 +130,6 @@ public class ReadPointCompactionPerformer } finally { QueryResourceManager.getInstance().endQuery(queryId); - TsFileMetricManager.getInstance() - .addCompactionTempFileNum(seqFiles.isEmpty(), false, -targetFiles.size()); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(seqFiles.isEmpty(), false, tempFileSize); } } @@ -186,11 +181,6 @@ public class ReadPointCompactionPerformer // check whether to flush chunk metadata or not compactionWriter.checkAndMayFlushChunkMetadata(); } - // add temp file metrics - long currentWriterSize = compactionWriter.getWriterSize(); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(seqFiles.isEmpty(), false, currentWriterSize - tempFileSize); - tempFileSize = currentWriterSize; } private void compactNonAlignedSeries( @@ -238,12 +228,6 @@ public class ReadPointCompactionPerformer // check whether to flush chunk metadata or not compactionWriter.checkAndMayFlushChunkMetadata(); } - - // add temp file metrics - long currentWriterSize = compactionWriter.getWriterSize(); - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(seqFiles.isEmpty(), false, currentWriterSize - tempFileSize); - tempFileSize = currentWriterSize; } /** diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/AbstractCompactionTask.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/AbstractCompactionTask.java index 19e4a0efaa..7fb2d23dc8 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/AbstractCompactionTask.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/AbstractCompactionTask.java @@ -161,6 +161,10 @@ public abstract class AbstractCompactionTask { return crossTask; } + public long getTemporalFileSize() { + return summary.getTemporalFileSize(); + } + public boolean isInnerSeqTask() { return innerSeqTask; } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/CompactionTaskSummary.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/CompactionTaskSummary.java index bcfc60675b..0e1675d47b 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/CompactionTaskSummary.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/task/CompactionTaskSummary.java @@ -32,6 +32,8 @@ public class CompactionTaskSummary { protected int deserializePageCount = 0; protected int mergedChunkNum = 0; protected long processPointNum = 0; + protected long temporalFileSize = 0; + protected int temporalFileNum = 0; public CompactionTaskSummary() {} @@ -134,6 +136,22 @@ public class CompactionTaskSummary { CANCELED } + public void setTemporalFileSize(long temporalFileSize) { + this.temporalFileSize = temporalFileSize; + } + + public long getTemporalFileSize() { + return temporalFileSize; + } + + public void setTemporalFileNum(int temporalFileNum) { + this.temporalFileNum = temporalFileNum; + } + + public int getTemporalFileNum() { + return temporalFileNum; + } + @Override public String toString() { String startTimeInStr = new SimpleDateFormat().format(new Date(startTime)); diff --git a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/utils/executor/readchunk/AlignedSeriesCompactionExecutor.java b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/utils/executor/readchunk/AlignedSeriesCompactionExecutor.java index 2ec407ad29..9bcf8b83b0 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/utils/executor/readchunk/AlignedSeriesCompactionExecutor.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/compaction/execute/utils/executor/readchunk/AlignedSeriesCompactionExecutor.java @@ -19,7 +19,6 @@ package org.apache.iotdb.db.engine.compaction.execute.utils.executor.readchunk; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.engine.TsFileMetricManager; import org.apache.iotdb.db.engine.compaction.execute.task.CompactionTaskSummary; import org.apache.iotdb.db.engine.compaction.schedule.CompactionTaskManager; import org.apache.iotdb.db.engine.compaction.schedule.constant.CompactionType; @@ -162,10 +161,6 @@ public class AlignedSeriesCompactionExecutor { chunkWriter.writeToFileWriter(writer); } writer.checkMetadataSizeAndMayFlush(); - - // update temporal file metrics - TsFileMetricManager.getInstance() - .addCompactionTempFileSize(true, true, writer.getPos() - originTempFileSize); } private void compactOneAlignedChunk(AlignedChunkReader chunkReader, int notNullChunkNum)
