This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch rel/1.2
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/1.2 by this push:
new 1de8f2944aa Revert "[To rel/1.2] add memory estimator on inner space
compaction (#10691)" (#10701)
1de8f2944aa is described below
commit 1de8f2944aa04ab3dc9ac3b05042a051cf1b9236
Author: shuwenwei <[email protected]>
AuthorDate: Thu Jul 27 16:17:10 2023 +0800
Revert "[To rel/1.2] add memory estimator on inner space compaction
(#10691)" (#10701)
---
.../execute/task/AbstractCompactionTask.java | 2 -
.../execute/task/CrossSpaceCompactionTask.java | 1 +
.../execute/task/InnerSpaceCompactionTask.java | 37 +--------
.../estimator/AbstractCompactionEstimator.java | 3 +-
.../estimator/AbstractInnerSpaceEstimator.java | 74 +----------------
.../estimator/CompactionEstimateUtils.java | 96 ----------------------
.../FastCompactionInnerCompactionEstimator.java | 52 ------------
.../compaction/selector/estimator/FileInfo.java | 48 -----------
.../ReadChunkInnerCompactionEstimator.java | 59 -------------
.../ReadPointCrossCompactionEstimator.java | 81 +++++++++++++++++-
.../compaction/CompactionSchedulerTest.java | 1 -
.../utils/CompactionTaskMemCostEstimatorTest.java | 94 ---------------------
.../iotdb/tsfile/read/TsFileSequenceReader.java | 12 +--
13 files changed, 84 insertions(+), 476 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/AbstractCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/AbstractCompactionTask.java
index 51c9de96889..4502f678b5a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/AbstractCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/AbstractCompactionTask.java
@@ -57,8 +57,6 @@ public abstract class AbstractCompactionTask {
protected boolean crossTask;
protected boolean innerSeqTask;
- protected long memoryCost = 0L;
-
protected AbstractCompactionTask(
String storageGroupName,
String dataRegionId,
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/CrossSpaceCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/CrossSpaceCompactionTask.java
index cfc42efda61..a8a74e443da 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/CrossSpaceCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/CrossSpaceCompactionTask.java
@@ -61,6 +61,7 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
protected List<TsFileResource> holdWriteLockList = new ArrayList<>();
protected double selectedSeqFileSize = 0;
protected double selectedUnseqFileSize = 0;
+ protected long memoryCost = 0L;
@SuppressWarnings("squid:S107")
public CrossSpaceCompactionTask(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InnerSpaceCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InnerSpaceCompactionTask.java
index 17d70879032..6a843710332 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InnerSpaceCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InnerSpaceCompactionTask.java
@@ -20,29 +20,21 @@
package org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task;
import org.apache.iotdb.commons.conf.IoTDBConstant;
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.service.metrics.CompactionMetrics;
import org.apache.iotdb.db.service.metrics.FileMetrics;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionExceptionHandler;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionFileCountExceededException;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionMemoryNotEnoughException;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionValidationFailedException;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.ICompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.impl.FastCompactionPerformer;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.impl.ReadChunkCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.subtask.FastCompactionTaskSummary;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.CompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.validator.CompactionValidator;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.AbstractInnerSpaceEstimator;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.FastCompactionInnerCompactionEstimator;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.ReadChunkInnerCompactionEstimator;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResourceList;
import
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResourceStatus;
import
org.apache.iotdb.db.storageengine.dataregion.tsfile.generator.TsFileNameGenerator;
-import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
import org.apache.iotdb.tsfile.exception.write.TsFileNotCompleteException;
@@ -78,8 +70,6 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
protected long maxModsFileSize;
- protected AbstractInnerSpaceEstimator innerSpaceEstimator;
-
public InnerSpaceCompactionTask(
long timePartition,
TsFileManager tsFileManager,
@@ -98,13 +88,6 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
this.selectedTsFileResourceList = selectedTsFileResourceList;
this.sequence = sequence;
this.performer = performer;
- if
(IoTDBDescriptor.getInstance().getConfig().isEnableCompactionMemControl()) {
- if (this.performer instanceof ReadChunkCompactionPerformer) {
- innerSpaceEstimator = new ReadChunkInnerCompactionEstimator();
- } else if (!sequence && this.performer instanceof
FastCompactionInnerCompactionEstimator) {
- innerSpaceEstimator = new FastCompactionInnerCompactionEstimator();
- }
- }
isHoldingReadLock = new boolean[selectedTsFileResourceList.size()];
isHoldingWriteLock = new boolean[selectedTsFileResourceList.size()];
for (int i = 0; i < selectedTsFileResourceList.size(); ++i) {
@@ -129,7 +112,6 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
if (!tsFileManager.isAllowCompaction()) {
return true;
}
-
long startTime = System.currentTimeMillis();
// get resource of target file
String dataDirectory =
selectedTsFileResourceList.get(0).getTsFile().getParent();
@@ -338,8 +320,6 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
isSequence());
}
} finally {
- SystemInfo.getInstance().resetCompactionMemoryCost(memoryCost);
-
SystemInfo.getInstance().decreaseCompactionFileNumCost(selectedTsFileResourceList.size());
releaseAllLocksAndResetStatus();
}
return isSuccess;
@@ -476,24 +456,9 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
return false;
}
}
- if (innerSpaceEstimator != null) {
- memoryCost =
innerSpaceEstimator.estimateInnerCompactionMemory(selectedTsFileResourceList);
- }
- SystemInfo.getInstance().addCompactionMemoryCost(memoryCost, 60);
-
SystemInfo.getInstance().addCompactionFileNum(selectedTsFileResourceList.size(),
60);
} catch (Exception e) {
- if (e instanceof InterruptedException) {
- LOGGER.warn("Interrupted when allocating memory for compaction", e);
- Thread.currentThread().interrupt();
- } else if (e instanceof CompactionMemoryNotEnoughException) {
- LOGGER.info("No enough memory for current compaction task {}", this,
e);
- } else if (e instanceof CompactionFileCountExceededException) {
- LOGGER.info("No enough file num for current compaction task {}", this,
e);
- SystemInfo.getInstance().resetCompactionMemoryCost(memoryCost);
- }
- resetCompactionCandidateStatusForAllSourceFiles();
releaseAllLocksAndResetStatus();
- return false;
+ throw e;
}
return true;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractCompactionEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractCompactionEstimator.java
index c4450637344..3d53867232b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractCompactionEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractCompactionEstimator.java
@@ -52,8 +52,7 @@ public abstract class AbstractCompactionEstimator {
List<TsFileResource> seqResources, TsFileResource unseqResource) throws
IOException;
/** Estimate the memory cost of compacting the source files in inner space
compaction task. */
- public abstract long estimateInnerCompactionMemory(List<TsFileResource>
resources)
- throws IOException;
+ public abstract long estimateInnerCompactionMemory(List<TsFileResource>
resources);
/**
* Construct a new or get an existing TsFileSequenceReader of a TsFile.
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractInnerSpaceEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractInnerSpaceEstimator.java
index a5309aaff30..cc34bd28cc4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractInnerSpaceEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/AbstractInnerSpaceEstimator.java
@@ -19,14 +19,9 @@
package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator;
-import org.apache.iotdb.db.conf.IoTDBConfig;
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
import java.io.IOException;
-import java.util.ArrayList;
import java.util.List;
/**
@@ -34,78 +29,11 @@ import java.util.List;
* its corresponding implementation.
*/
public abstract class AbstractInnerSpaceEstimator extends
AbstractCompactionEstimator {
- protected IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
-
- public long estimateInnerCompactionMemory(List<TsFileResource> resources)
throws IOException {
- InnerCompactionTaskInfo taskInfo =
calculatingCompactionTaskInfo(resources);
- long cost = calculatingMetadataMemoryCost(taskInfo);
- cost += calculatingDataMemoryCost(taskInfo);
- return cost;
- }
-
- public abstract long calculatingMetadataMemoryCost(InnerCompactionTaskInfo
taskInfo);
-
- public abstract long calculatingDataMemoryCost(InnerCompactionTaskInfo
taskInfo);
+ public abstract long estimateInnerCompactionMemory(List<TsFileResource>
resources);
public long estimateCrossCompactionMemory(
List<TsFileResource> seqResources, TsFileResource unseqResource) throws
IOException {
throw new RuntimeException(
"This kind of estimator cannot be used to estimate cross space
compaction task");
}
-
- protected InnerCompactionTaskInfo
calculatingCompactionTaskInfo(List<TsFileResource> resources)
- throws IOException {
- List<FileInfo> fileInfoList = new ArrayList<>();
- for (TsFileResource resource : resources) {
- TsFileSequenceReader reader = getFileReader(resource);
- FileInfo fileInfo =
CompactionEstimateUtils.getSeriesAndDeviceChunkNum(reader);
- fileInfoList.add(fileInfo);
- }
- return new InnerCompactionTaskInfo(resources, fileInfoList);
- }
-
- protected static class InnerCompactionTaskInfo {
- private final List<FileInfo> fileInfoList;
- private int maxConcurrentSeriesNum = 1;
- private long maxChunkMetadataSize = 0;
- private int maxChunkMetadataNumInDevice = 0;
- private long modificationFileSize = 0;
-
- protected InnerCompactionTaskInfo(List<TsFileResource> resources,
List<FileInfo> fileInfoList) {
- this.fileInfoList = fileInfoList;
- for (TsFileResource resource : resources) {
- ModificationFile modificationFile = resource.getModFile();
- if (modificationFile.exists()) {
- modificationFileSize += modificationFile.getSize();
- }
- }
- for (FileInfo fileInfo : fileInfoList) {
- maxConcurrentSeriesNum =
- Math.max(maxConcurrentSeriesNum,
fileInfo.maxAlignedSeriesNumInDevice);
- maxChunkMetadataNumInDevice =
- Math.max(maxChunkMetadataNumInDevice, fileInfo.maxDeviceChunkNum);
- maxChunkMetadataSize = Math.max(maxChunkMetadataSize,
fileInfo.averageChunkMetadataSize);
- }
- }
-
- public int getMaxChunkMetadataNumInDevice() {
- return maxChunkMetadataNumInDevice;
- }
-
- public long getMaxChunkMetadataSize() {
- return maxChunkMetadataSize;
- }
-
- public List<FileInfo> getFileInfoList() {
- return fileInfoList;
- }
-
- public int getMaxConcurrentSeriesNum() {
- return maxConcurrentSeriesNum;
- }
-
- public long getModificationFileSize() {
- return modificationFileSize;
- }
- }
}
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
deleted file mode 100644
index fe6f24fcfb4..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
+++ /dev/null
@@ -1,96 +0,0 @@
-/*
- * 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.storageengine.dataregion.compaction.selector.estimator;
-
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata;
-import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
-
-import java.io.IOException;
-import java.util.List;
-import java.util.Map;
-
-public class CompactionEstimateUtils {
-
- /**
- * Get the details of the tsfile, the returned array contains the following
elements in sequence:
- *
- * <p>total chunk num in this tsfile
- *
- * <p>max chunk num of one timeseries in this tsfile
- *
- * <p>max aligned series num in one device. If there is no aligned series in
this file, then it
- * turns to be -1.
- *
- * <p>max chunk num of one device in this tsfile
- *
- * @throws IOException if io errors occurred
- */
- public static FileInfo getSeriesAndDeviceChunkNum(TsFileSequenceReader
reader)
- throws IOException {
- int totalChunkNum = 0;
- int maxChunkNum = 0;
- int maxAlignedSeriesNumInDevice = -1;
- int maxDeviceChunkNum = 0;
- Map<String, List<TimeseriesMetadata>> deviceMetadata =
reader.getAllTimeseriesMetadata(true);
- for (Map.Entry<String, List<TimeseriesMetadata>> entry :
deviceMetadata.entrySet()) {
- int deviceChunkNum = 0;
- List<TimeseriesMetadata> deviceTimeseriesMetadata = entry.getValue();
- if (deviceTimeseriesMetadata.get(0).getMeasurementId().equals("")) {
- // aligned device
- maxAlignedSeriesNumInDevice =
- Math.max(maxAlignedSeriesNumInDevice,
deviceTimeseriesMetadata.size());
- }
- for (TimeseriesMetadata timeseriesMetadata : deviceTimeseriesMetadata) {
- deviceChunkNum += timeseriesMetadata.getChunkMetadataList().size();
- totalChunkNum += timeseriesMetadata.getChunkMetadataList().size();
- maxChunkNum = Math.max(maxChunkNum,
timeseriesMetadata.getChunkMetadataList().size());
- }
- maxDeviceChunkNum = Math.max(maxDeviceChunkNum, deviceChunkNum);
- }
- long averageChunkMetadataSize =
- totalChunkNum == 0 ? 0 : reader.getAllMetadataSize() / totalChunkNum;
- return new FileInfo(
- totalChunkNum,
- maxChunkNum,
- maxAlignedSeriesNumInDevice,
- maxDeviceChunkNum,
- averageChunkMetadataSize);
- }
-
- public static boolean addReadLock(List<TsFileResource> resources) {
- for (int i = 0; i < resources.size(); i++) {
- TsFileResource resource = resources.get(i);
- resource.readLock();
- if (resource.isDeleted()) {
- // release read lock
- for (int j = 0; j <= i; j++) {
- resources.get(j).readUnlock();
- }
- return false;
- }
- }
- return true;
- }
-
- public static void releaseReadLock(List<TsFileResource> resources) {
- resources.forEach(TsFileResource::readUnlock);
- }
-}
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
deleted file mode 100644
index 9de3f58af9b..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FastCompactionInnerCompactionEstimator.java
+++ /dev/null
@@ -1,52 +0,0 @@
-/*
- * 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.storageengine.dataregion.compaction.selector.estimator;
-
-public class FastCompactionInnerCompactionEstimator extends
AbstractInnerSpaceEstimator {
-
- /**
- * The metadata algorithm is: maxChunkMetaDataSize * maxChunkNumber *
fileSize * maxSeriesNumber
- *
- * @return estimate metadata memory cost
- */
- @Override
- public long calculatingMetadataMemoryCost(InnerCompactionTaskInfo taskInfo) {
- return taskInfo.getFileInfoList().size()
- * taskInfo.getMaxChunkMetadataNumInDevice()
- * taskInfo.getMaxChunkMetadataSize()
- * Math.max(config.getSubCompactionTaskNum(),
taskInfo.getMaxConcurrentSeriesNum());
- }
-
- /**
- * The data algorithm is: (targetChunkSize * fileSize * compressionRatio *
maxSeriesNumber) +
- * modsFileSize
- *
- * @return estimate data memory cost
- */
- @Override
- public long calculatingDataMemoryCost(InnerCompactionTaskInfo taskInfo) {
- long cost =
- config.getTargetChunkSize()
- * taskInfo.getFileInfoList().size()
- * Math.max(config.getSubCompactionTaskNum(),
taskInfo.getMaxConcurrentSeriesNum());
- cost += taskInfo.getModificationFileSize();
- return cost;
- }
-}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FileInfo.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FileInfo.java
deleted file mode 100644
index 7ff1ae9be3d..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/FileInfo.java
+++ /dev/null
@@ -1,48 +0,0 @@
-/*
- * 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.storageengine.dataregion.compaction.selector.estimator;
-
-public class FileInfo {
- // total chunk num in this tsfile
- int totalChunkNum = 0;
- // max chunk num of one timeseries in this tsfile
- int maxSeriesChunkNum = 0;
- // max aligned series num in one device. If there is no aligned series in
this file, then it
- // turns to be -1.
- int maxAlignedSeriesNumInDevice = -1;
- // max chunk num of one device in this tsfile
- @SuppressWarnings("squid:S1068")
- int maxDeviceChunkNum = 0;
-
- long averageChunkMetadataSize = 0;
-
- public FileInfo(
- int totalChunkNum,
- int maxSeriesChunkNum,
- int maxAlignedSeriesNumInDevice,
- int maxDeviceChunkNum,
- long averageChunkMetadataSize) {
- this.totalChunkNum = totalChunkNum;
- this.maxSeriesChunkNum = maxSeriesChunkNum;
- this.maxAlignedSeriesNumInDevice = maxAlignedSeriesNumInDevice;
- this.maxDeviceChunkNum = maxDeviceChunkNum;
- this.averageChunkMetadataSize = averageChunkMetadataSize;
- }
-}
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
deleted file mode 100644
index 02924768952..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadChunkInnerCompactionEstimator.java
+++ /dev/null
@@ -1,59 +0,0 @@
-/*
- * 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.storageengine.dataregion.compaction.selector.estimator;
-
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
-
-public class ReadChunkInnerCompactionEstimator extends
AbstractInnerSpaceEstimator {
-
- @Override
- public long calculatingMetadataMemoryCost(InnerCompactionTaskInfo taskInfo) {
- long cost = 0;
- // add ChunkMetadata size of MultiTsFileDeviceIterator
- cost +=
- taskInfo.getFileInfoList().size()
- * taskInfo.getMaxChunkMetadataNumInDevice()
- * taskInfo.getMaxChunkMetadataSize();
-
- // add ChunkMetadata size of targetFileWriter
- long sizeForFileWriter =
- (long)
- ((double) SystemInfo.getInstance().getMemorySizeForCompaction()
- /
IoTDBDescriptor.getInstance().getConfig().getCompactionThreadCount()
- *
IoTDBDescriptor.getInstance().getConfig().getChunkMetadataSizeProportion());
- cost += sizeForFileWriter;
-
- return cost;
- }
-
- @Override
- public long calculatingDataMemoryCost(InnerCompactionTaskInfo taskInfo) {
- // add max target chunk size and max source chunk size
- long cost =
- 2
- * taskInfo.getMaxConcurrentSeriesNum()
- * IoTDBDescriptor.getInstance().getConfig().getTargetChunkSize();
-
- // add modification file size
- cost += taskInfo.getModificationFileSize();
- return cost;
- }
-}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadPointCrossCompactionEstimator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadPointCrossCompactionEstimator.java
index d8908cfd0e6..318420a831f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadPointCrossCompactionEstimator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/ReadPointCrossCompactionEstimator.java
@@ -22,6 +22,7 @@ package
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimat
import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata;
import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
import org.slf4j.Logger;
@@ -30,6 +31,7 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
public class ReadPointCrossCompactionEstimator extends
AbstractCrossSpaceEstimator {
private static final Logger logger =
@@ -75,7 +77,18 @@ public class ReadPointCrossCompactionEstimator extends
AbstractCrossSpaceEstimat
private boolean addReadLock(List<TsFileResource> seqResources,
TsFileResource unseqResource) {
List<TsFileResource> allResources = new ArrayList<>(seqResources);
allResources.add(unseqResource);
- return CompactionEstimateUtils.addReadLock(allResources);
+ for (int i = 0; i < allResources.size(); i++) {
+ TsFileResource resource = allResources.get(i);
+ resource.readLock();
+ if (resource.isDeleted()) {
+ // release read lock
+ for (int j = 0; j <= i; j++) {
+ allResources.get(j).readUnlock();
+ }
+ return false;
+ }
+ }
+ return true;
}
private void releaseReadLock(List<TsFileResource> seqResources,
TsFileResource unseqResource) {
@@ -91,7 +104,7 @@ public class ReadPointCrossCompactionEstimator extends
AbstractCrossSpaceEstimat
*/
private long calculateReadingUnseqFile(TsFileResource unseqResource) throws
IOException {
TsFileSequenceReader reader = getFileReader(unseqResource);
- FileInfo fileInfo =
CompactionEstimateUtils.getSeriesAndDeviceChunkNum(reader);
+ FileInfo fileInfo = getSeriesAndDeviceChunkNum(reader);
// it is max aligned series num of one device when tsfile contains aligned
series,
// else is sub compaction task num.
int concurrentSeriesNum =
@@ -128,7 +141,7 @@ public class ReadPointCrossCompactionEstimator extends
AbstractCrossSpaceEstimat
long cost = 0;
for (TsFileResource seqResource : seqResources) {
TsFileSequenceReader reader = getFileReader(seqResource);
- FileInfo fileInfo =
CompactionEstimateUtils.getSeriesAndDeviceChunkNum(reader);
+ FileInfo fileInfo = getSeriesAndDeviceChunkNum(reader);
// it is max aligned series num of one device when tsfile contains
aligned series,
// else is sub compaction task num.
int concurrentSeriesNum =
@@ -194,4 +207,66 @@ public class ReadPointCrossCompactionEstimator extends
AbstractCrossSpaceEstimat
return cost;
}
+
+ /**
+ * Get the details of the tsfile, the returned array contains the following
elements in sequence:
+ *
+ * <p>total chunk num in this tsfile
+ *
+ * <p>max chunk num of one timeseries in this tsfile
+ *
+ * <p>max aligned series num in one device. If there is no aligned series in
this file, then it
+ * turns to be -1.
+ *
+ * <p>max chunk num of one device in this tsfile
+ *
+ * @throws IOException if io errors occurred
+ */
+ private FileInfo getSeriesAndDeviceChunkNum(TsFileSequenceReader reader)
throws IOException {
+ int totalChunkNum = 0;
+ int maxChunkNum = 0;
+ int maxAlignedSeriesNumInDevice = -1;
+ int maxDeviceChunkNum = 0;
+ Map<String, List<TimeseriesMetadata>> deviceMetadata =
reader.getAllTimeseriesMetadata(true);
+ for (Map.Entry<String, List<TimeseriesMetadata>> entry :
deviceMetadata.entrySet()) {
+ int deviceChunkNum = 0;
+ List<TimeseriesMetadata> deviceTimeseriesMetadata = entry.getValue();
+ if (deviceTimeseriesMetadata.get(0).getMeasurementId().equals("")) {
+ // aligned device
+ maxAlignedSeriesNumInDevice =
+ Math.max(maxAlignedSeriesNumInDevice,
deviceTimeseriesMetadata.size());
+ }
+ for (TimeseriesMetadata timeseriesMetadata : deviceTimeseriesMetadata) {
+ deviceChunkNum += timeseriesMetadata.getChunkMetadataList().size();
+ totalChunkNum += timeseriesMetadata.getChunkMetadataList().size();
+ maxChunkNum = Math.max(maxChunkNum,
timeseriesMetadata.getChunkMetadataList().size());
+ }
+ maxDeviceChunkNum = Math.max(maxDeviceChunkNum, deviceChunkNum);
+ }
+ return new FileInfo(totalChunkNum, maxChunkNum,
maxAlignedSeriesNumInDevice, maxDeviceChunkNum);
+ }
+
+ private class FileInfo {
+ // total chunk num in this tsfile
+ private int totalChunkNum = 0;
+ // max chunk num of one timeseries in this tsfile
+ private int maxSeriesChunkNum = 0;
+ // max aligned series num in one device. If there is no aligned series in
this file, then it
+ // turns to be -1.
+ private int maxAlignedSeriesNumInDevice = -1;
+ // max chunk num of one device in this tsfile
+ @SuppressWarnings("squid:S1068")
+ private int maxDeviceChunkNum = 0;
+
+ public FileInfo(
+ int totalChunkNum,
+ int maxSeriesChunkNum,
+ int maxAlignedSeriesNumInDevice,
+ int maxDeviceChunkNum) {
+ this.totalChunkNum = totalChunkNum;
+ this.maxSeriesChunkNum = maxSeriesChunkNum;
+ this.maxAlignedSeriesNumInDevice = maxAlignedSeriesNumInDevice;
+ this.maxDeviceChunkNum = maxDeviceChunkNum;
+ }
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionSchedulerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionSchedulerTest.java
index d417eab7f7f..198e865dd07 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionSchedulerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionSchedulerTest.java
@@ -101,7 +101,6 @@ public class CompactionSchedulerTest {
.getConfig()
.setInnerUnseqCompactionPerformer(InnerUnseqCompactionPerformer.READ_POINT);
IoTDBDescriptor.getInstance().getConfig().setMinCrossCompactionUnseqFileLevel(0);
-
IoTDBDescriptor.getInstance().getConfig().setEnableCompactionMemControl(false);
CompactionTaskManager.getInstance().start();
while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTaskMemCostEstimatorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTaskMemCostEstimatorTest.java
deleted file mode 100644
index fc0e5d3d4c4..00000000000
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTaskMemCostEstimatorTest.java
+++ /dev/null
@@ -1,94 +0,0 @@
-/*
- * 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.storageengine.dataregion.compaction.utils;
-
-import org.apache.iotdb.commons.exception.MetadataException;
-import org.apache.iotdb.db.exception.StorageEngineException;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.AbstractCompactionTest;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.FastCompactionInnerCompactionEstimator;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.ReadChunkInnerCompactionEstimator;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
-
-import org.junit.After;
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
-
-import java.io.IOException;
-import java.util.List;
-
-public class CompactionTaskMemCostEstimatorTest extends AbstractCompactionTest
{
-
- @Before
- public void setUp()
- throws IOException, WriteProcessException, MetadataException,
InterruptedException {
- super.setUp();
- }
-
- @After
- public void tearDown() throws IOException, StorageEngineException {
- super.tearDown();
- }
-
- @Test
- public void testEstimateReadChunkInnerSpaceCompactionTaskMemCost()
- throws IOException, MetadataException, WriteProcessException {
- createFiles(3, 10, 5, 100000, 0, 0, 50, 50, true, true);
- tsFileManager.addAll(seqResources, true);
- List<TsFileResource> tsFileList = tsFileManager.getTsFileList(true);
- System.out.println(tsFileList.get(0).getTsFile().getAbsolutePath());
- long cost = new
ReadChunkInnerCompactionEstimator().estimateInnerCompactionMemory(tsFileList);
- Assert.assertTrue(cost > 0);
- }
-
- @Test
- public void testEstimateReadChunkInnerSpaceCompactionTaskMemCost2()
- throws IOException, MetadataException, WriteProcessException {
- createFiles(3, 10, 5, 100, 0, 0, 50, 50, false, true);
- tsFileManager.addAll(seqResources, true);
- List<TsFileResource> tsFileList = tsFileManager.getTsFileList(true);
- long cost = new
ReadChunkInnerCompactionEstimator().estimateInnerCompactionMemory(tsFileList);
- Assert.assertTrue(cost > 0);
- }
-
- @Test
- public void testEstimateFastCompactionInnerSpaceCompactionTaskMemCost()
- throws IOException, MetadataException, WriteProcessException {
- createFiles(3, 10, 5, 100000, 0, 0, 50, 50, true, true);
- tsFileManager.addAll(seqResources, true);
- List<TsFileResource> tsFileList = tsFileManager.getTsFileList(true);
- System.out.println(tsFileList.get(0).getTsFile().getAbsolutePath());
- long cost =
- new
FastCompactionInnerCompactionEstimator().estimateInnerCompactionMemory(tsFileList);
- Assert.assertTrue(cost > 0);
- }
-
- @Test
- public void testEstimateFastCompactionInnerSpaceCompactionTaskMemCost2()
- throws IOException, MetadataException, WriteProcessException {
- createFiles(3, 10, 5, 100, 0, 0, 50, 50, false, true);
- tsFileManager.addAll(seqResources, true);
- List<TsFileResource> tsFileList = tsFileManager.getTsFileList(true);
- long cost =
- new
FastCompactionInnerCompactionEstimator().estimateInnerCompactionMemory(tsFileList);
- Assert.assertTrue(cost > 0);
- }
-}
diff --git
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java
index b1742330c2c..56602d07b09 100644
---
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java
+++
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java
@@ -228,20 +228,12 @@ public class TsFileSequenceReader implements
AutoCloseable {
return fileMetadataSize;
}
- /** Return the tsfile meta data size of this tsfile. */
- public long getFileMetadataSize() throws IOException {
- return tsFileInput.size() - getFileMetadataPos();
- }
-
/**
* Return the whole meta data size of this tsfile, including ChunkMetadata,
TimeseriesMetadata and
* etc.
*/
- public long getAllMetadataSize() throws IOException {
- if (tsFileMetaData == null) {
- readFileMetadata();
- }
- return tsFileInput.size() - tsFileMetaData.getMetaOffset();
+ public long getFileMetadataSize() throws IOException {
+ return tsFileInput.size() - getFileMetadataPos();
}
/** this function does not modify the position of the file reader. */