This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch rc/1.3.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rc/1.3.3 by this push:
new cab3e79da09 [to rc/1.3.3] Calculate metadata memory when estimate
compaction memory (#13426)
cab3e79da09 is described below
commit cab3e79da0984e0db4e70c310c7ad29d94714eb6
Author: shuwenwei <[email protected]>
AuthorDate: Fri Sep 6 14:12:44 2024 +0800
[to rc/1.3.3] Calculate metadata memory when estimate compaction memory
(#13426)
* calculate metadata memory
* fix bug
* fix bug
* modify estimate method
* fix bug
---
.../compaction/io/CompactionTsFileReader.java | 27 ++++++++++++
.../estimator/CompactionEstimateUtils.java | 48 ++++++++++++++++++++++
.../FastCompactionInnerCompactionEstimator.java | 13 +++++-
.../FastCrossSpaceCompactionEstimator.java | 12 +++++-
.../ReadChunkInnerCompactionEstimator.java | 15 ++++++-
5 files changed, 111 insertions(+), 4 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/io/CompactionTsFileReader.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/io/CompactionTsFileReader.java
index de769e073f1..279a5f2be0b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/io/CompactionTsFileReader.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/io/CompactionTsFileReader.java
@@ -165,6 +165,33 @@ public class CompactionTsFileReader extends
TsFileSequenceReader {
return timeseriesMetadataOffsetMap;
}
+ public Map<String, Pair<Long, Long>> getTimeseriesMetadataOffsetByDevice(
+ MetadataIndexNode measurementNode) throws IOException {
+ Map<String, Pair<Long, Long>> timeseriesMetadataOffsetMap = new
LinkedHashMap<>();
+ List<IMetadataIndexEntry> childrenEntryList =
measurementNode.getChildren();
+ for (int i = 0; i < childrenEntryList.size(); i++) {
+ long startOffset = childrenEntryList.get(i).getOffset();
+ long endOffset =
+ i == childrenEntryList.size() - 1
+ ? measurementNode.getEndOffset()
+ : childrenEntryList.get(i + 1).getOffset();
+ if
(measurementNode.getNodeType().equals(MetadataIndexNodeType.LEAF_MEASUREMENT)) {
+ // leaf measurement node
+ timeseriesMetadataOffsetMap.put(
+ childrenEntryList.get(i).getCompareKey().toString(),
+ new Pair<>(startOffset, endOffset));
+ } else {
+ // internal measurement node
+ ByteBuffer nextBuffer = readData(startOffset, endOffset);
+ MetadataIndexNode nextLayerMeasurementNode =
+ MetadataIndexNode.deserializeFrom(nextBuffer, false);
+ timeseriesMetadataOffsetMap.putAll(
+ getTimeseriesMetadataOffsetByDevice(nextLayerMeasurementNode));
+ }
+ }
+ return timeseriesMetadataOffsetMap;
+ }
+
private void acquireReadDataSizeWithCompactionReadRateLimiter(int
readDataSize) {
CompactionTaskManager.getInstance().getCompactionReadOperationRateLimiter().acquire(1);
CompactionTaskManager.getInstance().getCompactionReadRateLimiter().acquire(readDataSize);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
index aab182a21af..d1e242952f4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
@@ -20,16 +20,20 @@
package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.io.CompactionTsFileReader;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.constant.CompactionType;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
import org.apache.tsfile.file.metadata.ChunkMetadata;
import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.file.metadata.MetadataIndexNode;
import org.apache.tsfile.read.TsFileDeviceIterator;
import org.apache.tsfile.read.TsFileSequenceReader;
import org.apache.tsfile.utils.Pair;
import java.io.IOException;
+import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
@@ -97,6 +101,50 @@ public class CompactionEstimateUtils {
averageChunkMetadataSize);
}
+ public static long roughEstimateMetadataCostInCompaction(
+ List<TsFileResource> resources, CompactionType taskType) throws
IOException {
+ if (!CompactionEstimateUtils.addReadLock(resources)) {
+ return -1L;
+ }
+ long cost = 0L;
+ Map<IDeviceID, Long> deviceMetadataSizeMap = new HashMap<>();
+ try {
+ for (TsFileResource resource : resources) {
+ if (resource.modFileExists()) {
+ cost += resource.getModFile().getSize();
+ }
+ try (CompactionTsFileReader reader =
+ new CompactionTsFileReader(resource.getTsFilePath(), taskType)) {
+ for (Map.Entry<IDeviceID, Long> entry :
getDeviceMetadataSizeMap(reader).entrySet()) {
+ deviceMetadataSizeMap.merge(entry.getKey(), entry.getValue(),
Long::sum);
+ }
+ }
+ }
+ return cost +
deviceMetadataSizeMap.values().stream().max(Long::compareTo).orElse(0L);
+ } finally {
+ CompactionEstimateUtils.releaseReadLock(resources);
+ }
+ }
+
+ public static Map<IDeviceID, Long>
getDeviceMetadataSizeMap(CompactionTsFileReader reader)
+ throws IOException {
+ Map<IDeviceID, Long> deviceMetadataSizeMap = new HashMap<>();
+ TsFileDeviceIterator deviceIterator =
reader.getAllDevicesIteratorWithIsAligned();
+ while (deviceIterator.hasNext()) {
+ IDeviceID deviceID = deviceIterator.next().getLeft();
+ MetadataIndexNode firstMeasurementNodeOfCurrentDevice =
+ deviceIterator.getFirstMeasurementNodeOfCurrentDevice();
+ long totalTimeseriesMetadataSizeOfCurrentDevice = 0;
+ Map<String, Pair<Long, Long>> timeseriesMetadataOffsetByDevice =
+
reader.getTimeseriesMetadataOffsetByDevice(firstMeasurementNodeOfCurrentDevice);
+ for (Pair<Long, Long> offsetPair :
timeseriesMetadataOffsetByDevice.values()) {
+ totalTimeseriesMetadataSizeOfCurrentDevice += (offsetPair.right -
offsetPair.left);
+ }
+ deviceMetadataSizeMap.put(deviceID,
totalTimeseriesMetadataSizeOfCurrentDevice);
+ }
+ return deviceMetadataSizeMap;
+ }
+
public static boolean shouldAccurateEstimate(long roughEstimatedMemCost) {
return roughEstimatedMemCost > 0
&& IoTDBDescriptor.getInstance().getConfig().getCompactionThreadCount()
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCompactionInnerCompactionEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCompactionInnerCompactionEstimator.java
index 32140f5a4e5..9c1fe4e0515 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCompactionInnerCompactionEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCompactionInnerCompactionEstimator.java
@@ -19,6 +19,7 @@
package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.constant.CompactionType;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import java.io.IOException;
@@ -80,6 +81,15 @@ public class FastCompactionInnerCompactionEstimator extends
AbstractInnerSpaceEs
@Override
public long roughEstimateInnerCompactionMemory(List<TsFileResource>
resources)
throws IOException {
+ long metadataCost =
+ CompactionEstimateUtils.roughEstimateMetadataCostInCompaction(
+ resources,
+ resources.get(0).isSeq()
+ ? CompactionType.INNER_SEQ_COMPACTION
+ : CompactionType.INNER_UNSEQ_COMPACTION);
+ if (metadataCost < 0) {
+ return metadataCost;
+ }
int maxConcurrentSeriesNum =
Math.max(
config.getCompactionMaxAlignedSeriesNumInOneBatch(),
config.getSubCompactionTaskNum());
@@ -89,6 +99,7 @@ public class FastCompactionInnerCompactionEstimator extends
AbstractInnerSpaceEs
// source files (chunk + uncompressed page) * overlap file num
// target file (chunk + unsealed page writer)
return (maxOverlapFileNum + 1) * maxConcurrentSeriesNum * (maxChunkSize +
maxPageSize)
- + memoryBudgetForFileWriter;
+ + memoryBudgetForFileWriter
+ + metadataCost;
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCrossSpaceCompactionEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCrossSpaceCompactionEstimator.java
index ffd4ceba1be..a4d077646ba 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCrossSpaceCompactionEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCrossSpaceCompactionEstimator.java
@@ -19,6 +19,7 @@
package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.constant.CompactionType;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import java.io.IOException;
@@ -85,6 +86,14 @@ public class FastCrossSpaceCompactionEstimator extends
AbstractCrossSpaceEstimat
List<TsFileResource> sourceFiles = new ArrayList<>(seqResources.size() +
unseqResources.size());
sourceFiles.addAll(seqResources);
sourceFiles.addAll(unseqResources);
+
+ long metadataCost =
+ CompactionEstimateUtils.roughEstimateMetadataCostInCompaction(
+ sourceFiles, CompactionType.CROSS_COMPACTION);
+ if (metadataCost < 0) {
+ return metadataCost;
+ }
+
int maxConcurrentSeriesNum =
Math.max(
config.getCompactionMaxAlignedSeriesNumInOneBatch(),
config.getSubCompactionTaskNum());
@@ -94,6 +103,7 @@ public class FastCrossSpaceCompactionEstimator extends
AbstractCrossSpaceEstimat
// source files (chunk + uncompressed page) * overlap file num
// target files (chunk + unsealed page writer)
return (maxOverlapFileNum + 1) * maxConcurrentSeriesNum * (maxChunkSize +
maxPageSize)
- + memoryBudgetForFileWriter;
+ + memoryBudgetForFileWriter
+ + metadataCost;
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadChunkInnerCompactionEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadChunkInnerCompactionEstimator.java
index cd90dd334f3..9d126b86810 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadChunkInnerCompactionEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadChunkInnerCompactionEstimator.java
@@ -19,8 +19,10 @@
package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.constant.CompactionType;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import java.io.IOException;
import java.util.List;
public class ReadChunkInnerCompactionEstimator extends
AbstractInnerSpaceEstimator {
@@ -71,7 +73,14 @@ public class ReadChunkInnerCompactionEstimator extends
AbstractInnerSpaceEstimat
}
@Override
- public long roughEstimateInnerCompactionMemory(List<TsFileResource>
resources) {
+ public long roughEstimateInnerCompactionMemory(List<TsFileResource>
resources)
+ throws IOException {
+ long metadataCost =
+ CompactionEstimateUtils.roughEstimateMetadataCostInCompaction(
+ resources, CompactionType.INNER_SEQ_COMPACTION);
+ if (metadataCost < 0) {
+ return metadataCost;
+ }
int maxConcurrentSeriesNum =
Math.max(
config.getCompactionMaxAlignedSeriesNumInOneBatch(),
config.getSubCompactionTaskNum());
@@ -79,6 +88,8 @@ public class ReadChunkInnerCompactionEstimator extends
AbstractInnerSpaceEstimat
long maxPageSize = tsFileConfig.getPageSizeInByte();
// source files (chunk + uncompressed page)
// target file (chunk + unsealed page writer)
- return 2 * maxConcurrentSeriesNum * (maxChunkSize + maxPageSize) +
memoryBudgetForFileWriter;
+ return 2 * maxConcurrentSeriesNum * (maxChunkSize + maxPageSize)
+ + memoryBudgetForFileWriter
+ + metadataCost;
}
}