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;
   }
 }

Reply via email to