This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 2feb202ca9d add tsfile correctness validation for write compaction and
load (#11392)
2feb202ca9d is described below
commit 2feb202ca9de8fce4fe4c3af4ae37742bd9b791f
Author: Zhijia Cao <[email protected]>
AuthorDate: Tue Oct 31 21:43:03 2023 +0800
add tsfile correctness validation for write compaction and load (#11392)
---
.../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 20 +-
.../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 16 +-
.../db/storageengine/dataregion/DataRegion.java | 10 +-
.../constant/CompactionValidationLevel.java | 26 --
.../execute/task/AbstractCompactionTask.java | 20 ++
.../execute/task/CrossSpaceCompactionTask.java | 16 +-
.../execute/task/InnerSpaceCompactionTask.java | 12 +-
.../task/InsertionCrossSpaceCompactionTask.java | 12 +-
.../compaction/execute/utils/CompactionUtils.java | 163 ---------
.../utils/validator/CompactionValidator.java | 51 ---
.../ResourceAndTsfileCompactionValidator.java | 57 ----
.../dataregion/utils/TsFileResourceUtils.java | 378 +++++++++++++++++++++
.../validate/TsFileResourceAndDataValidator.java | 58 ++++
.../validate/TsFileResourceValidator.java} | 43 +--
.../validate/TsFileValidator.java} | 36 +-
.../compaction/CompactionValidationTest.java | 42 +--
.../CrossSpaceCompactionWithUnusualCasesTest.java | 120 +------
.../FastCrossCompactionPerformerTest.java | 4 +-
.../FastInnerCompactionPerformerTest.java | 2 +-
.../TsFileValidationCorrectnessTests.java | 296 ++++++++++++++++
...eCompactionWithFastPerformerValidationTest.java | 10 +-
...eCrossSpaceCompactionWithFastPerformerTest.java | 12 +-
...sSpaceCompactionWithReadPointPerformerTest.java | 4 +-
.../inner/InnerCompactionEmptyTsFileTest.java | 2 +-
.../compaction/utils/CompactionTestFileWriter.java | 4 +
.../compaction/utils/TsFileGeneratorUtils.java | 123 +++++++
.../resources/conf/iotdb-common.properties | 19 +-
27 files changed, 997 insertions(+), 559 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index 4522d899a5d..2e0a06e912e 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -29,7 +29,6 @@ import org.apache.iotdb.db.audit.AuditLogOperation;
import org.apache.iotdb.db.audit.AuditLogStorage;
import org.apache.iotdb.db.exception.LoadConfigurationException;
import org.apache.iotdb.db.protocol.thrift.impl.ClientRPCServiceImpl;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.constant.CompactionValidationLevel;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.CrossCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.InnerSeqCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.InnerUnseqCompactionPerformer;
@@ -534,8 +533,7 @@ public class IoTDBConfig {
*/
private int subCompactionTaskNum = 4;
- private CompactionValidationLevel compactionValidationLevel =
- CompactionValidationLevel.RESOURCE_ONLY;
+ private boolean enableTsFileValidation = false;
/** The size of candidate compaction task queue. */
private int candidateCompactionTaskQueueSize = 50;
@@ -3675,14 +3673,6 @@ public class IoTDBConfig {
this.schemaRatisLogMax = schemaRatisLogMax;
}
- public CompactionValidationLevel getCompactionValidationLevel() {
- return this.compactionValidationLevel;
- }
-
- public void setCompactionValidationLevel(CompactionValidationLevel level) {
- this.compactionValidationLevel = level;
- }
-
public int getCandidateCompactionTaskQueueSize() {
return candidateCompactionTaskQueueSize;
}
@@ -3808,4 +3798,12 @@ public class IoTDBConfig {
public void setSchemaRatisPeriodicSnapshotInterval(long
schemaRatisPeriodicSnapshotInterval) {
this.schemaRatisPeriodicSnapshotInterval =
schemaRatisPeriodicSnapshotInterval;
}
+
+ public boolean isEnableTsFileValidation() {
+ return enableTsFileValidation;
+ }
+
+ public void setEnableTsFileValidation(boolean enableTsFileValidation) {
+ this.enableTsFileValidation = enableTsFileValidation;
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 19c00c4c941..e172782d929 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -32,7 +32,6 @@ import
org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.schemaengine.rescon.DataNodeSchemaQuotaManager;
import org.apache.iotdb.db.service.metrics.IoTDBInternalLocalReporter;
import org.apache.iotdb.db.storageengine.StorageEngine;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.constant.CompactionValidationLevel;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.CrossCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.InnerSeqCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.constant.InnerUnseqCompactionPerformer;
@@ -694,10 +693,10 @@ public class IoTDBDescriptor {
"compaction_write_throughput_mb_per_sec",
Integer.toString(conf.getCompactionWriteThroughputMbPerSec()))));
- conf.setCompactionValidationLevel(
- CompactionValidationLevel.valueOf(
+ conf.setEnableTsFileValidation(
+ Boolean.parseBoolean(
properties.getProperty(
- "compaction_validation_level",
conf.getCompactionValidationLevel().toString())));
+ "enable_tsfile_validation",
String.valueOf(conf.isEnableTsFileValidation()))));
conf.setCandidateCompactionTaskQueueSize(
Integer.parseInt(
properties.getProperty(
@@ -1121,10 +1120,6 @@ public class IoTDBDescriptor {
}
private void loadCompactionHotModifiedProps(Properties properties) throws
InterruptedException {
- conf.setCompactionValidationLevel(
- CompactionValidationLevel.valueOf(
- properties.getProperty(
- "compaction_validation_level",
conf.getCompactionValidationLevel().toString())));
loadCompactionIsEnabledHotModifiedProps(properties);
@@ -1626,6 +1621,11 @@ public class IoTDBDescriptor {
"enable_query_memory_estimation",
Boolean.toString(conf.isEnableQueryMemoryEstimation()))));
+ conf.setEnableTsFileValidation(
+ Boolean.parseBoolean(
+ properties.getProperty(
+ "enable_tsfile_validation",
String.valueOf(conf.isEnableTsFileValidation()))));
+
// update wal config
long prevDeleteWalFilesPeriodInMs = conf.getDeleteWalFilesPeriodInMs();
loadWALHotModifiedProps(properties);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index 5dd46accb01..0dadb508bb1 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -85,6 +85,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
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.dataregion.tsfile.generator.VersionController;
+import
org.apache.iotdb.db.storageengine.dataregion.utils.validate.TsFileValidator;
import org.apache.iotdb.db.storageengine.dataregion.wal.WALManager;
import org.apache.iotdb.db.storageengine.dataregion.wal.node.IWALNode;
import
org.apache.iotdb.db.storageengine.dataregion.wal.recover.WALRecoverManager;
@@ -2070,7 +2071,8 @@ public class DataRegion implements IDataRegionForQuery {
closeQueryLock.writeLock().lock();
try {
tsFileProcessor.close();
- if (tsFileProcessor.isEmpty() ||
tsFileProcessor.getTsFileResource().isEmpty()) {
+ if (tsFileProcessor.isEmpty()
+ ||
!TsFileValidator.getInstance().validateTsFile(tsFileProcessor.getTsFileResource()))
{
tsFileProcessor.getTsFileResource().remove();
tsFileManager.remove(tsFileProcessor.getTsFileResource(),
tsFileProcessor.isSequence());
} else {
@@ -2222,6 +2224,12 @@ public class DataRegion implements IDataRegionForQuery {
throws LoadFileException {
File tsfileToBeInserted = newTsFileResource.getTsFile();
long newFilePartitionId = newTsFileResource.getTimePartitionWithCheck();
+
+ if (!TsFileValidator.getInstance().validateTsFile(newTsFileResource)) {
+ throw new LoadFileException(
+ "tsfile validate failed, " +
newTsFileResource.getTsFile().getName());
+ }
+
writeLock("loadNewTsFile");
try {
newTsFileResource.setSeq(false);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/constant/CompactionValidationLevel.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/constant/CompactionValidationLevel.java
deleted file mode 100644
index fcd2f4c4614..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/constant/CompactionValidationLevel.java
+++ /dev/null
@@ -1,26 +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.constant;
-
-public enum CompactionValidationLevel {
- NONE,
- RESOURCE_ONLY,
- RESOURCE_AND_TSFILE
-}
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 05f9d005dce..a7571b3abe1 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
@@ -25,6 +25,7 @@ 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.storageengine.dataregion.compaction.constant.CompactionTaskType;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionValidationFailedException;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.FileCannotTransitToCompactingException;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.ICompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
@@ -35,6 +36,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFil
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.TsFileResourceStatus;
+import
org.apache.iotdb.db.storageengine.dataregion.utils.validate.TsFileValidator;
import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
import org.slf4j.Logger;
@@ -370,6 +372,24 @@ public abstract class AbstractCompactionTask {
return CompactionUtils.isDiskHasSpace();
}
+ protected void validateTsFileResource(
+ List<TsFileResource> targetTsFileList, boolean needValidateOverlap) {
+ TsFileValidator validator = TsFileValidator.getInstance();
+ if (!validator.validateTsFiles(targetTsFileList)) {
+ LOGGER.error("Failed to pass compaction validation, target files is {}",
targetTsFileList);
+ throw new CompactionValidationFailedException(
+ "Failed to pass compaction validation, .resources file or tsfile
data is wrong");
+ }
+ if (needValidateOverlap
+ && !validator.validateTsFilesIsHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(timePartition).getArrayList()))
{
+ LOGGER.error("Failed to pass compaction validation, target files is {}",
targetTsFileList);
+ throw new CompactionValidationFailedException(
+ "Failed to pass compaction validation, sequence files has overlap,
time partition id is "
+ + timePartition);
+ }
+ }
+
public CompactionTaskType getCompactionTaskType() {
if (this instanceof CrossSpaceCompactionTask) {
return CompactionTaskType.CROSS;
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 e12df3249a9..a833907a4a8 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
@@ -23,7 +23,6 @@ import org.apache.iotdb.commons.conf.IoTDBConstant;
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.CompactionRecoverException;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.exception.CompactionValidationFailedException;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.ICrossCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.impl.FastCompactionPerformer;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.subtask.FastCompactionTaskSummary;
@@ -32,7 +31,6 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.CompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.SimpleCompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.TsFileIdentifier;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.validator.CompactionValidator;
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.TsFileResourceStatus;
@@ -215,19 +213,7 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
}
}
- CompactionValidator validator = CompactionValidator.getInstance();
- if (!validator.validateCompaction(
- tsFileManager, targetTsfileResourceList, storageGroupName,
timePartition, false)) {
- LOGGER.error(
- "Failed to pass compaction validation, "
- + "source sequence files is: {}, "
- + "unsequence files is {}, "
- + "target files is {}",
- selectedSequenceFiles,
- selectedUnsequenceFiles,
- targetTsfileResourceList);
- throw new CompactionValidationFailedException("Failed to pass
compaction validation");
- }
+ validateTsFileResource(targetTsfileResourceList, true);
lockWrite(selectedSequenceFiles);
lockWrite(selectedUnsequenceFiles);
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 d721f59a658..6221d5c1821 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
@@ -24,7 +24,6 @@ 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.CompactionRecoverException;
-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;
@@ -34,7 +33,6 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.CompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.SimpleCompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.TsFileIdentifier;
-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;
@@ -244,15 +242,7 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
compactionLogger.force();
}
- CompactionValidator validator = CompactionValidator.getInstance();
- if (!validator.validateCompaction(
- tsFileManager, targetTsFileList, storageGroupName, timePartition,
!sequence)) {
- LOGGER.error(
- "Failed to pass compaction validation, source files is: {},
target files is {}",
- selectedTsFileResourceList,
- targetTsFileList);
- throw new CompactionValidationFailedException("Failed to pass
compaction validation");
- }
+ validateTsFileResource(targetTsFileList, sequence);
LOGGER.info(
"{}-{} [Compaction] Compacted target files, try to get the write
lock of source files",
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InsertionCrossSpaceCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InsertionCrossSpaceCompactionTask.java
index 26db2bc7961..95ad2a68efd 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InsertionCrossSpaceCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/InsertionCrossSpaceCompactionTask.java
@@ -28,13 +28,13 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.CompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.SimpleCompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.TsFileIdentifier;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.validator.CompactionValidator;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.utils.InsertionCrossCompactionTaskResource;
import
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
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.TsFileResourceStatus;
import
org.apache.iotdb.db.storageengine.dataregion.tsfile.generator.TsFileNameGenerator;
+import
org.apache.iotdb.db.storageengine.dataregion.utils.validate.TsFileValidator;
import java.io.File;
import java.io.IOException;
@@ -152,13 +152,9 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
replaceTsFileInMemory(
Collections.singletonList(unseqFileToInsert),
Collections.singletonList(targetFile));
- CompactionValidator validator = CompactionValidator.getInstance();
- if (!validator.validateCompaction(
- tsFileManager,
- Collections.singletonList(targetFile),
- storageGroupName,
- timePartition,
- false)) {
+ if (!TsFileValidator.getInstance()
+ .validateTsFilesIsHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(timePartition))) {
LOGGER.error(
"Failed to pass compaction validation, source un seq files is: {},
target files is {}",
unseqFileToInsert,
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/CompactionUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/CompactionUtils.java
index 78e4bbb938d..0aec4782f91 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/CompactionUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/CompactionUtils.java
@@ -26,30 +26,12 @@ import org.apache.iotdb.commons.service.metric.enums.Tag;
import org.apache.iotdb.db.service.metrics.FileMetrics;
import org.apache.iotdb.db.storageengine.dataregion.modification.Modification;
import
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
-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.timeindex.DeviceTimeIndex;
import org.apache.iotdb.metrics.utils.MetricLevel;
import org.apache.iotdb.metrics.utils.SystemMetric;
-import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
-import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
-import org.apache.iotdb.tsfile.encoding.decoder.Decoder;
-import org.apache.iotdb.tsfile.enums.TSDataType;
-import org.apache.iotdb.tsfile.file.MetaMarker;
-import org.apache.iotdb.tsfile.file.header.ChunkGroupHeader;
-import org.apache.iotdb.tsfile.file.header.ChunkHeader;
-import org.apache.iotdb.tsfile.file.header.PageHeader;
import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
-import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
-import org.apache.iotdb.tsfile.read.common.BatchData;
-import org.apache.iotdb.tsfile.read.reader.page.PageReader;
-import org.apache.iotdb.tsfile.read.reader.page.TimePageReader;
-import org.apache.iotdb.tsfile.read.reader.page.ValuePageReader;
-import org.apache.iotdb.tsfile.utils.Pair;
-import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
import org.slf4j.Logger;
@@ -57,13 +39,10 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.IOException;
-import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collection;
-import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
-import java.util.Map;
import java.util.Set;
/**
@@ -301,148 +280,6 @@ public class CompactionUtils {
}
}
- public static boolean validateTsFileResources(
- TsFileManager manager, String storageGroupName, long timePartition)
throws IOException {
- List<TsFileResource> resources =
-
manager.getOrCreateSequenceListByTimePartition(timePartition).getArrayList();
-
- // deviceID -> <TsFileResource, last end time>
- Map<String, Pair<TsFileResource, Long>> lastEndTimeMap = new HashMap<>();
- for (TsFileResource resource : resources) {
- DeviceTimeIndex timeIndex;
- if (resource.getTimeIndexType() != 1) {
- // if time index is not device time index, then deserialize it from
resource file
- timeIndex = resource.buildDeviceTimeIndex();
- } else {
- timeIndex = (DeviceTimeIndex) resource.getTimeIndex();
- }
- Set<String> devices = timeIndex.getDevices();
- for (String device : devices) {
- long currentStartTime = timeIndex.getStartTime(device);
- long currentEndTime = timeIndex.getEndTime(device);
- Pair<TsFileResource, Long> lastDeviceInfo =
- lastEndTimeMap.computeIfAbsent(device, x -> new Pair<>(null,
Long.MIN_VALUE));
- long lastEndTime = lastDeviceInfo.right;
- if (lastEndTime >= currentStartTime) {
- logger.error(
- "{} Device {} is overlapped between {} and {}, "
- + "end time in {} is {}, start time in {} is {}",
- storageGroupName,
- device,
- lastDeviceInfo.left,
- resource,
- lastDeviceInfo.left,
- lastEndTime,
- resource,
- currentStartTime);
- return false;
- }
- lastDeviceInfo.left = resource;
- lastDeviceInfo.right = currentEndTime;
- lastEndTimeMap.put(device, lastDeviceInfo);
- }
- }
- return true;
- }
-
- /**
- * Validate TsFiles by reading them sequentially. This method should be fast
because the read is
- * sequential.
- *
- * @param tsFileResourceList the tsfiles to be checked
- * @return true if all tsfiles are valid, false if any of the tsfiles is
invalid
- */
- public static boolean validateTsFiles(List<TsFileResource>
tsFileResourceList) {
- for (TsFileResource tsFileResource : tsFileResourceList) {
- if (!validateSingleTsFiles(tsFileResource)) {
- return false;
- }
- }
- return true;
- }
-
- @SuppressWarnings({"squid:S6541", "squid:S3776", "squid:S1481"}) // do not
warn about brain method
- public static boolean validateSingleTsFiles(TsFileResource resource) {
- try (TsFileSequenceReader reader = new
TsFileSequenceReader(resource.getTsFilePath())) {
- reader.readHeadMagic();
- reader.readTailMagic();
-
- reader.position((long) TSFileConfig.MAGIC_STRING.getBytes().length + 1);
- List<long[]> timeBatch = new ArrayList<>();
- int pageIndex = 0;
- byte marker;
- while ((marker = reader.readMarker()) != MetaMarker.SEPARATOR) {
- switch (marker) {
- case MetaMarker.CHUNK_HEADER:
- case MetaMarker.TIME_CHUNK_HEADER:
- case MetaMarker.VALUE_CHUNK_HEADER:
- case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER:
- case MetaMarker.ONLY_ONE_PAGE_TIME_CHUNK_HEADER:
- case MetaMarker.ONLY_ONE_PAGE_VALUE_CHUNK_HEADER:
- ChunkHeader header = reader.readChunkHeader(marker);
- if (header.getDataSize() == 0) {
- // empty value chunk
- break;
- }
- Decoder defaultTimeDecoder =
- Decoder.getDecoderByType(
-
TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()),
- TSDataType.INT64);
- Decoder valueDecoder =
- Decoder.getDecoderByType(header.getEncodingType(),
header.getDataType());
- int dataSize = header.getDataSize();
- pageIndex = 0;
- if (header.getDataType() == TSDataType.VECTOR) {
- timeBatch.clear();
- }
- while (dataSize > 0) {
- valueDecoder.reset();
- PageHeader pageHeader =
- reader.readPageHeader(
- header.getDataType(),
- (header.getChunkType() & 0x3F) ==
MetaMarker.CHUNK_HEADER);
- ByteBuffer pageData = reader.readPage(pageHeader,
header.getCompressionType());
- if ((header.getChunkType() & TsFileConstant.TIME_COLUMN_MASK)
- == TsFileConstant.TIME_COLUMN_MASK) { // Time Chunk
- TimePageReader timePageReader =
- new TimePageReader(pageHeader, pageData,
defaultTimeDecoder);
- timeBatch.add(timePageReader.getNextTimeBatch());
- } else if ((header.getChunkType() &
TsFileConstant.VALUE_COLUMN_MASK)
- == TsFileConstant.VALUE_COLUMN_MASK) { // Value Chunk
- ValuePageReader valuePageReader =
- new ValuePageReader(pageHeader, pageData,
header.getDataType(), valueDecoder);
- TsPrimitiveType[] valueBatch =
- valuePageReader.nextValueBatch(timeBatch.get(pageIndex));
- } else { // NonAligned Chunk
- PageReader pageReader =
- new PageReader(
- pageData, header.getDataType(), valueDecoder,
defaultTimeDecoder, null);
- BatchData batchData = pageReader.getAllSatisfiedPageData();
- }
- pageIndex++;
- dataSize -= pageHeader.getSerializedPageSize();
- }
- break;
- case MetaMarker.CHUNK_GROUP_HEADER:
- ChunkGroupHeader chunkGroupHeader = reader.readChunkGroupHeader();
- break;
- case MetaMarker.OPERATION_INDEX_RANGE:
- reader.readPlanIndex();
- break;
- default:
- MetaMarker.handleUnexpectedMarker(marker);
- }
- }
- for (String device : reader.getAllDevices()) {
- Map<String, List<ChunkMetadata>> seriesMetaData =
reader.readChunkMetadataInDevice(device);
- }
- } catch (Exception e) {
- logger.error("Meets error when validating TsFile {}, ",
resource.getTsFilePath(), e);
- return false;
- }
- return true;
- }
-
public static void deleteSourceTsFileAndUpdateFileMetrics(
List<TsFileResource> sourceSeqResourceList, List<TsFileResource>
sourceUnseqResourceList) {
deleteSourceTsFileAndUpdateFileMetrics(sourceSeqResourceList, true);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/CompactionValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/CompactionValidator.java
deleted file mode 100644
index 684554bc3b4..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/CompactionValidator.java
+++ /dev/null
@@ -1,51 +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.execute.utils.validator;
-
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.constant.CompactionValidationLevel;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-
-import java.io.IOException;
-import java.util.List;
-
-public interface CompactionValidator {
- boolean validateCompaction(
- TsFileManager manager,
- List<TsFileResource> targetTsFileList,
- String storageGroupName,
- long timePartition,
- boolean isInnerUnSequenceSpaceTask)
- throws IOException;
-
- static CompactionValidator getInstance() {
- CompactionValidationLevel level =
-
IoTDBDescriptor.getInstance().getConfig().getCompactionValidationLevel();
- switch (level) {
- case NONE:
- return NoneCompactionValidator.getInstance();
- case RESOURCE_ONLY:
- return ResourceOnlyCompactionValidator.getInstance();
- default:
- return ResourceAndTsfileCompactionValidator.getInstance();
- }
- }
-}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceAndTsfileCompactionValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceAndTsfileCompactionValidator.java
deleted file mode 100644
index 6c955e7c868..00000000000
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceAndTsfileCompactionValidator.java
+++ /dev/null
@@ -1,57 +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.execute.utils.validator;
-
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-
-import java.io.IOException;
-import java.util.List;
-
-@SuppressWarnings("squid:S6548")
-public class ResourceAndTsfileCompactionValidator implements
CompactionValidator {
-
- private ResourceAndTsfileCompactionValidator() {}
-
- public static ResourceAndTsfileCompactionValidator getInstance() {
- return ResourceAndTsfileCompactionValidatorHolder.INSTANCE;
- }
-
- @Override
- public boolean validateCompaction(
- TsFileManager manager,
- List<TsFileResource> targetTsFileList,
- String storageGroupName,
- long timePartition,
- boolean isInnerUnSequenceSpaceTask)
- throws IOException {
- if (isInnerUnSequenceSpaceTask) {
- return CompactionUtils.validateTsFiles(targetTsFileList);
- }
- return CompactionUtils.validateTsFileResources(manager, storageGroupName,
timePartition)
- && CompactionUtils.validateTsFiles(targetTsFileList);
- }
-
- private static class ResourceAndTsfileCompactionValidatorHolder {
- private static final ResourceAndTsfileCompactionValidator INSTANCE =
- new ResourceAndTsfileCompactionValidator();
- }
-}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/TsFileResourceUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/TsFileResourceUtils.java
new file mode 100644
index 00000000000..aaab876bf07
--- /dev/null
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/TsFileResourceUtils.java
@@ -0,0 +1,378 @@
+/*
+ * 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.utils;
+
+import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import
org.apache.iotdb.db.storageengine.dataregion.tsfile.timeindex.DeviceTimeIndex;
+import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
+import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
+import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
+import org.apache.iotdb.tsfile.encoding.decoder.Decoder;
+import org.apache.iotdb.tsfile.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.MetaMarker;
+import org.apache.iotdb.tsfile.file.header.ChunkGroupHeader;
+import org.apache.iotdb.tsfile.file.header.ChunkHeader;
+import org.apache.iotdb.tsfile.file.header.PageHeader;
+import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata;
+import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
+import org.apache.iotdb.tsfile.read.common.BatchData;
+import org.apache.iotdb.tsfile.read.reader.page.PageReader;
+import org.apache.iotdb.tsfile.read.reader.page.TimePageReader;
+import org.apache.iotdb.tsfile.read.reader.page.ValuePageReader;
+import org.apache.iotdb.tsfile.utils.Pair;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.LinkedList;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+public class TsFileResourceUtils {
+ private static final Logger logger =
LoggerFactory.getLogger(TsFileResourceUtils.class);
+ private static final String VALIDATE_FAILED = "validate failed,";
+
+ public static boolean validateTsFileResourceCorrectness(TsFileResource
resource) {
+ DeviceTimeIndex timeIndex;
+ try {
+ if (resource.getTimeIndexType() != 1) {
+ // if time index is not device time index, then deserialize it from
resource file
+ timeIndex = resource.buildDeviceTimeIndex();
+ } else {
+ timeIndex = (DeviceTimeIndex) resource.getTimeIndex();
+ }
+ if (timeIndex == null) {
+ logger.error("{} {} time index is null", resource.getTsFilePath(),
VALIDATE_FAILED);
+ return false;
+ }
+ Set<String> devices = timeIndex.getDevices();
+ if (devices.isEmpty()) {
+ logger.error("{} {} empty resource", resource.getTsFilePath(),
VALIDATE_FAILED);
+ return false;
+ }
+ for (String device : devices) {
+ long startTime = timeIndex.getStartTime(device);
+ long endTime = timeIndex.getEndTime(device);
+ if (startTime == Long.MAX_VALUE) {
+ logger.error(
+ "{} {} the start time of {} is {}",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED,
+ device,
+ Long.MAX_VALUE);
+ return false;
+ }
+ if (endTime == Long.MIN_VALUE) {
+ logger.error(
+ "{} {} the end time of {} is {}",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED,
+ device,
+ Long.MIN_VALUE);
+ return false;
+ }
+ if (startTime > endTime) {
+ logger.error(
+ "{} {} the start time of {} is greater than end time",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED,
+ device);
+ return false;
+ }
+ }
+ } catch (IOException e) {
+ logger.error("meet error when validate .resource file:{},e",
resource.getTsFilePath());
+ return false;
+ }
+ return true;
+ }
+
+ public static boolean validateTsFileDataCorrectness(TsFileResource resource)
{
+ try (TsFileSequenceReader reader = new
TsFileSequenceReader(resource.getTsFilePath())) {
+ if (!reader.isComplete()) {
+ logger.error("{} {} illegal tsfile", resource.getTsFilePath(),
VALIDATE_FAILED);
+ return false;
+ }
+
+ Map<Long, IChunkMetadata> chunkMetadataMap = getChunkMetadata(reader);
+ if (chunkMetadataMap.isEmpty()) {
+ logger.error(
+ "{} {} there is no data in the file", resource.getTsFilePath(),
VALIDATE_FAILED);
+ return false;
+ }
+
+ List<long[]> alignedTimeBatch = new ArrayList<>();
+ reader.position((long) TSFileConfig.MAGIC_STRING.getBytes().length + 1);
+ int pageIndex = 0;
+ byte marker;
+ while ((marker = reader.readMarker()) != MetaMarker.SEPARATOR) {
+ switch (marker) {
+ case MetaMarker.CHUNK_HEADER:
+ case MetaMarker.TIME_CHUNK_HEADER:
+ case MetaMarker.VALUE_CHUNK_HEADER:
+ case MetaMarker.ONLY_ONE_PAGE_CHUNK_HEADER:
+ case MetaMarker.ONLY_ONE_PAGE_TIME_CHUNK_HEADER:
+ case MetaMarker.ONLY_ONE_PAGE_VALUE_CHUNK_HEADER:
+ long chunkOffset = reader.position();
+ ChunkHeader header = reader.readChunkHeader(marker);
+ IChunkMetadata chunkMetadata = chunkMetadataMap.get(chunkOffset -
Byte.BYTES);
+ if
(!chunkMetadata.getMeasurementUid().equals(header.getMeasurementID())) {
+ logger.error(
+ "{} chunk start offset is inconsistent with the value in the
metadata.",
+ VALIDATE_FAILED);
+ return false;
+ }
+
+ // empty value chunk
+ int dataSize = header.getDataSize();
+ if (dataSize == 0) {
+ break;
+ }
+
+ boolean isHasStatistic = (header.getChunkType() & 0x3F) ==
MetaMarker.CHUNK_HEADER;
+ Decoder defaultTimeDecoder =
+ Decoder.getDecoderByType(
+
TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()),
+ TSDataType.INT64);
+ Decoder valueDecoder =
+ Decoder.getDecoderByType(header.getEncodingType(),
header.getDataType());
+
+ pageIndex = 0;
+
+ if (header.getDataType() == TSDataType.VECTOR) {
+ alignedTimeBatch.clear();
+ }
+ LinkedList<Long> lastNoAlignedPageTimeStamps = new LinkedList<>();
+ while (dataSize > 0) {
+ valueDecoder.reset();
+ PageHeader pageHeader =
reader.readPageHeader(header.getDataType(), isHasStatistic);
+ ByteBuffer pageData = reader.readPage(pageHeader,
header.getCompressionType());
+
+ if ((header.getChunkType() & TsFileConstant.TIME_COLUMN_MASK)
+ == TsFileConstant.TIME_COLUMN_MASK) { // Time Chunk
+ TimePageReader timePageReader =
+ new TimePageReader(pageHeader, pageData,
defaultTimeDecoder);
+ long[] pageTimestamps = timePageReader.getNextTimeBatch();
+ long pageHeaderStartTime =
+ isHasStatistic ? pageHeader.getStartTime() :
chunkMetadata.getStartTime();
+ long pageHeaderEndTime =
+ isHasStatistic ? pageHeader.getEndTime() :
chunkMetadata.getEndTime();
+ if (!validateTimeFrame(
+ alignedTimeBatch,
+ pageTimestamps,
+ pageHeaderStartTime,
+ pageHeaderEndTime,
+ resource)) {
+ return false;
+ }
+ alignedTimeBatch.add(pageTimestamps);
+ } else if ((header.getChunkType() &
TsFileConstant.VALUE_COLUMN_MASK)
+ == TsFileConstant.VALUE_COLUMN_MASK) { // Value Chunk
+ ValuePageReader valuePageReader =
+ new ValuePageReader(pageHeader, pageData,
header.getDataType(), valueDecoder);
+
valuePageReader.nextValueBatch(alignedTimeBatch.get(pageIndex));
+ } else { // NonAligned Chunk
+ PageReader pageReader =
+ new PageReader(
+ pageData, header.getDataType(), valueDecoder,
defaultTimeDecoder, null);
+ BatchData batchData = pageReader.getAllSatisfiedPageData();
+ long pageHeaderStartTime =
+ isHasStatistic ? pageHeader.getStartTime() :
chunkMetadata.getStartTime();
+ long pageHeaderEndTime =
+ isHasStatistic ? pageHeader.getEndTime() :
chunkMetadata.getEndTime();
+ long pageStartTime = Long.MAX_VALUE;
+ long previousTime = Long.MIN_VALUE;
+
+ while (batchData.hasCurrent()) {
+ long currentTime = batchData.currentTime();
+ if (!lastNoAlignedPageTimeStamps.isEmpty()
+ && currentTime <= lastNoAlignedPageTimeStamps.getLast())
{
+ logger.error(
+ "{} {} time ranges overlap between pages.",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ if (currentTime <= previousTime) {
+ logger.error(
+ "{} {} the timestamp in the page is repeated or not
incremental.",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ pageStartTime = Math.min(pageStartTime, currentTime);
+ previousTime = currentTime;
+ lastNoAlignedPageTimeStamps.add(currentTime);
+ batchData.next();
+ }
+ if (pageHeaderStartTime != pageStartTime) {
+ logger.error(
+ "{} {} the start time in page is different from that in
page header.",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ if (pageHeaderEndTime != previousTime) {
+ logger.error(
+ "{} {} the end time in page is different from that in
page header.",
+ resource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ }
+ pageIndex++;
+ dataSize -= pageHeader.getSerializedPageSize();
+ }
+
+ break;
+ case MetaMarker.CHUNK_GROUP_HEADER:
+ ChunkGroupHeader chunkGroupHeader = reader.readChunkGroupHeader();
+ if (chunkGroupHeader.getDeviceID() == null
+ || chunkGroupHeader.getDeviceID().isEmpty()) {
+ logger.error(
+ "{} {} device id is null or empty.",
resource.getTsFilePath(), VALIDATE_FAILED);
+ return false;
+ }
+ break;
+ case MetaMarker.OPERATION_INDEX_RANGE:
+ reader.readPlanIndex();
+ break;
+ default:
+ MetaMarker.handleUnexpectedMarker(marker);
+ }
+ }
+
+ } catch (Exception e) {
+ logger.error("Meets error when validating TsFile {}, ",
resource.getTsFilePath(), e);
+ return false;
+ }
+ return true;
+ }
+
+ private static boolean validateTimeFrame(
+ List<long[]> timeBatch,
+ long[] pageTimestamps,
+ long pageHeaderStartTime,
+ long pageHeaderEndTime,
+ TsFileResource tsFileResource) {
+ if (pageHeaderStartTime != pageTimestamps[0]) {
+ logger.error(
+ "{} {} the start time in page is different from that in page
header.",
+ tsFileResource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ if (pageHeaderEndTime != pageTimestamps[pageTimestamps.length - 1]) {
+ logger.error(
+ "{} {} the end time in page is different from that in page header.",
+ tsFileResource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ for (int i = 0; i < pageTimestamps.length - 1; i++) {
+ if (pageTimestamps[i + 1] <= pageTimestamps[i]) {
+ logger.error(
+ "{} {} the timestamp in the page is repeated or not incremental.",
+ tsFileResource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ }
+
+ if (timeBatch.size() >= 1) {
+ long[] lastPageTimes = timeBatch.get(timeBatch.size() - 1);
+ if (lastPageTimes[lastPageTimes.length - 1] >= pageTimestamps[0]) {
+ logger.error(
+ "{} {} time ranges overlap between pages.",
+ tsFileResource.getTsFilePath(),
+ VALIDATE_FAILED);
+ return false;
+ }
+ }
+ return true;
+ }
+
+ public static Map<Long, IChunkMetadata>
getChunkMetadata(TsFileSequenceReader reader)
+ throws IOException {
+ Map<Long, IChunkMetadata> offset2ChunkMetadata = new HashMap<>();
+ Map<String, List<TimeseriesMetadata>> device2Metadata =
reader.getAllTimeseriesMetadata(true);
+ for (Map.Entry<String, List<TimeseriesMetadata>> entry :
device2Metadata.entrySet()) {
+ for (TimeseriesMetadata timeseriesMetadata : entry.getValue()) {
+ for (IChunkMetadata chunkMetadata :
timeseriesMetadata.getChunkMetadataList()) {
+ offset2ChunkMetadata.put(chunkMetadata.getOffsetOfChunkHeader(),
chunkMetadata);
+ }
+ }
+ }
+ return offset2ChunkMetadata;
+ }
+
+ public static boolean
validateTsFileResourcesHasNoOverlap(List<TsFileResource> resources) {
+ try {
+ // deviceID -> <TsFileResource, last end time>
+ Map<String, Pair<TsFileResource, Long>> lastEndTimeMap = new HashMap<>();
+ for (TsFileResource resource : resources) {
+ DeviceTimeIndex timeIndex;
+ if (resource.getTimeIndexType() != 1) {
+ // if time index is not device time index, then deserialize it from
resource file
+ timeIndex = resource.buildDeviceTimeIndex();
+ } else {
+ timeIndex = (DeviceTimeIndex) resource.getTimeIndex();
+ }
+ if (timeIndex == null) {
+ return false;
+ }
+ Set<String> devices = timeIndex.getDevices();
+ for (String device : devices) {
+ long currentStartTime = timeIndex.getStartTime(device);
+ long currentEndTime = timeIndex.getEndTime(device);
+ Pair<TsFileResource, Long> lastDeviceInfo =
+ lastEndTimeMap.computeIfAbsent(device, x -> new Pair<>(null,
Long.MIN_VALUE));
+ long lastEndTime = lastDeviceInfo.right;
+ if (lastEndTime >= currentStartTime) {
+ logger.error(
+ "Device {} is overlapped between {} and {}, "
+ + "end time in {} is {}, start time in {} is {}",
+ device,
+ lastDeviceInfo.left,
+ resource,
+ lastDeviceInfo.left,
+ lastEndTime,
+ resource,
+ currentStartTime);
+ return false;
+ }
+ lastDeviceInfo.left = resource;
+ lastDeviceInfo.right = currentEndTime;
+ lastEndTimeMap.put(device, lastDeviceInfo);
+ }
+ }
+ return true;
+ } catch (IOException e) {
+ return true;
+ }
+ }
+}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceAndDataValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceAndDataValidator.java
new file mode 100644
index 00000000000..ae366a7af79
--- /dev/null
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceAndDataValidator.java
@@ -0,0 +1,58 @@
+/*
+ * 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.utils.validate;
+
+import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import org.apache.iotdb.db.storageengine.dataregion.utils.TsFileResourceUtils;
+
+import java.util.List;
+
+public class TsFileResourceAndDataValidator implements TsFileValidator {
+
+ @Override
+ public boolean validateTsFile(TsFileResource resource) {
+ return TsFileResourceUtils.validateTsFileResourceCorrectness(resource)
+ && TsFileResourceUtils.validateTsFileDataCorrectness(resource);
+ }
+
+ @Override
+ public boolean validateTsFiles(List<TsFileResource> resourceList) {
+ for (TsFileResource resource : resourceList) {
+ if (!validateTsFile(resource)) {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ @Override
+ public boolean validateTsFilesIsHasNoOverlap(List<TsFileResource>
resourceList) {
+ return
TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(resourceList);
+ }
+
+ public static TsFileResourceAndDataValidator getInstance() {
+ return TsFileDataCorrectnessValidatorHolder.INSTANCE;
+ }
+
+ private static class TsFileDataCorrectnessValidatorHolder {
+ private static final TsFileResourceAndDataValidator INSTANCE =
+ new TsFileResourceAndDataValidator();
+ }
+}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceOnlyCompactionValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceValidator.java
similarity index 53%
rename from
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceOnlyCompactionValidator.java
rename to
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceValidator.java
index 953d13a1b77..07373fc15e4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/ResourceOnlyCompactionValidator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileResourceValidator.java
@@ -17,40 +17,43 @@
* under the License.
*/
-package
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.validator;
+package org.apache.iotdb.db.storageengine.dataregion.utils.validate;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
-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.utils.TsFileResourceUtils;
-import java.io.IOException;
import java.util.List;
@SuppressWarnings("squid:S6548")
-public class ResourceOnlyCompactionValidator implements CompactionValidator {
+public class TsFileResourceValidator implements TsFileValidator {
- private ResourceOnlyCompactionValidator() {}
+ private TsFileResourceValidator() {}
- public static ResourceOnlyCompactionValidator getInstance() {
- return ResourceOnlyCompactionValidatorHolder.INSTANCE;
+ @Override
+ public boolean validateTsFile(TsFileResource resource) {
+ return TsFileResourceUtils.validateTsFileResourceCorrectness(resource);
}
@Override
- public boolean validateCompaction(
- TsFileManager manager,
- List<TsFileResource> targetTsFileList,
- String storageGroupName,
- long timePartition,
- boolean isInnerUnSequenceSpaceTask)
- throws IOException {
- if (isInnerUnSequenceSpaceTask) {
- return true;
+ public boolean validateTsFiles(List<TsFileResource> resourceList) {
+ for (TsFileResource resource : resourceList) {
+ if (!validateTsFile(resource)) {
+ return false;
+ }
}
- return CompactionUtils.validateTsFileResources(manager, storageGroupName,
timePartition);
+ return true;
+ }
+
+ @Override
+ public boolean validateTsFilesIsHasNoOverlap(List<TsFileResource>
resourceList) {
+ return
TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(resourceList);
+ }
+
+ public static TsFileResourceValidator getInstance() {
+ return ResourceOnlyCompactionValidatorHolder.INSTANCE;
}
private static class ResourceOnlyCompactionValidatorHolder {
- private static final ResourceOnlyCompactionValidator INSTANCE =
- new ResourceOnlyCompactionValidator();
+ private static final TsFileResourceValidator INSTANCE = new
TsFileResourceValidator();
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/NoneCompactionValidator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileValidator.java
similarity index 53%
rename from
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/NoneCompactionValidator.java
rename to
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileValidator.java
index 33c5ef91364..b3d855cb203 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/validator/NoneCompactionValidator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/utils/validate/TsFileValidator.java
@@ -17,33 +17,31 @@
* under the License.
*/
-package
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.validator;
+package org.apache.iotdb.db.storageengine.dataregion.utils.validate;
-import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import java.util.List;
-@SuppressWarnings("squid:S6548")
-public class NoneCompactionValidator implements CompactionValidator {
+public interface TsFileValidator {
- private NoneCompactionValidator() {}
+ /** validate that a single .resource file is correct */
+ boolean validateTsFile(TsFileResource resource);
- public static NoneCompactionValidator getInstance() {
- return NoneCompactionValidatorHolder.INSTANCE;
- }
+ /** validate that multiple .resource files are correct */
+ boolean validateTsFiles(List<TsFileResource> resourceList);
- @Override
- public boolean validateCompaction(
- TsFileManager manager,
- List<TsFileResource> targetTsFileList,
- String storageGroupName,
- long timePartition,
- boolean isInnerUnSequenceSpaceTask) {
- return true;
- }
+ /** */
+ boolean validateTsFilesIsHasNoOverlap(List<TsFileResource> resourceList);
- private static class NoneCompactionValidatorHolder {
- private static final NoneCompactionValidator INSTANCE = new
NoneCompactionValidator();
+ static TsFileValidator getInstance() {
+ boolean enableTsFileValidation =
+ IoTDBDescriptor.getInstance().getConfig().isEnableTsFileValidation();
+ if (enableTsFileValidation) {
+ return TsFileResourceAndDataValidator.getInstance();
+ } else {
+ return TsFileResourceValidator.getInstance();
+ }
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionValidationTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionValidationTest.java
index 6d9d403df39..712ec4c5ca3 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionValidationTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CompactionValidationTest.java
@@ -19,8 +19,8 @@
package org.apache.iotdb.db.storageengine.dataregion.compaction;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import org.apache.iotdb.db.storageengine.dataregion.utils.TsFileResourceUtils;
import org.apache.iotdb.db.utils.constant.TestConstant;
import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
import org.apache.iotdb.tsfile.enums.TSDataType;
@@ -38,13 +38,11 @@ import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
-import org.mockito.Mockito;
import java.io.File;
import java.io.IOException;
import java.io.RandomAccessFile;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
public class CompactionValidationTest {
@@ -137,22 +135,18 @@ public class CompactionValidationTest {
public void testSingleCompleteFile() {
String path = dir + File.separator + "test.tsfile";
writeOneFile(path);
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
-
Assert.assertTrue(CompactionUtils.validateTsFiles(Collections.singletonList(mockTsFile)));
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+
Assert.assertTrue(TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile));
}
@Test
public void testMultiCompleteFile() {
- List<TsFileResource> resources = new ArrayList<>();
for (int i = 0; i < 10; i++) {
String path = dir + File.separator + "test" + i + ".tsfile";
writeOneFile(path);
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
- resources.add(mockTsFile);
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+
Assert.assertTrue(TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile));
}
- Assert.assertTrue(CompactionUtils.validateTsFiles(resources));
}
@Test
@@ -163,14 +157,12 @@ public class CompactionValidationTest {
randomAccessFile.seek(1024);
randomAccessFile.write(new byte[] {1, 2, 3, 4, 5, 6, 7, 8});
randomAccessFile.close();
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
-
Assert.assertFalse(CompactionUtils.validateTsFiles(Collections.singletonList(mockTsFile)));
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+
Assert.assertFalse(TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile));
}
@Test
public void testMultiUncompletedFiles() throws IOException {
- List<TsFileResource> resources = new ArrayList<>();
for (int i = 0; i < 10; i++) {
String path = dir + File.separator + "test" + i + ".tsfile";
writeOneFile(path);
@@ -178,16 +170,13 @@ public class CompactionValidationTest {
randomAccessFile.seek(1024);
randomAccessFile.write(new byte[] {1, 2, 3, 4, 5, 6, 7, 8});
randomAccessFile.close();
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
- resources.add(mockTsFile);
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+
Assert.assertFalse(TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile));
}
- Assert.assertFalse(CompactionUtils.validateTsFiles(resources));
}
@Test // broken in chunk
public void testOneUncompletedInMultiCompletedFiles1() throws IOException {
- List<TsFileResource> resources = new ArrayList<>();
for (int i = 0; i < 10; i++) {
String path = dir + File.separator + "test" + i + ".tsfile";
writeOneFile(path);
@@ -197,16 +186,13 @@ public class CompactionValidationTest {
randomAccessFile.write(new byte[] {1, 2, 3, 4, 5, 6, 7, 8});
randomAccessFile.close();
}
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
- resources.add(mockTsFile);
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+ TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile);
}
- Assert.assertFalse(CompactionUtils.validateTsFiles(resources));
}
@Test // broken in metadata
public void testOneUncompletedInMultiCompletedFiles2() throws IOException {
- List<TsFileResource> resources = new ArrayList<>();
for (int i = 0; i < 10; i++) {
String path = dir + File.separator + "test" + i + ".tsfile";
writeOneFile(path);
@@ -216,10 +202,8 @@ public class CompactionValidationTest {
randomAccessFile.write(new byte[] {1, 2, 3, 4, 5, 6, 7, 8});
randomAccessFile.close();
}
- TsFileResource mockTsFile = Mockito.mock(TsFileResource.class);
- Mockito.when(mockTsFile.getTsFilePath()).thenReturn(path);
- resources.add(mockTsFile);
+ TsFileResource mockTsFile = new TsFileResource(new File(path));
+ TsFileResourceUtils.validateTsFileDataCorrectness(mockTsFile);
}
- Assert.assertFalse(CompactionUtils.validateTsFiles(resources));
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CrossSpaceCompactionWithUnusualCasesTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CrossSpaceCompactionWithUnusualCasesTest.java
index 923d56c47a5..7469536d023 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CrossSpaceCompactionWithUnusualCasesTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/CrossSpaceCompactionWithUnusualCasesTest.java
@@ -25,10 +25,6 @@ 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.exception.StorageEngineException;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.constant.CompactionValidationLevel;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer.impl.FastCompactionPerformer;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.AbstractCompactionTask;
-import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.CrossSpaceCompactionTask;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.impl.RewriteCrossSpaceCompactionSelector;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.utils.CrossCompactionTaskResource;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.utils.TsFileGeneratorUtils;
@@ -66,9 +62,7 @@ public class CrossSpaceCompactionWithUnusualCasesTest extends
AbstractCompaction
public void setUp()
throws IOException, WriteProcessException, MetadataException,
InterruptedException {
super.setUp();
- IoTDBDescriptor.getInstance()
- .getConfig()
-
.setCompactionValidationLevel(CompactionValidationLevel.RESOURCE_AND_TSFILE);
+ IoTDBDescriptor.getInstance().getConfig().setEnableTsFileValidation(true);
IoTDBDescriptor.getInstance().getConfig().setMinCrossCompactionUnseqFileLevel(0);
TSFileDescriptor.getInstance().getConfig().setMaxNumberOfPointsInPage(10);
TSFileDescriptor.getInstance().getConfig().setMaxDegreeOfIndexNode(3);
@@ -124,20 +118,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
selector.selectCrossSpaceTask(seqResources, unseqResources);
Assert.assertEquals(1, result.size());
Assert.assertEquals(3, result.get(0).getTotalFileNums());
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -196,20 +176,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertEquals(3, result.get(0).getTotalFileNums());
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource2));
Assert.assertTrue(result.get(0).getSeqFiles().contains(seqTsFileResource3));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -279,20 +245,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource2));
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource4));
Assert.assertTrue(result.get(0).getSeqFiles().contains(seqTsFileResource3));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -372,20 +324,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource2));
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource4));
Assert.assertTrue(result.get(0).getSeqFiles().contains(seqTsFileResource5));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -467,20 +405,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertEquals(1, result.size());
Assert.assertEquals(5, result.get(0).getTotalFileNums());
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource2));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -528,20 +452,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
selector.selectCrossSpaceTask(seqResources, unseqResources);
Assert.assertEquals(1, result.size());
Assert.assertEquals(2, result.get(0).getTotalFileNums());
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -631,20 +541,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource2));
Assert.assertFalse(result.get(0).getSeqFiles().contains(seqTsFileResource4));
Assert.assertTrue(result.get(0).getSeqFiles().contains(seqTsFileResource5));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
@@ -873,20 +769,6 @@ public class CrossSpaceCompactionWithUnusualCasesTest
extends AbstractCompaction
Assert.assertEquals(2, result.get(0).getSeqFiles().size());
Assert.assertEquals(seqTsFileResource1,
result.get(0).getSeqFiles().get(0));
Assert.assertEquals(seqTsFileResource4,
result.get(0).getSeqFiles().get(1));
-
- // execution
- FastCompactionPerformer performer = new FastCompactionPerformer(true);
- AbstractCompactionTask task =
- new CrossSpaceCompactionTask(
- 0,
- tsFileManager,
- result.get(0).getSeqFiles(),
- result.get(0).getUnseqFiles(),
- performer,
- 0,
- 0);
- Assert.assertTrue(task.start());
- validateSeqFiles(true);
}
@Test
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastCrossCompactionPerformerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastCrossCompactionPerformerTest.java
index b3b7a58b1db..e88508357c1 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastCrossCompactionPerformerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastCrossCompactionPerformerTest.java
@@ -37,6 +37,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.Compacti
import
org.apache.iotdb.db.storageengine.dataregion.compaction.utils.CompactionFileGeneratorUtils;
import
org.apache.iotdb.db.storageengine.dataregion.read.control.FileReaderManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import org.apache.iotdb.db.storageengine.dataregion.utils.TsFileResourceUtils;
import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
import org.apache.iotdb.db.tools.validate.TsFileValidationTool;
import org.apache.iotdb.db.utils.EnvironmentUtils;
@@ -3934,7 +3935,8 @@ public class FastCrossCompactionPerformerTest extends
AbstractCompactionTest {
targetResources.get(3).degradeTimeIndex();
targetResources.get(2).degradeTimeIndex();
Assert.assertTrue(
- CompactionUtils.validateTsFileResources(tsFileManager,
COMPACTION_TEST_SG, 0));
+ TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(0).getArrayList()));
List<String> deviceIdList = new ArrayList<>();
deviceIdList.add(COMPACTION_TEST_SG + PATH_SEPARATOR + "d0");
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastInnerCompactionPerformerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastInnerCompactionPerformerTest.java
index f6005213f44..675375ab399 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastInnerCompactionPerformerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/FastInnerCompactionPerformerTest.java
@@ -949,7 +949,7 @@ public class FastInnerCompactionPerformerTest extends
AbstractCompactionTest {
InnerSpaceCompactionTask task =
new InnerSpaceCompactionTask(
0, tsFileManager, unseqResources, false, new
FastCompactionPerformer(false), 0);
- Assert.assertTrue(task.start());
+ Assert.assertFalse(task.start());
Assert.assertEquals(0,
FileReaderManager.getInstance().getClosedFileReaderMap().size());
Assert.assertEquals(0,
FileReaderManager.getInstance().getUnclosedFileReaderMap().size());
validateSeqFiles(true);
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/TsFileValidationCorrectnessTests.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/TsFileValidationCorrectnessTests.java
new file mode 100644
index 00000000000..95426d2c184
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/TsFileValidationCorrectnessTests.java
@@ -0,0 +1,296 @@
+/*
+ * 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;
+
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.utils.CompactionTestFileWriter;
+import
org.apache.iotdb.db.storageengine.dataregion.compaction.utils.TsFileGeneratorUtils;
+import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
+import
org.apache.iotdb.db.storageengine.dataregion.utils.validate.TsFileValidator;
+import org.apache.iotdb.db.utils.constant.TestConstant;
+import org.apache.iotdb.tsfile.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.read.common.TimeRange;
+import org.apache.iotdb.tsfile.write.chunk.AlignedChunkWriterImpl;
+import org.apache.iotdb.tsfile.write.chunk.ChunkWriterImpl;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
+import org.apache.iotdb.tsfile.write.schema.VectorMeasurementSchema;
+
+import org.apache.commons.io.FileUtils;
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.Collections;
+
+public class TsFileValidationCorrectnessTests {
+ final String dir = TestConstant.OUTPUT_DATA_DIR + "test-validation";
+
+ @Before
+ public void setUp() throws IOException {
+ FileUtils.forceMkdir(new File(dir));
+ }
+
+ @After
+ public void tearDown() throws IOException {
+ File[] files = new File(dir).listFiles();
+ if (files != null) {
+ for (File f : files) {
+ FileUtils.delete(f);
+ }
+ }
+
+ FileUtils.forceDelete(new File(dir));
+ }
+ // 1. empty tsfile
+ @Test
+ public void testTsFileHasNoData() throws IOException {
+ TsFileResource tsFileResource =
+ new TsFileResource(new File(dir + File.separator + "test1.tsfile"));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.endFile();
+ }
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertFalse(success);
+ }
+
+ @Test
+ public void testAlignedTsFileHasOnePageData() throws IOException {
+ String path = dir + File.separator + "test2.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ TsFileGeneratorUtils.generateSingleAlignedSeriesFile(
+ "d1",
+ Collections.singletonList("s1"),
+ new TimeRange[] {new TimeRange(1, 100)},
+ TSEncoding.PLAIN,
+ CompressionType.SNAPPY,
+ path);
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 100);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+
+ @Test
+ public void testAlignedTsFileHasManyPage() throws IOException {
+ String path = dir + File.separator + "test3.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ TsFileGeneratorUtils.generateSingleAlignedSeriesFile(
+ "d1",
+ Collections.singletonList("s1"),
+ new TimeRange[] {new TimeRange(1, 100), new TimeRange(22, 110)},
+ TSEncoding.PLAIN,
+ CompressionType.SNAPPY,
+ path);
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 110);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+
+ @Test
+ public void testAlignedTimestampRepeatedOrNotIncremented() throws
IOException {
+ String path = dir + File.separator + "test4.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ VectorMeasurementSchema vectorMeasurementSchema =
+ new VectorMeasurementSchema(
+ "d1", new String[] {"s1"}, new TSDataType[] {TSDataType.INT32});
+ AlignedChunkWriterImpl chunkWriter = new
AlignedChunkWriterImpl(vectorMeasurementSchema);
+ chunkWriter.getTimeChunkWriter().write(1);
+ chunkWriter.getTimeChunkWriter().write(2);
+ chunkWriter.getTimeChunkWriter().write(2);
+ chunkWriter.getTimeChunkWriter().write(3);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertFalse(success);
+ }
+
+ @Test
+ public void testAlignedTimestampHasOverlapBetweenPages() throws IOException {
+ String path = dir + File.separator + "test5.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ VectorMeasurementSchema vectorMeasurementSchema =
+ new VectorMeasurementSchema(
+ "d1", new String[] {"s1"}, new TSDataType[] {TSDataType.INT32});
+ AlignedChunkWriterImpl chunkWriter = new
AlignedChunkWriterImpl(vectorMeasurementSchema);
+ chunkWriter.getTimeChunkWriter().write(1);
+ chunkWriter.getTimeChunkWriter().write(2);
+ chunkWriter.getTimeChunkWriter().write(3);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.getTimeChunkWriter().write(3);
+ chunkWriter.getTimeChunkWriter().write(4);
+ chunkWriter.getTimeChunkWriter().write(5);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertFalse(success);
+ }
+
+ @Test
+ public void testAlignedTimestampTimeChunkOffsetEqualsMetadata() throws
IOException {
+ String path = dir + File.separator + "test6.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ VectorMeasurementSchema vectorMeasurementSchema =
+ new VectorMeasurementSchema(
+ "d1", new String[] {"s1"}, new TSDataType[] {TSDataType.INT32});
+ AlignedChunkWriterImpl chunkWriter = new
AlignedChunkWriterImpl(vectorMeasurementSchema);
+ chunkWriter.getTimeChunkWriter().write(1);
+ chunkWriter.getTimeChunkWriter().write(2);
+ chunkWriter.getTimeChunkWriter().write(3);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(1, 1,
false);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(2, 1,
false);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(3, 1,
false);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 3);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+
+ @Test
+ public void testNonAlignedTsFileHasOnePageData() throws IOException {
+ String path = dir + File.separator + "test7.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ TsFileGeneratorUtils.generateSingleNonAlignedSeriesFile(
+ "d1",
+ "s1",
+ new TimeRange[] {new TimeRange(1, 100)},
+ TSEncoding.PLAIN,
+ CompressionType.SNAPPY,
+ path);
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 100);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+
+ @Test
+ public void testNonAlignedTsFileHasManyPage() throws IOException {
+ String path = dir + File.separator + "test8.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ TsFileGeneratorUtils.generateSingleNonAlignedSeriesFile(
+ "d1",
+ "s1",
+ new TimeRange[] {new TimeRange(1, 100), new TimeRange(22, 110)},
+ TSEncoding.PLAIN,
+ CompressionType.SNAPPY,
+ path);
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 110);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+
+ @Test
+ public void testNonAlignedTimestampRepeatedOrNotIncremented() throws
IOException {
+ String path = dir + File.separator + "test9.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ ChunkWriterImpl chunkWriter =
+ new ChunkWriterImpl(new MeasurementSchema("s1", TSDataType.INT32));
+ chunkWriter.getPageWriter().write(1, 2);
+ chunkWriter.getPageWriter().write(2, 2);
+ chunkWriter.getPageWriter().write(2, 2);
+ chunkWriter.getPageWriter().write(3, 2);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertFalse(success);
+ }
+
+ @Test
+ public void testNonAlignedTimestampHasOverlapBetweenPages() throws
IOException {
+ String path = dir + File.separator + "test10.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ ChunkWriterImpl chunkWriter =
+ new ChunkWriterImpl(new MeasurementSchema("s1", TSDataType.INT32));
+ chunkWriter.getPageWriter().write(1, 2);
+ chunkWriter.getPageWriter().write(2, 2);
+ chunkWriter.getPageWriter().write(3, 2);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.getPageWriter().write(3, 4);
+ chunkWriter.getPageWriter().write(4, 4);
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertFalse(success);
+ }
+
+ @Test
+ public void testNonAlignedTimestampTimeChunkOffsetEqualsMetadata() throws
IOException {
+ String path = dir + File.separator + "test11.tsfile";
+ TsFileResource tsFileResource = new TsFileResource(new File(path));
+ try (CompactionTestFileWriter writer = new
CompactionTestFileWriter(tsFileResource)) {
+ writer.startChunkGroup("d1");
+ VectorMeasurementSchema vectorMeasurementSchema =
+ new VectorMeasurementSchema(
+ "d1", new String[] {"s1"}, new TSDataType[] {TSDataType.INT32});
+ AlignedChunkWriterImpl chunkWriter = new
AlignedChunkWriterImpl(vectorMeasurementSchema);
+ chunkWriter.getTimeChunkWriter().write(1);
+ chunkWriter.getTimeChunkWriter().write(2);
+ chunkWriter.getTimeChunkWriter().write(3);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(1, 1,
false);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(2, 1,
false);
+ chunkWriter.getValueChunkWriterByIndex(0).getPageWriter().write(3, 1,
false);
+ chunkWriter.sealCurrentPage();
+ chunkWriter.writeToFileWriter(writer.getFileWriter());
+ writer.endChunkGroup();
+ writer.endFile();
+ }
+ tsFileResource.updateStartTime("d1", 1);
+ tsFileResource.updateEndTime("d1", 3);
+ tsFileResource.serialize();
+ boolean success =
TsFileValidator.getInstance().validateTsFile(tsFileResource);
+ Assert.assertTrue(success);
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/CrossSpaceCompactionWithFastPerformerValidationTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/CrossSpaceCompactionWithFastPerformerValidationTest.java
index 6d4fc3eade0..1d2a0020895 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/CrossSpaceCompactionWithFastPerformerValidationTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/CrossSpaceCompactionWithFastPerformerValidationTest.java
@@ -45,6 +45,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.read.control.FileReaderManag
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.TsFileResourceStatus;
+import org.apache.iotdb.db.storageengine.dataregion.utils.TsFileResourceUtils;
import org.apache.iotdb.db.tools.validate.TsFileValidationTool;
import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
@@ -2154,14 +2155,16 @@ public class
CrossSpaceCompactionWithFastPerformerValidationTest extends Abstrac
// meet overlap files
Assert.assertFalse(
- CompactionUtils.validateTsFileResources(tsFileManager,
COMPACTION_TEST_SG, 0));
+ TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(0).getArrayList()));
tsFileManager.getTsFileList(true).get(0).deserialize();
tsFileManager.getTsFileList(true).get(1).deserialize();
tsFileManager.getTsFileList(true).get(0).degradeTimeIndex();
tsFileManager.getTsFileList(true).get(1).degradeTimeIndex();
Assert.assertTrue(
- CompactionUtils.validateTsFileResources(tsFileManager,
COMPACTION_TEST_SG, 0));
+ TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(0).getArrayList()));
// seq file 4,5 and 6 are being compacted by inner space compaction
List<TsFileResource> sourceFiles = new ArrayList<>();
@@ -2173,7 +2176,8 @@ public class
CrossSpaceCompactionWithFastPerformerValidationTest extends Abstrac
InnerSpaceCompactionTask innerSpaceCompactionTask =
new InnerSpaceCompactionTask(0, tsFileManager, sourceFiles, true,
performer, 0);
Assert.assertTrue(
- CompactionUtils.validateTsFileResources(tsFileManager,
COMPACTION_TEST_SG, 0));
+ TsFileResourceUtils.validateTsFileResourcesHasNoOverlap(
+
tsFileManager.getOrCreateSequenceListByTimePartition(0).getArrayList()));
innerSpaceCompactionTask.start();
validateSeqFiles(true);
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithFastPerformerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithFastPerformerTest.java
index 2fa8907b9da..1135e305047 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithFastPerformerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithFastPerformerTest.java
@@ -446,6 +446,12 @@ public class
RewriteCrossSpaceCompactionWithFastPerformerTest extends AbstractCo
new TsFileManager(COMPACTION_TEST_SG, "0",
STORAGE_GROUP_DIR.getPath());
tsFileManager.addAll(seqResources, true);
tsFileManager.addAll(unseqResources, false);
+ for (TsFileResource resource : seqResources) {
+ Assert.assertTrue(resource.getModFile().exists());
+ }
+ for (TsFileResource resource : unseqResources) {
+ Assert.assertTrue(resource.getModFile().exists());
+ }
CrossSpaceCompactionTask task =
new CrossSpaceCompactionTask(
0,
@@ -457,12 +463,6 @@ public class
RewriteCrossSpaceCompactionWithFastPerformerTest extends AbstractCo
0);
task.start();
- for (TsFileResource resource : seqResources) {
- Assert.assertFalse(resource.getModFile().exists());
- }
- for (TsFileResource resource : unseqResources) {
- Assert.assertFalse(resource.getModFile().exists());
- }
for (TsFileResource resource : targetResources) {
resource.setFile(
new File(
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithReadPointPerformerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithReadPointPerformerTest.java
index 4ed2d5a9b81..3245e29357b 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithReadPointPerformerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/RewriteCrossSpaceCompactionWithReadPointPerformerTest.java
@@ -458,10 +458,10 @@ public class
RewriteCrossSpaceCompactionWithReadPointPerformerTest extends Abstr
task.start();
for (TsFileResource resource : seqResources) {
- Assert.assertFalse(resource.getModFile().exists());
+ Assert.assertTrue(resource.getModFile().exists());
}
for (TsFileResource resource : unseqResources) {
- Assert.assertFalse(resource.getModFile().exists());
+ Assert.assertTrue(resource.getModFile().exists());
}
for (TsFileResource resource : targetResources) {
resource.setFile(
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/InnerCompactionEmptyTsFileTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/InnerCompactionEmptyTsFileTest.java
index 86dbeddf3a5..d1182bdc20d 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/InnerCompactionEmptyTsFileTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/InnerCompactionEmptyTsFileTest.java
@@ -89,6 +89,6 @@ public class InnerCompactionEmptyTsFileTest extends
InnerCompactionTest {
Future<CompactionTaskSummary> future =
CompactionTaskManager.getInstance().getCompactionTaskFutureMayBlock(task);
unseqResources.get(0).readUnlock();
- Assert.assertTrue(future.get().isSuccess());
+ Assert.assertFalse(future.get().isSuccess());
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTestFileWriter.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTestFileWriter.java
index cea4fd2c609..012336a77ca 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTestFileWriter.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/CompactionTestFileWriter.java
@@ -332,4 +332,8 @@ public class CompactionTestFileWriter implements Closeable {
alignedChunkWriter.writeToFileWriter(fileWriter);
}
}
+
+ public TsFileIOWriter getFileWriter() {
+ return fileWriter;
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/TsFileGeneratorUtils.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/TsFileGeneratorUtils.java
index d53963d957c..20c68b17b5f 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/TsFileGeneratorUtils.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/utils/TsFileGeneratorUtils.java
@@ -22,6 +22,7 @@ import
org.apache.iotdb.commons.exception.IllegalPathException;
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.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
import org.apache.iotdb.tsfile.enums.TSDataType;
@@ -40,6 +41,7 @@ import
org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
import org.apache.iotdb.tsfile.write.writer.TsFileIOWriter;
+import java.io.File;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
@@ -373,4 +375,125 @@ public class TsFileGeneratorUtils {
}
return dataTypes;
}
+
+ public static TsFileResource generateSingleNonAlignedSeriesFile(
+ String device,
+ String measurement,
+ TimeRange[] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleNonAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
+
+ public static TsFileResource generateSingleNonAlignedSeriesFile(
+ String device,
+ String measurement,
+ TimeRange[][] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleNonAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
+
+ public static TsFileResource generateSingleNonAlignedSeriesFile(
+ String device,
+ String measurement,
+ TimeRange[][][] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ boolean isSeq,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleNonAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
+
+ public static TsFileResource generateSingleAlignedSeriesFile(
+ String device,
+ List<String> measurement,
+ TimeRange[] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
+
+ public static TsFileResource generateSingleAlignedSeriesFile(
+ String device,
+ List<String> measurement,
+ TimeRange[][] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
+
+ public static TsFileResource generateSingleAlignedSeriesFile(
+ String device,
+ List<String> measurement,
+ TimeRange[][][] chunkTimeRanges,
+ TSEncoding encoding,
+ CompressionType compressionType,
+ String filePath)
+ throws IOException {
+ TsFileResource seqResource1 = new TsFileResource(new File(filePath));
+
+ CompactionTestFileWriter writer1 = new
CompactionTestFileWriter(seqResource1);
+ writer1.startChunkGroup(device);
+ writer1.generateSimpleAlignedSeriesToCurrentDevice(
+ measurement, chunkTimeRanges, encoding, compressionType);
+ writer1.endChunkGroup();
+ writer1.endFile();
+ writer1.close();
+ return seqResource1;
+ }
}
diff --git
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-common.properties
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-common.properties
index a67d3ada788..079cd39f907 100644
---
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-common.properties
+++
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-common.properties
@@ -300,6 +300,8 @@ cluster_name=defaultCluster
# -1 means unlimited.
# database_limit_threshold = -1
+
+
####################
### Configurations for creating schema automatically
####################
@@ -546,6 +548,15 @@ timestamp_precision=ms
# Datatype: int
# device_path_cache_size=500000
+
+# Verify that TSfiles generated by Flush, Load, and Compaction are correct.
The following is verified:
+# 1. Check whether the file contains a header and a tail
+# 2. Check whether files can be deserialized successfully
+# 3. Check whether the file contains data
+# 4. Whether there is time range overlap between data, whether it is
increased, and whether the metadata index offset of the sequence is correct
+# Datatype: boolean
+# enable_tsfile_validation=false
+
####################
### Compaction Configurations
####################
@@ -671,13 +682,7 @@ timestamp_precision=ms
# Datatype: int
# sub_compaction_thread_count=4
-# The level of validation after compaction
-# The details of these three levels are as follows:
-# 1. NONE: the validation after compaction is disabled.
-# 2. RESOURCE_ONLY: the validation after compaction check tsfile resource only.
-# 3. RESOURCE_AND_TSFILE: the validation after compaction check resource and
file.
-# Datatype: String
-# compaction_validation_level=RESOURCE_ONLY
+
####################
### Write Ahead Log Configuration