This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch deletion_expr
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/deletion_expr by this push:
new d7651b7aadf temp save
d7651b7aadf is described below
commit d7651b7aadfcec4212f0a62d2539ee4ff47c50c6
Author: Tian Jiang <[email protected]>
AuthorDate: Tue Sep 24 10:11:41 2024 +0800
temp save
---
.../db/storageengine/dataregion/DataRegion.java | 50 +++++++++---
.../RepairUnsortedFileCompactionPerformer.java | 2 +-
.../execute/task/AbstractCompactionTask.java | 48 ++++++++---
.../execute/task/CrossSpaceCompactionTask.java | 21 ++---
.../execute/task/InnerSpaceCompactionTask.java | 28 +++++--
.../task/InsertionCrossSpaceCompactionTask.java | 18 +++--
.../task/RepairUnsortedFileCompactionTask.java | 27 ++++---
.../execute/task/SettleCompactionTask.java | 6 +-
.../compaction/execute/utils/CompactionUtils.java | 35 --------
.../compaction/repair/RepairTimePartition.java | 1 +
.../repair/RepairTimePartitionScanTask.java | 1 +
.../estimator/CompactionEstimateUtils.java | 2 +-
.../dataregion/modification/ModFileManager.java | 35 +++++---
.../dataregion/modification/ModificationFile.java | 2 +-
.../dataregion/modification/TreeDeletionEntry.java | 4 +
.../dataregion/tsfile/TsFileResource.java | 92 ++++++++++++++--------
.../TsFileOverlapValidationAndRepairTool.java | 2 +-
.../InsertionCrossSpaceCompactionRecoverTest.java | 8 +-
.../NewSizeTieredCompactionSelectorTest.java | 9 ++-
.../repair/RepairUnsortedFileCompactionTest.java | 2 +-
.../settle/SettleCompactionRecoverTest.java | 66 ++++++++--------
21 files changed, 273 insertions(+), 186 deletions(-)
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 e34abbaba14..eb9a8304a74 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
@@ -82,6 +82,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.flush.TsFileFlushPolicy;
import org.apache.iotdb.db.storageengine.dataregion.memtable.IMemTable;
import org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessor;
import
org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessorInfo;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
import org.apache.iotdb.db.storageengine.dataregion.modification.v1.Deletion;
import
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
import org.apache.iotdb.db.storageengine.dataregion.read.IQueryDataSource;
@@ -265,6 +266,10 @@ public class DataRegion implements IDataRegionForQuery {
*/
private final HashMap<Long, VersionController>
timePartitionIdVersionControllerMap =
new HashMap<>();
+ /**
+ * time partition id -> ModFileManager
+ */
+ private final Map<Long, ModFileManager> timePartitionModFileManagerMap = new
ConcurrentHashMap<>();
/** file system factory (local or hdfs). */
private final FSFactory fsFactory = FSFactoryProducer.getFSFactory();
@@ -2155,7 +2160,7 @@ public class DataRegion implements IDataRegionForQuery {
}
/** Seperate tsfiles in TsFileManager to sealedList and unsealedList. */
- private void getTwoKindsOfTsFiles(
+ private void separateTsFileBySealStatus(
List<TsFileResource> sealedResource,
List<TsFileResource> unsealedResource,
long startTime,
@@ -2187,6 +2192,7 @@ public class DataRegion implements IDataRegionForQuery {
writeLock("delete");
boolean hasReleasedLock = false;
+ boolean anyFileRemoved = false;
try {
DataNodeSchemaCache.getInstance().invalidateLastCache(pattern);
@@ -2206,10 +2212,22 @@ public class DataRegion implements IDataRegionForQuery {
List<TsFileResource> sealedTsFileResource = new ArrayList<>();
List<TsFileResource> unsealedTsFileResource = new ArrayList<>();
- getTwoKindsOfTsFiles(sealedTsFileResource, unsealedTsFileResource,
startTime, endTime);
+ separateTsFileBySealStatus(sealedTsFileResource, unsealedTsFileResource,
startTime, endTime);
+ Set<String> deviceMatchInfo = new HashSet<>();
+ sealedTsFileResource.removeIf(f -> canSkipDelete(
+ f,
+ devicePaths,
+ deletion.getStartTime(),
+ deletion.getEndTime(),
+ deviceMatchInfo));
+ unsealedTsFileResource.removeIf(f -> canSkipDelete(
+ f,
+ devicePaths,
+ deletion.getStartTime(),
+ deletion.getEndTime(),
+ deviceMatchInfo));
// deviceMatchInfo is used for filter the matched deviceId in
TsFileResource
// deviceMatchInfo contains the DeviceId means this device matched the
pattern
- Set<String> deviceMatchInfo = new HashSet<>();
deleteDataInFiles(unsealedTsFileResource, deletion, devicePaths,
deviceMatchInfo);
writeUnlock();
hasReleasedLock = true;
@@ -2249,7 +2267,7 @@ public class DataRegion implements IDataRegionForQuery {
}
List<TsFileResource> sealedTsFileResource = new ArrayList<>();
List<TsFileResource> unsealedTsFileResource = new ArrayList<>();
- getTwoKindsOfTsFiles(sealedTsFileResource, unsealedTsFileResource,
startTime, endTime);
+ separateTsFileBySealStatus(sealedTsFileResource, unsealedTsFileResource,
startTime, endTime);
deleteDataDirectlyInFile(unsealedTsFileResource, pathToDelete,
startTime, endTime);
writeUnlock();
releasedLock = true;
@@ -2374,23 +2392,22 @@ public class DataRegion implements IDataRegionForQuery {
return true;
}
+ private void deleteDataInMemory(Collection<TsFileResource>
tsFileResourceList,
+ Deletion deletion,
+ Set<PartialPath> devicePaths,
+ List<TsFileResource> nonWritableFiles) {
+
+ }
+
// suppress warn of Throwable catch
@SuppressWarnings("java:S1181")
private void deleteDataInFiles(
Collection<TsFileResource> tsFileResourceList,
Deletion deletion,
Set<PartialPath> devicePaths,
- Set<String> deviceMatchInfo)
+ List<TsFileResource> nonWritableFiles)
throws IOException {
for (TsFileResource tsFileResource : tsFileResourceList) {
- if (canSkipDelete(
- tsFileResource,
- devicePaths,
- deletion.getStartTime(),
- deletion.getEndTime(),
- deviceMatchInfo)) {
- continue;
- }
ModificationFileV1 modFile = tsFileResource.getOldModFile();
if (tsFileResource.isClosed()) {
@@ -3799,4 +3816,11 @@ public class DataRegion implements IDataRegionForQuery {
public TsFileManager getTsFileManager() {
return tsFileManager;
}
+
+ public ModFileManager getModFileManager(long timePartitionId) {
+ // TODO: conclude better arguments from the previous partition
+ long singleModFileSizeThreshold = config.getSingleModFileSizeThreshold();
+ int levelModFileCntThreshold = config.getLevelModFileCntThreshold();
+ return timePartitionModFileManagerMap.computeIfAbsent(timePartitionId, tid
-> new ModFileManager(levelModFileCntThreshold, singleModFileSizeThreshold));
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/performer/impl/RepairUnsortedFileCompactionPerformer.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/performer/impl/RepairUnsortedFileCompactionPerformer.java
index fd0fc274050..31fadd917bc 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/performer/impl/RepairUnsortedFileCompactionPerformer.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/performer/impl/RepairUnsortedFileCompactionPerformer.java
@@ -69,7 +69,7 @@ public class RepairUnsortedFileCompactionPerformer extends
ReadPointCompactionPe
} else {
targetFile.setTimeIndex(CompactionUtils.buildDeviceTimeIndex(seqSourceFile));
}
- if (seqSourceFile.modFileExists()) {
+ if (seqSourceFile.oldModFileExists()) {
Files.createLink(
new
File(seqSourceFile.getCompactionModFile().getFilePath()).toPath(),
new File(seqSourceFile.getOldModFile().getFilePath()).toPath());
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 6f71924b9ec..f41dc90dea4 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
@@ -35,7 +35,8 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.TsFileIdentifier;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.repair.RepairDataFileScanUtil;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.CompactionTaskManager;
-import
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
+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.TsFileRepairStatus;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
@@ -84,12 +85,15 @@ public abstract class AbstractCompactionTask {
private boolean fileHandleAcquired = false;
protected long compactionConfigVersion = Long.MAX_VALUE;
+ protected ModFileManager modFileManager;
+
protected AbstractCompactionTask(
String storageGroupName,
String dataRegionId,
long timePartition,
TsFileManager tsFileManager,
- long serialId) {
+ long serialId,
+ ModFileManager modFileManager) {
this.storageGroupName = storageGroupName;
this.dataRegionId = dataRegionId;
this.timePartition = timePartition;
@@ -373,13 +377,15 @@ public abstract class AbstractCompactionTask {
return null;
}
- protected void deleteCompactionModsFile(List<TsFileResource>
tsFileResourceList)
- throws IOException {
- for (TsFileResource tsFile : tsFileResourceList) {
- ModificationFileV1 modificationFile = tsFile.getCompactionModFile();
- if (modificationFile.exists()) {
- modificationFile.remove();
- }
+ protected void unsetCompactionModsFile(List<TsFileResource> sourceFileList) {
+ for (TsFileResource tsFile : sourceFileList) {
+ tsFile.setCompactionModFile(null);
+ }
+ }
+
+ protected void deleteCompactionModsFile(List<TsFileResource> targetFileList)
throws IOException {
+ for (TsFileResource resource : targetFileList) {
+ resource.removeModFile();
}
}
@@ -501,4 +507,28 @@ public abstract class AbstractCompactionTask {
public void setRecoverMemoryStatus(boolean recoverMemoryStatus) {
this.recoverMemoryStatus = recoverMemoryStatus;
}
+
+ @SafeVarargs
+ protected final void allocateModFile(List<TsFileResource> targetFiles,
+ List<TsFileResource>... allSourceFiles)
+ throws IOException {
+ // allocate the same mod file for all target files
+ ModificationFile modificationFile =
modFileManager.allocate(targetFiles.get(0));
+ for (int i = 1; i < targetFiles.size(); i++) {
+ // do not persist the modification file path now because the resource is
incomplete
+ targetFiles.get(i).setModFile(modificationFile, false);
+ }
+ // mark the mod file as the compaction mod file for all source files
+ for (List<TsFileResource> sourceFiles : allSourceFiles) {
+ for (TsFileResource resource : sourceFiles) {
+ resource.writeLock();
+ try {
+ // make sure the compaction mod file can be seen by deletion threads
+ resource.setCompactionModFile(modificationFile);
+ } finally {
+ resource.writeUnlock();
+ }
+ }
+ }
+ }
}
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 b80f66737a0..91831c5dc31 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
@@ -34,6 +34,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
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.schedule.CompactionTaskManager;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
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;
@@ -71,13 +72,15 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
List<TsFileResource> selectedUnsequenceFiles,
ICrossCompactionPerformer performer,
long memoryCost,
- long serialId) {
+ long serialId,
+ ModFileManager modFileManager) {
super(
tsFileManager.getStorageGroupName(),
tsFileManager.getDataRegionId(),
timePartition,
tsFileManager,
- serialId);
+ serialId,
+ modFileManager);
this.selectedSequenceFiles = selectedSequenceFiles;
this.selectedUnsequenceFiles = selectedUnsequenceFiles;
this.emptyTargetTsFileResourceList = new ArrayList<>();
@@ -88,8 +91,8 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
}
public CrossSpaceCompactionTask(
- String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile) {
- super(databaseName, dataRegionId, 0L, tsFileManager, 0L);
+ String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile, ModFileManager modFileManager) {
+ super(databaseName, dataRegionId, 0L, tsFileManager, 0L, modFileManager);
this.logFile = logFile;
this.needRecoverTaskInfoFromLogFile = true;
}
@@ -122,7 +125,7 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
}
@Override
- @SuppressWarnings({"squid:S6541", "squid:S3776", "squid:S2142"})
+ @SuppressWarnings({"squid:S6541", "squid:S3776", "squid:S2142", "unchecked"})
public boolean doCompaction() {
recoverMemoryStatus = true;
boolean isSuccess = true;
@@ -171,6 +174,7 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
long startTime = System.currentTimeMillis();
targetTsfileResourceList =
TsFileNameGenerator.getCrossCompactionTargetFileResources(selectedSequenceFiles);
+ allocateModFile(targetTsfileResourceList, selectedSequenceFiles,
selectedUnsequenceFiles);
logFile =
new File(
@@ -197,8 +201,6 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
targetTsfileResourceList,
CompactionTaskType.CROSS,
storageGroupName + "-" + dataRegionId);
- CompactionUtils.combineModsInCrossCompaction(
- selectedSequenceFiles, selectedUnsequenceFiles,
targetTsfileResourceList);
validateCompactionResult(
selectedSequenceFiles, selectedUnsequenceFiles,
targetTsfileResourceList);
@@ -307,8 +309,9 @@ public class CrossSpaceCompactionTask extends
AbstractCompactionTask {
Stream.concat(selectedSequenceFiles.stream(),
selectedUnsequenceFiles.stream())
.collect(Collectors.toList()));
}
- deleteCompactionModsFile(selectedSequenceFiles);
- deleteCompactionModsFile(selectedUnsequenceFiles);
+ unsetCompactionModsFile(selectedSequenceFiles);
+ unsetCompactionModsFile(selectedUnsequenceFiles);
+ deleteCompactionModsFile(targetTsfileResourceList);
// delete target file
if (targetTsfileResourceList != null &&
!deleteTsFilesOnDisk(targetTsfileResourceList)) {
throw new CompactionRecoverException("failed to delete target file %s");
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 48a9ad5b629..c78dd790498 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
@@ -40,6 +40,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimato
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.CompactionEstimateUtils;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.FastCompactionInnerCompactionEstimator;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.ReadChunkInnerCompactionEstimator;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
import
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
@@ -73,13 +74,15 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
List<TsFileResource> selectedTsFileResourceList,
boolean sequence,
ICompactionPerformer performer,
- long serialId) {
+ long serialId,
+ ModFileManager modFileManager) {
super(
tsFileManager.getStorageGroupName(),
tsFileManager.getDataRegionId(),
timePartition,
tsFileManager,
- serialId);
+ serialId,
+ modFileManager);
filesView = new InnerCompactionTaskFilesView(selectedTsFileResourceList,
sequence);
this.performer = performer;
this.hashCode = this.hashCode();
@@ -93,13 +96,15 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
List<TsFileResource> skippedTsFileResourceList,
boolean sequence,
ICompactionPerformer performer,
- long serialId) {
+ long serialId,
+ ModFileManager modFileManager) {
super(
tsFileManager.getStorageGroupName(),
tsFileManager.getDataRegionId(),
timePartition,
tsFileManager,
- serialId);
+ serialId,
+ modFileManager);
this.filesView =
new InnerCompactionTaskFilesView(
selectedTsFileResourceList, skippedTsFileResourceList, sequence);
@@ -110,8 +115,8 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
}
public InnerSpaceCompactionTask(
- String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile) {
- super(databaseName, dataRegionId, 0L, tsFileManager, 0L);
+ String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile, ModFileManager modFileManager) {
+ super(databaseName, dataRegionId, 0L, tsFileManager, 0L, modFileManager);
this.logFile = logFile;
this.needRecoverTaskInfoFromLogFile = true;
this.filesView = new InnerCompactionTaskFilesView();
@@ -183,6 +188,7 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
sortedAllSourceFilesInTask = sourceFiles;
}
+ @TestOnly
protected void setTargetFileForRecover(TsFileResource resource) {
targetFilesInLog = Collections.singletonList(resource);
targetFilesInPerformer = targetFilesInLog;
@@ -292,6 +298,7 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
return isSuccess;
}
+ @SuppressWarnings("unchecked")
protected void calculateSourceFilesAndTargetFiles()
throws DiskSpaceInsufficientException, IOException {
LinkedList<TsFileResource> availablePositionForTargetFiles = new
LinkedList<>();
@@ -345,6 +352,9 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
TsFileNameGenerator.getNewInnerCompactionTargetFileResources(
availablePositionForTargetFiles.subList(0, requiredPositionNum),
filesView.sequence);
}
+
+ allocateModFile(filesView.targetFilesInPerformer,
filesView.sourceFilesInLog);
+
filesView.targetFilesInLog =
new ArrayList<>(
filesView.targetFilesInPerformer.size() +
filesView.renamedTargetFiles.size());
@@ -451,11 +461,12 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
Files.createLink(
new File(newFile.getTsFilePath() +
TsFileResource.RESOURCE_SUFFIX).toPath(),
new File(oldFile.getTsFilePath() +
TsFileResource.RESOURCE_SUFFIX).toPath());
- if (oldFile.modFileExists()) {
+ if (oldFile.oldModFileExists()) {
Files.createLink(
new File(newFile.getTsFilePath() +
ModificationFileV1.FILE_SUFFIX).toPath(),
new File(oldFile.getTsFilePath() +
ModificationFileV1.FILE_SUFFIX).toPath());
}
+ newFile.inheritModFile(oldFile);
newFile.deserialize();
}
CompactionUtils.moveTargetFile(
@@ -540,7 +551,8 @@ public class InnerSpaceCompactionTask extends
AbstractCompactionTask {
if (recoverMemoryStatus) {
replaceTsFileInMemory(targetFiles, filesView.sourceFilesInLog);
}
- deleteCompactionModsFile(filesView.sortedAllSourceFilesInTask);
+ unsetCompactionModsFile(filesView.sortedAllSourceFilesInTask);
+ deleteCompactionModsFile(targetFiles);
// delete target file
for (TsFileResource targetTsFileResource : targetFiles) {
if (targetTsFileResource != null &&
!deleteTsFileOnDisk(targetTsFileResource)) {
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 6ae111048d0..5f14cbdcc74 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
@@ -30,6 +30,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log
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.selector.utils.InsertionCrossCompactionTaskResource;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
import
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
@@ -62,7 +63,8 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
tsFileManager.getDataRegionId(),
timePartition,
tsFileManager,
- serialId);
+ serialId,
+ null);
this.phaser = phaser;
this.selectedSeqFiles = new ArrayList<>();
this.selectedUnseqFiles = new ArrayList<>();
@@ -85,7 +87,7 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
public InsertionCrossSpaceCompactionTask(
String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile) {
- super(databaseName, dataRegionId, 0L, tsFileManager, 0L);
+ super(databaseName, dataRegionId, 0L, tsFileManager, 0L, null);
this.logFile = logFile;
this.needRecoverTaskInfoFromLogFile = true;
this.selectedSeqFiles = Collections.emptyList();
@@ -233,9 +235,8 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
new File(targetTsFile.getPath() +
ModificationFileV1.FILE_SUFFIX).toPath(),
new File(sourceTsFile.getPath() +
ModificationFileV1.FILE_SUFFIX).toPath());
}
-
targetFile.setProgressIndex(unseqFileToInsert.getMaxProgressIndexAfterClose());
+ targetFile.inheritModFile(unseqFileToInsert);
targetFile.deserialize();
-
targetFile.setProgressIndex(unseqFileToInsert.getMaxProgressIndexAfterClose());
}
private boolean recoverTaskInfoFromLogFile() throws IOException {
@@ -297,8 +298,8 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
|| !targetFile.tsFileExists()
|| !targetFile.resourceFileExists()
|| (unseqFileToInsert != null
- && unseqFileToInsert.modFileExists()
- && !targetFile.modFileExists())
+ && unseqFileToInsert.oldModFileExists()
+ && !targetFile.oldModFileExists())
|| failToPassOverlapValidation;
}
@@ -308,7 +309,8 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
replaceTsFileInMemory(
Collections.singletonList(targetFile),
Collections.singletonList(unseqFileToInsert));
}
- deleteCompactionModsFile(Collections.singletonList(unseqFileToInsert));
+ unsetCompactionModsFile(Collections.singletonList(unseqFileToInsert));
+ deleteCompactionModsFile(Collections.singletonList(targetFile));
if (targetFile == null) {
return;
}
@@ -333,7 +335,7 @@ public class InsertionCrossSpaceCompactionTask extends
AbstractCompactionTask {
throw new CompactionRecoverException("source files cannot be deleted
successfully");
}
- deleteCompactionModsFile(Collections.singletonList(unseqFileToInsert));
+ unsetCompactionModsFile(Collections.singletonList(unseqFileToInsert));
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/RepairUnsortedFileCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/RepairUnsortedFileCompactionTask.java
index 08c81808f84..52c4e5242a9 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/RepairUnsortedFileCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/RepairUnsortedFileCompactionTask.java
@@ -25,6 +25,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.performer
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.CompactionUtils;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.utils.log.CompactionLogger;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.estimator.RepairUnsortedFileCompactionEstimator;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileRepairStatus;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
@@ -70,7 +71,8 @@ public class RepairUnsortedFileCompactionTask extends
InnerSpaceCompactionTask {
Collections.singletonList(sourceFile),
sequence,
new RepairUnsortedFileCompactionPerformer(true),
- serialId);
+ serialId,
+ null);
this.sourceFile = sourceFile;
this.innerSpaceEstimator = new RepairUnsortedFileCompactionEstimator();
this.rewriteFile = false;
@@ -89,7 +91,8 @@ public class RepairUnsortedFileCompactionTask extends
InnerSpaceCompactionTask {
Collections.singletonList(sourceFile),
sequence,
new RepairUnsortedFileCompactionPerformer(rewriteFile),
- serialId);
+ serialId,
+ null);
this.sourceFile = sourceFile;
if (rewriteFile) {
this.innerSpaceEstimator = new RepairUnsortedFileCompactionEstimator();
@@ -110,7 +113,8 @@ public class RepairUnsortedFileCompactionTask extends
InnerSpaceCompactionTask {
Collections.singletonList(sourceFile),
sequence,
new RepairUnsortedFileCompactionPerformer(true),
- serialId);
+ serialId,
+ null);
this.sourceFile = sourceFile;
this.innerSpaceEstimator = new RepairUnsortedFileCompactionEstimator();
this.latch = latch;
@@ -131,7 +135,8 @@ public class RepairUnsortedFileCompactionTask extends
InnerSpaceCompactionTask {
Collections.singletonList(sourceFile),
sequence,
new RepairUnsortedFileCompactionPerformer(rewriteFile),
- serialId);
+ serialId,
+ null);
this.sourceFile = sourceFile;
if (rewriteFile) {
this.innerSpaceEstimator = new RepairUnsortedFileCompactionEstimator();
@@ -201,15 +206,11 @@ public class RepairUnsortedFileCompactionTask extends
InnerSpaceCompactionTask {
storageGroupName,
dataRegionId);
- if (rewriteFile) {
- CompactionUtils.combineModsInInnerCompaction(
- filesView.sourceFilesInCompactionPerformer,
filesView.targetFilesInPerformer);
- } else {
- if (sourceFile.modFileExists()) {
- Files.createLink(
- new
File(filesView.targetFilesInPerformer.get(0).getOldModFile().getFilePath()).toPath(),
- new File(sourceFile.getOldModFile().getFilePath()).toPath());
- }
+ filesView.targetFilesInPerformer.get(0).inheritModFile(sourceFile);
+ if (!rewriteFile && sourceFile.oldModFileExists()) {
+ Files.createLink(
+ new
File(filesView.targetFilesInPerformer.get(0).getOldModFile().getFilePath()).toPath(),
+ new File(sourceFile.getOldModFile().getFilePath()).toPath());
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/SettleCompactionTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/SettleCompactionTask.java
index 7ff1134640d..0076dae560a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/SettleCompactionTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/task/SettleCompactionTask.java
@@ -63,7 +63,7 @@ public class SettleCompactionTask extends
InnerSpaceCompactionTask {
boolean isSequence,
ICompactionPerformer performer,
long serialId) {
- super(timePartition, tsFileManager, partiallyDirtyFiles, isSequence,
performer, serialId);
+ super(timePartition, tsFileManager, partiallyDirtyFiles, isSequence,
performer, serialId, null);
this.fullyDirtyFiles = fullyDirtyFiles;
fullyDirtyFiles.forEach(x -> fullyDirtyFileSize += x.getTsFileSize());
partiallyDirtyFiles.forEach(
@@ -76,7 +76,7 @@ public class SettleCompactionTask extends
InnerSpaceCompactionTask {
public SettleCompactionTask(
String databaseName, String dataRegionId, TsFileManager tsFileManager,
File logFile) {
- super(databaseName, dataRegionId, tsFileManager, logFile);
+ super(databaseName, dataRegionId, tsFileManager, logFile, null);
}
@Override
@@ -86,6 +86,7 @@ public class SettleCompactionTask extends
InnerSpaceCompactionTask {
return allSourceFiles;
}
+ @SuppressWarnings("unchecked")
@Override
protected void calculateSourceFilesAndTargetFiles() throws IOException {
filesView.renamedTargetFiles = Collections.emptyList();
@@ -96,6 +97,7 @@ public class SettleCompactionTask extends
InnerSpaceCompactionTask {
TsFileNameGenerator.getSettleCompactionTargetFileResources(
filesView.sourceFilesInCompactionPerformer,
filesView.sequence));
filesView.targetFilesInPerformer = filesView.targetFilesInLog;
+ allocateModFile(filesView.targetFilesInPerformer,
filesView.sourceFilesInCompactionPerformer);
}
@Override
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 8cb6f06319d..091ed37222f 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
@@ -117,41 +117,6 @@ public class CompactionUtils {
targetResource.closeWithoutSettingStatus();
}
- /**
- * Collect all the compaction modification files of source files, and
combines them as the
- * modification file of target file.
- *
- * @throws IOException if io errors occurred
- */
- public static void combineModsInCrossCompaction(
- List<TsFileResource> seqResources,
- List<TsFileResource> unseqResources,
- List<TsFileResource> targetResources)
- throws IOException {
- Set<Modification> modifications = new HashSet<>();
- // get compaction mods from all source unseq files
- for (TsFileResource unseqFile : unseqResources) {
-
modifications.addAll(ModificationFileV1.getCompactionMods(unseqFile).getModifications());
- }
-
- // write target mods file
- for (int i = 0; i < targetResources.size(); i++) {
- TsFileResource targetResource = targetResources.get(i);
- if (targetResource == null) {
- continue;
- }
- Set<Modification> seqModifications =
- new
HashSet<>(ModificationFileV1.getCompactionMods(seqResources.get(i)).getModifications());
- modifications.addAll(seqModifications);
- updateOneTargetMods(targetResource, modifications);
- if (!modifications.isEmpty()) {
- FileMetrics.getInstance().increaseModFileNum(1);
-
FileMetrics.getInstance().increaseModFileSize(targetResource.getOldModFile().getSize());
- }
- modifications.removeAll(seqModifications);
- }
- }
-
/**
* Collect all the compaction modification files of source files, and
combines them as the
* modification file of target file.
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartition.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartition.java
index 87a4d0f7e35..d2b42b4a686 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartition.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartition.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.storageengine.dataregion.compaction.repair;
import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
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;
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartitionScanTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartitionScanTask.java
index 828307e6397..00261dbe4be 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartitionScanTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairTimePartitionScanTask.java
@@ -21,6 +21,7 @@ package
org.apache.iotdb.db.storageengine.dataregion.compaction.repair;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.RepairUnsortedFileCompactionTask;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.CompactionTaskManager;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.ModFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileRepairStatus;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
index 2ab84f56eb9..3c8cfedd19a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/selector/estimator/CompactionEstimateUtils.java
@@ -110,7 +110,7 @@ public class CompactionEstimateUtils {
Map<IDeviceID, Long> deviceMetadataSizeMap = new HashMap<>();
try {
for (TsFileResource resource : resources) {
- if (resource.modFileExists()) {
+ if (resource.oldModFileExists()) {
cost += resource.getOldModFile().getSize();
}
try (CompactionTsFileReader reader =
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModFileManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModFileManager.java
index 79feb967128..6ecdd356900 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModFileManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModFileManager.java
@@ -71,25 +71,30 @@ public class ModFileManager {
}
}
+ /**
+ * Allocate a Mod File by newing or sharing to the give TsFile.
+ * @param resource tsFile to allocate
+ */
@SuppressWarnings("DuplicatedCode")
- public void allocate(TsFileResource resource) throws IOException {
+ public ModificationFile allocate(TsFileResource resource) throws IOException
{
TsFileID tsFileID = resource.getTsFileID();
long levelNum = tsFileID.getInnerCompactionCount();
// find the nearest TsFile that has a mod, i.e, share candidate
TsFileResource prev = resource.getPrev();
TsFileResource next = resource.getNext();
+ ModificationFile allocatedModFile;
while (prev != null || next != null) {
if (prev != null) {
ModificationFile prevModFile = prev.getModFile();
if (prevModFile != null) {
if (shouldAllocateNew(prevModFile, levelNum)) {
- ModificationFile newModFile = allocateNew(resource);
- resource.setModFile(newModFile);
+ allocatedModFile = allocateNew(resource);
} else {
- resource.setModFile(prevModFile);
+ allocatedModFile = prevModFile;
+ allocatedModFile.addReference(resource);
}
- return;
+ return allocatedModFile;
}
prev = prev.getPrev();
}
@@ -98,16 +103,19 @@ public class ModFileManager {
ModificationFile nextModFile = next.getModFile();
if (nextModFile != null) {
if (shouldAllocateNew(nextModFile, levelNum)) {
- ModificationFile newModFile = allocateNew(resource);
- resource.setModFile(newModFile);
+ allocatedModFile = allocateNew(resource);
} else {
- resource.setModFile(nextModFile);
+ allocatedModFile = nextModFile;
+ allocatedModFile.addReference(resource);
}
- return;
+ return allocatedModFile;
}
next = next.getNext();
}
}
+
+ // no mod file found, allocate a new one
+ return allocateNew(resource);
}
private boolean shouldAllocateNew(ModificationFile modificationFile, long
levelNum) {
@@ -122,7 +130,13 @@ public class ModFileManager {
return fileLength > singleModFileSizeThreshold;
}
- public ModificationFile allocateNew(TsFileResource resource) {
+ /**
+ * Force to allocate a new Mod File for the TsFile.
+ * This will NOT set any fields of the TsFileResource.
+ * @param resource TsFile to allocate.
+ * @return the newly allocated Mod File.
+ */
+ private ModificationFile allocateNew(TsFileResource resource) {
TsFileID tsFileID = resource.getTsFileID();
long levelNum = tsFileID.getInnerCompactionCount();
long nextModNum = maxModNum(levelNum) + 1;
@@ -131,6 +145,7 @@ public class ModFileManager {
levelNum,
k -> new TreeMap<>());
synchronized (levelModsFileMap) {
+ // use the provided file as the initial reference to avoid a newly
created Mod File being cleaned
return levelModsFileMap.computeIfAbsent(nextModNum, k -> new
ModificationFile(file, resource));
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModificationFile.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModificationFile.java
index 6ee348047fd..f2273d3e218 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModificationFile.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/ModificationFile.java
@@ -69,7 +69,7 @@ public class ModificationFile implements AutoCloseable {
}
@Override
- public void close() throws Exception {
+ public void close() throws IOException {
fileOutputStream.close();
fileOutputStream = null;
channel.close();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/TreeDeletionEntry.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/TreeDeletionEntry.java
index a2607f49940..bf33bd882e2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/TreeDeletionEntry.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/modification/TreeDeletionEntry.java
@@ -34,6 +34,10 @@ public class TreeDeletionEntry extends ModEntry {
super(ModType.TREE_DELETION);
}
+ public TreeDeletionEntry(PartialPath path, long start, long end) {
+ this(path, new TimeRange(start, end));
+ }
+
public TreeDeletionEntry(PartialPath path, TimeRange timeRange) {
this();
this.pathPattern = path;
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
index 4ae4e4a748a..56268af6fc2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.storageengine.dataregion.tsfile;
import java.nio.channels.FileChannel;
+import java.util.Collections;
import org.apache.iotdb.commons.consensus.index.ProgressIndex;
import org.apache.iotdb.commons.consensus.index.ProgressIndexType;
import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
@@ -77,7 +78,7 @@ import java.util.concurrent.atomic.AtomicReference;
import static org.apache.iotdb.commons.conf.IoTDBConstant.FILE_NAME_SEPARATOR;
import static org.apache.tsfile.common.constant.TsFileConstant.TSFILE_SUFFIX;
-@SuppressWarnings("java:S1135") // ignore todos
+@SuppressWarnings({"java:S1135", "resource"}) // ignore todos
public class TsFileResource {
private static final long INSTANCE_SIZE =
@@ -113,9 +114,10 @@ public class TsFileResource {
private long modFileOffset;
@SuppressWarnings("squid:S3077")
private volatile ModificationFileV1 oldModFile;
+ private volatile boolean oldModFileChecked = false;
@SuppressWarnings("squid:S3077")
- private volatile ModificationFileV1 compactionModFile;
+ private volatile ModificationFile compactionModFile;
private ModFileManager modFileManager;
// the start pos of mod file path in this TsFileResource
private long modFilePathOffset = -1;
@@ -339,9 +341,19 @@ public class TsFileResource {
}
}
- public void setModFile(ModificationFile modFile) throws IOException {
+ public void inheritModFile(TsFileResource parent) {
+ this.modFile = parent.modFile;
+ this.modFilePathOffset = parent.modFilePathOffset;
+ this.modFileOffset = parent.modFileOffset;
+ }
+
+ public void setModFile(ModificationFile modFile, boolean persist) throws
IOException {
this.modFile = modFile;
this.modFileOffset = modFile.getFile().length();
+ if (!persist) {
+ return;
+ }
+
File resFile = new File(file + RESOURCE_SUFFIX);
if (!resFile.exists()) {
// the file has not been serialized, just serialize it
@@ -423,12 +435,12 @@ public class TsFileResource {
return file != null && file.exists();
}
- public boolean modFileExists() {
- return getOldModFile().exists();
- }
-
- public boolean compactionModFileExists() {
- return getCompactionModFile().exists();
+ public boolean oldModFileExists() {
+ ModificationFileV1 oldModFile = getOldModFile();
+ if (oldModFile == null) {
+ return false;
+ }
+ return oldModFile.exists();
}
public List<IChunkMetadata> getChunkMetadataList(PartialPath seriesPath) {
@@ -441,27 +453,32 @@ public class TsFileResource {
@SuppressWarnings("squid:S2886")
public ModificationFileV1 getOldModFile() {
- if (oldModFile == null) {
- synchronized (this) {
- if (oldModFile == null) {
- oldModFile = ModificationFileV1.getNormalMods(this);
- }
+ if (oldModFileChecked) {
+ return oldModFile;
+ }
+ writeLock();
+ // avoid creating different old mod files
+ try {
+ File oldModFile = new File(getTsFilePath() +
ModificationFileV1.FILE_SUFFIX);
+ if (oldModFile.exists()) {
+ this.oldModFile = new ModificationFileV1(oldModFile.getPath());
}
+ oldModFileChecked = true;
+ } finally {
+ writeUnlock();
}
- return oldModFile;
+ return this.oldModFile;
}
- public ModificationFileV1 getCompactionModFile() {
- if (compactionModFile == null) {
- synchronized (this) {
- if (compactionModFile == null) {
- compactionModFile = ModificationFileV1.getCompactionMods(this);
- }
- }
- }
+ public ModificationFile getCompactionModFile() {
return compactionModFile;
}
+ public void setCompactionModFile(
+ ModificationFile compactionModFile) {
+ this.compactionModFile = compactionModFile;
+ }
+
public void resetModFile() throws IOException {
if (oldModFile != null) {
synchronized (this) {
@@ -609,14 +626,15 @@ public class TsFileResource {
/** Used for compaction. */
public void closeWithoutSettingStatus() throws IOException {
+ if (modFile != null) {
+ modFile.close();
+ modFile = null;
+ }
if (oldModFile != null) {
oldModFile.close();
oldModFile = null;
}
- if (compactionModFile != null) {
- compactionModFile.close();
- compactionModFile = null;
- }
+
processor = null;
pathToChunkMetadataListMap = null;
pathToReadOnlyMemChunkMap = null;
@@ -681,8 +699,16 @@ public class TsFileResource {
}
public void removeModFile() throws IOException {
- getOldModFile().remove();
- oldModFile = null;
+ ModificationFileV1 oldModFile = getOldModFile();
+ if (oldModFile != null) {
+ oldModFile.remove();
+ setOldModFile(null);
+ }
+
+ if (modFile != null) {
+ modFile.removeReferences(Collections.singletonList(this));
+ modFile = null;
+ }
}
/**
@@ -691,6 +717,7 @@ public class TsFileResource {
*/
public boolean remove() {
forceMarkDeleted();
+
try {
fsFactory.deleteIfExists(file);
fsFactory.deleteIfExists(
@@ -702,10 +729,9 @@ public class TsFileResource {
if (!removeResourceFile()) {
return false;
}
+
try {
- fsFactory.deleteIfExists(fsFactory.getFile(file.getPath() +
ModificationFileV1.FILE_SUFFIX));
- fsFactory.deleteIfExists(
- fsFactory.getFile(file.getPath() +
ModificationFileV1.COMPACTION_FILE_SUFFIX));
+ removeModFile();
} catch (IOException e) {
LOGGER.error("ModificationFile {} cannot be deleted: {}", file,
e.getMessage());
return false;
@@ -1298,7 +1324,7 @@ public class TsFileResource {
public ModificationFile getModFileMayAllocate() throws IOException {
if (modFile == null) {
- modFileManager.allocate(this);
+ setModFile(modFileManager.allocate(this), true);
}
return modFile;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/tools/validate/TsFileOverlapValidationAndRepairTool.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/tools/validate/TsFileOverlapValidationAndRepairTool.java
index d4fa4628e13..8e1f0a86809 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/tools/validate/TsFileOverlapValidationAndRepairTool.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/tools/validate/TsFileOverlapValidationAndRepairTool.java
@@ -130,7 +130,7 @@ public class TsFileOverlapValidationAndRepairTool {
moveFile(
new File(tsfile.getAbsolutePath() + TsFileResource.RESOURCE_SUFFIX),
new File(targetFile.getAbsolutePath() +
TsFileResource.RESOURCE_SUFFIX));
- if (resource.modFileExists()) {
+ if (resource.oldModFileExists()) {
moveFile(
new File(tsfile.getAbsolutePath() + ModificationFileV1.FILE_SUFFIX),
new File(targetFile.getAbsolutePath() +
ModificationFileV1.FILE_SUFFIX));
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/InsertionCrossSpaceCompactionRecoverTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/InsertionCrossSpaceCompactionRecoverTest.java
index a801fd3eb06..49f8f2a96f5 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/InsertionCrossSpaceCompactionRecoverTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/cross/InsertionCrossSpaceCompactionRecoverTest.java
@@ -159,7 +159,7 @@ public class InsertionCrossSpaceCompactionRecoverTest
extends AbstractCompaction
Assert.assertFalse(targetFile.tsFileExists());
Assert.assertFalse(targetFile.resourceFileExists());
- Assert.assertFalse(targetFile.modFileExists());
+ Assert.assertFalse(targetFile.oldModFileExists());
}
@Test
@@ -249,7 +249,7 @@ public class InsertionCrossSpaceCompactionRecoverTest
extends AbstractCompaction
Assert.assertTrue(targetFile.tsFileExists());
Assert.assertTrue(targetFile.resourceFileExists());
- Assert.assertFalse(targetFile.modFileExists());
+ Assert.assertFalse(targetFile.oldModFileExists());
}
@Test
@@ -342,7 +342,7 @@ public class InsertionCrossSpaceCompactionRecoverTest
extends AbstractCompaction
Assert.assertTrue(targetFile.tsFileExists());
Assert.assertTrue(targetFile.resourceFileExists());
- Assert.assertTrue(targetFile.modFileExists());
+ Assert.assertTrue(targetFile.oldModFileExists());
}
@Test
@@ -424,7 +424,7 @@ public class InsertionCrossSpaceCompactionRecoverTest
extends AbstractCompaction
Assert.assertFalse(targetFile.tsFileExists());
Assert.assertFalse(targetFile.resourceFileExists());
- Assert.assertFalse(targetFile.modFileExists());
+ Assert.assertFalse(targetFile.oldModFileExists());
}
private TsFileResource createTsFileResource(String name, boolean seq) {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/sizetiered/NewSizeTieredCompactionSelectorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/sizetiered/NewSizeTieredCompactionSelectorTest.java
index 37b12a7c4fd..f5a3dd7beea 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/sizetiered/NewSizeTieredCompactionSelectorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/inner/sizetiered/NewSizeTieredCompactionSelectorTest.java
@@ -30,6 +30,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.compaction.execute.task.Inne
import
org.apache.iotdb.db.storageengine.dataregion.compaction.schedule.CompactionScheduleContext;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.selector.impl.NewSizeTieredCompactionSelector;
import
org.apache.iotdb.db.storageengine.dataregion.compaction.utils.CompactionTestFileWriter;
+import
org.apache.iotdb.db.storageengine.dataregion.modification.TreeDeletionEntry;
import org.apache.iotdb.db.storageengine.dataregion.modification.v1.Deletion;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResourceStatus;
@@ -356,7 +357,7 @@ public class NewSizeTieredCompactionSelectorTest extends
AbstractCompactionTest
}
resource
.getCompactionModFile()
- .write(new Deletion(new PartialPath("root.**"), Long.MAX_VALUE,
Long.MAX_VALUE));
+ .write(new TreeDeletionEntry(new PartialPath("root.**"),
Long.MAX_VALUE, Long.MAX_VALUE));
resource.getCompactionModFile().close();
seqResources.add(resource);
}
@@ -378,11 +379,11 @@ public class NewSizeTieredCompactionSelectorTest extends
AbstractCompactionTest
for (int i = 0; i < filesAfterCompaction.size(); i++) {
TsFileResource resource = filesAfterCompaction.get(i);
if (i == 8) {
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
} else {
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
}
- Assert.assertFalse(resource.compactionModFileExists());
+ Assert.assertFalse(resource.getCompactionModFile().getFile().exists());
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairUnsortedFileCompactionTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairUnsortedFileCompactionTest.java
index 9cc7aef40b1..25e33054b6a 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairUnsortedFileCompactionTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/repair/RepairUnsortedFileCompactionTest.java
@@ -579,7 +579,7 @@ public class RepairUnsortedFileCompactionTest extends
AbstractRepairDataTest {
TsFileResource targetResource = tsFileManager.getTsFileList(false).get(0);
Assert.assertTrue(TsFileResourceUtils.validateTsFileDataCorrectness(targetResource));
Assert.assertTrue(TsFileResourceUtils.validateTsFileResourceCorrectness(targetResource));
- Assert.assertTrue(targetResource.modFileExists());
+ Assert.assertTrue(targetResource.oldModFileExists());
Assert.assertEquals(1,
targetResource.getOldModFile().getModifications().size());
Deletion modification =
(Deletion)
targetResource.getOldModFile().getModifications().iterator().next();
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/settle/SettleCompactionRecoverTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/settle/SettleCompactionRecoverTest.java
index a8b92828bc9..4253d8165ac 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/settle/SettleCompactionRecoverTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/compaction/settle/SettleCompactionRecoverTest.java
@@ -112,13 +112,13 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -167,13 +167,13 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -249,7 +249,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -261,7 +261,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// resource file exist
for (TsFileResource resource : tsFileManager.getTsFileList(false)) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -269,7 +269,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// target resource not exist
Assert.assertFalse(targetResource.resourceFileExists());
Assert.assertFalse(targetResource.tsFileExists());
- Assert.assertFalse(targetResource.modFileExists());
+ Assert.assertFalse(targetResource.oldModFileExists());
}
@Test
@@ -342,25 +342,25 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
// target file exist
Assert.assertTrue(targetResource.resourceFileExists());
Assert.assertTrue(targetResource.tsFileExists());
- Assert.assertTrue(targetResource.modFileExists());
+ Assert.assertTrue(targetResource.oldModFileExists());
Assert.assertEquals(1, tsFileManager.getTsFileList(false).size());
for (TsFileResource resource : tsFileManager.getTsFileList(false)) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -439,20 +439,20 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
// target file is deleted after compaction
Assert.assertFalse(targetResource.resourceFileExists());
Assert.assertFalse(targetResource.tsFileExists());
- Assert.assertFalse(targetResource.modFileExists());
+ Assert.assertFalse(targetResource.oldModFileExists());
Assert.assertFalse(targetResource.getCompactionModFile().exists());
Assert.assertEquals(0, tsFileManager.getTsFileList(false).size());
@@ -538,7 +538,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -596,13 +596,13 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -655,7 +655,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertTrue(resource.getCompactionModFile().exists());
}
@@ -665,13 +665,13 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -747,7 +747,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -755,7 +755,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// resource file exist
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -763,7 +763,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// target resource not exist
Assert.assertFalse(targetResource.resourceFileExists());
Assert.assertFalse(targetResource.tsFileExists());
- Assert.assertFalse(targetResource.modFileExists());
+ Assert.assertFalse(targetResource.oldModFileExists());
Assert.assertFalse(logFile.exists());
}
@@ -826,7 +826,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -834,7 +834,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// resource file exist
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertTrue(resource.tsFileExists());
- Assert.assertTrue(resource.modFileExists());
+ Assert.assertTrue(resource.oldModFileExists());
Assert.assertTrue(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
@@ -842,7 +842,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// target resource not exist
Assert.assertFalse(targetResource.resourceFileExists());
Assert.assertFalse(targetResource.tsFileExists());
- Assert.assertFalse(targetResource.modFileExists());
+ Assert.assertFalse(targetResource.oldModFileExists());
Assert.assertFalse(logFile.exists());
}
@@ -917,20 +917,20 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
// target file exist
Assert.assertTrue(targetResource.resourceFileExists());
Assert.assertTrue(targetResource.tsFileExists());
- Assert.assertTrue(targetResource.modFileExists());
+ Assert.assertTrue(targetResource.oldModFileExists());
Assert.assertFalse(logFile.exists());
}
@@ -1011,20 +1011,20 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
for (TsFileResource resource : partialDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}
// target file is deleted after compaction
Assert.assertFalse(targetResource.resourceFileExists());
Assert.assertFalse(targetResource.tsFileExists());
- Assert.assertFalse(targetResource.modFileExists());
+ Assert.assertFalse(targetResource.oldModFileExists());
Assert.assertFalse(targetResource.getCompactionModFile().exists());
Assert.assertFalse(logFile.exists());
@@ -1111,7 +1111,7 @@ public class SettleCompactionRecoverTest extends
AbstractCompactionTest {
// source files not exist
for (TsFileResource resource : allDeletedFiles) {
Assert.assertFalse(resource.tsFileExists());
- Assert.assertFalse(resource.modFileExists());
+ Assert.assertFalse(resource.oldModFileExists());
Assert.assertFalse(resource.resourceFileExists());
Assert.assertFalse(resource.getCompactionModFile().exists());
}