This is an automated email from the ASF dual-hosted git repository.
qiaojialin 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 5d9beb4 [IOTDB-2039] Fix data redundant after too many open files
exception occurs during compaction (#4498)
5d9beb4 is described below
commit 5d9beb446854cc654404d0a8e1b141c91e5878f3
Author: liuxuxin <[email protected]>
AuthorDate: Wed Dec 15 13:39:43 2021 +0800
[IOTDB-2039] Fix data redundant after too many open files exception occurs
during compaction (#4498)
---
.../db/engine/compaction/CompactionScheduler.java | 3 +
.../inner/AbstractInnerSpaceCompactionTask.java | 4 -
.../InnerSpaceCompactionExceptionHandler.java | 246 ++++++++++++
.../SizeTieredCompactionRecoverTask.java | 12 +-
.../inner/sizetiered/SizeTieredCompactionTask.java | 118 +++---
.../inner/utils/InnerSpaceCompactionUtils.java | 49 ++-
.../db/engine/storagegroup/TsFileManager.java | 10 +
.../db/engine/storagegroup/TsFileResource.java | 21 +-
.../db/engine/storagegroup/TsFileResourceList.java | 59 +++
.../inner/AbstractInnerSpaceCompactionTest.java | 295 ++++++++++++++
.../compaction/inner/InnerSeqCompactionTest.java | 12 +-
.../inner/InnerSpaceCompactionExceptionTest.java | 427 +++++++++++++++++++++
.../compaction/inner/InnerUnseqCompactionTest.java | 3 +-
.../SizeTieredCompactionHandleExceptionTest.java | 197 ++++++++++
.../SizeTieredCompactionRecoverTest.java | 250 +-----------
.../storagegroup/TsFileResourceListTest.java | 30 ++
16 files changed, 1424 insertions(+), 312 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionScheduler.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionScheduler.java
index 6bccf4d..8a04a0c 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionScheduler.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionScheduler.java
@@ -55,6 +55,9 @@ public class CompactionScheduler {
new ConcurrentHashMap<>();
public static void scheduleCompaction(TsFileManager tsFileManager, long
timePartition) {
+ if (!tsFileManager.isAllowCompaction()) {
+ return;
+ }
tsFileManager.readLock();
try {
TsFileResourceList sequenceFileList =
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTask.java
index fa1e830..88a615e 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTask.java
@@ -98,10 +98,6 @@ public abstract class AbstractInnerSpaceCompactionTask
extends AbstractCompactio
return maxFileVersion;
}
- public int getMaxCompactionCount() {
- return maxCompactionCount;
- }
-
@Override
public boolean checkValidAndSetMerging() {
for (TsFileResource resource : selectedTsFileResourceList) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionHandler.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionHandler.java
new file mode 100644
index 0000000..abcb489
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionHandler.java
@@ -0,0 +1,246 @@
+/*
+ * 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.engine.compaction.inner;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
+import org.apache.iotdb.db.engine.modification.Modification;
+import org.apache.iotdb.db.engine.modification.ModificationFile;
+import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResourceList;
+import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
+
+import org.apache.commons.io.FileUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * This class is used to handle exception (including OOM error) occurred
during compaction. The
+ * <i>allowCompaction</i> flag in {@link
org.apache.iotdb.db.engine.storagegroup.TsFileManager} may
+ * be set to false if exception cannot be handled correctly(such as OOM during
handling exception),
+ * after which the subsequent compaction in this storage group will not be
carried out. Under some
+ * serious circumstances(such as data lost), the system will be set to
read-only.
+ */
+public class InnerSpaceCompactionExceptionHandler {
+ private static final Logger LOGGER = LoggerFactory.getLogger("COMPACTION");
+
+ public static void handleException(
+ String fullStorageGroupName,
+ File logFile,
+ TsFileResource targetTsFile,
+ List<TsFileResource> selectedTsFileResourceList,
+ TsFileManager tsFileManager,
+ TsFileResourceList tsFileResourceList) {
+
+ if (!logFile.exists()) {
+ // log file does not exist
+ // it means the compaction has not started yet
+ // we need not to handle it
+ return;
+ }
+
+ boolean handleSuccess = true;
+
+ List<TsFileResource> lostSourceFiles = new ArrayList<>();
+ boolean allSourceFileExist =
+ checkAllSourceFilesExist(selectedTsFileResourceList, lostSourceFiles);
+
+ if (allSourceFileExist) {
+ handleSuccess =
+ handleWhenAllSourceFilesExist(
+ fullStorageGroupName, targetTsFile, selectedTsFileResourceList,
tsFileResourceList);
+ } else {
+ // some source file does not exists
+ // it means we start to delete source file
+ LOGGER.info(
+ "{} [Compaction][ExceptionHandler] some source files {} is lost",
+ fullStorageGroupName,
+ lostSourceFiles);
+ if (!targetTsFile.getTsFile().exists()) {
+ // some source files are missed, and target file not exists
+ // some data is lost, set the system to read-only
+ LOGGER.warn(
+ "{} [Compaction][ExceptionHandler] target file {} does not exist
either, do nothing. Set system to read-only",
+ fullStorageGroupName,
+ targetTsFile);
+ IoTDBDescriptor.getInstance().getConfig().setReadOnly(true);
+ handleSuccess = false;
+ } else {
+ handleSuccess =
+ handleWhenSomeSourceFilesLost(
+ fullStorageGroupName,
+ targetTsFile,
+ selectedTsFileResourceList,
+ tsFileResourceList,
+ lostSourceFiles);
+ }
+ }
+
+ if (!handleSuccess) {
+ LOGGER.error(
+ "{} [Compaction][ExceptionHandler] Failed to handle exception, set
allowCompaction to false",
+ fullStorageGroupName);
+ tsFileManager.setAllowCompaction(false);
+ } else {
+ LOGGER.info(
+ "{} [Compaction][ExceptionHandler] Handle exception successfully,
delete log file {}",
+ fullStorageGroupName,
+ logFile);
+ try {
+ FileUtils.delete(logFile);
+ } catch (IOException e) {
+ LOGGER.error(
+ "{} [Compaction][ExceptionHandler] Exception occurs while deleting
log file {}, set allowCompaction to false",
+ fullStorageGroupName,
+ logFile,
+ e);
+ tsFileManager.setAllowCompaction(false);
+ }
+ }
+ }
+
+ private static boolean checkAllSourceFilesExist(
+ List<TsFileResource> sourceFiles, List<TsFileResource> lostSourceFiles) {
+ boolean allSourceFileExist = true;
+ for (TsFileResource sourceTsFile : sourceFiles) {
+ if (!sourceTsFile.getTsFile().exists()) {
+ allSourceFileExist = false;
+ lostSourceFiles.add(sourceTsFile);
+ }
+ }
+ return allSourceFileExist;
+ }
+
+ private static boolean handleWhenAllSourceFilesExist(
+ String fullStorageGroupName,
+ TsFileResource targetTsFile,
+ List<TsFileResource> selectedTsFileResourceList,
+ TsFileResourceList tsFileResourceList) {
+ // all source file exists, delete the target file
+ LOGGER.info(
+ "{} [Compaction][ExceptionHandler] all source files {} exists, delete
target file {}",
+ fullStorageGroupName,
+ selectedTsFileResourceList,
+ targetTsFile);
+ if (!targetTsFile.remove()) {
+ // failed to remove target tsfile
+ // system should not carry out the subsequent compaction in case of data
redundant
+ LOGGER.warn(
+ "{} [Compaction][ExceptionHandler] failed to remove target file {}",
+ fullStorageGroupName,
+ targetTsFile);
+ return false;
+ }
+ // deal with compaction modification
+ try {
+ for (TsFileResource sourceFile : selectedTsFileResourceList) {
+ if (sourceFile.getCompactionModFile().exists()) {
+ ModificationFile compactionModificationFile =
+ ModificationFile.getCompactionMods(sourceFile);
+ Collection<Modification> newModification =
compactionModificationFile.getModifications();
+ compactionModificationFile.close();
+ // write the modifications to a new modification file
+ sourceFile.resetModFile();
+ try (ModificationFile newModificationFile = sourceFile.getModFile())
{
+ for (Modification modification : newModification) {
+ newModificationFile.write(modification);
+ }
+ }
+ FileUtils.delete(new
File(ModificationFile.getCompactionMods(sourceFile).getFilePath()));
+ }
+ }
+ } catch (Throwable e) {
+ LOGGER.error(
+ "{} Exception occurs while handling exception, set allowCompaction
to false",
+ fullStorageGroupName,
+ e);
+ return false;
+ }
+ return true;
+ }
+
+ private static boolean handleWhenSomeSourceFilesLost(
+ String fullStorageGroupName,
+ TsFileResource targetTsFile,
+ List<TsFileResource> selectedTsFileResourceList,
+ TsFileResourceList tsFileResourceList,
+ List<TsFileResource> lostSourceFiles) {
+ boolean handleSuccess = true;
+ try {
+ RestorableTsFileIOWriter writer = new
RestorableTsFileIOWriter(targetTsFile.getTsFile());
+ writer.close();
+ if (!writer.hasCrashed()) {
+ // target file is complete, delete source files
+ LOGGER.info(
+ "{} [Compaction][ExceptionHandler] target file {} is complete,
delete remaining source files",
+ fullStorageGroupName,
+ targetTsFile);
+ for (TsFileResource sourceFile : selectedTsFileResourceList) {
+ if (!sourceFile.remove()) {
+ LOGGER.warn(
+ "{} [Compaction][ExceptionHandler] failed to remove source
file {}",
+ fullStorageGroupName,
+ sourceFile);
+ handleSuccess = false;
+ } else {
+ tsFileResourceList.remove(sourceFile);
+ }
+ }
+
+ if (targetTsFile.getModFile().exists()) {
+ // if origin mods file exists, remove it, and generate a new mods
file
+ FileUtils.delete(new File(targetTsFile.getModFile().getFilePath()));
+ }
+
+
InnerSpaceCompactionUtils.combineModsInCompaction(selectedTsFileResourceList,
targetTsFile);
+ InnerSpaceCompactionUtils.deleteModificationForSourceFile(
+ selectedTsFileResourceList, fullStorageGroupName);
+
+ if (!tsFileResourceList.contains(targetTsFile)) {
+ tsFileResourceList.keepOrderInsert(targetTsFile);
+ }
+ } else {
+ // target file is not complete, and some source file is lost
+ // some data is lost
+ LOGGER.warn(
+ "{} [Compaction][ExceptionHandler] target file {} is not complete,
and some source files {} is lost, do nothing. Set allowCompaction to false",
+ fullStorageGroupName,
+ targetTsFile,
+ lostSourceFiles);
+ IoTDBDescriptor.getInstance().getConfig().setReadOnly(true);
+ handleSuccess = false;
+ }
+ } catch (Throwable e) {
+ // If we are handling OOM error, OOM may occur again during checking
target file
+ LOGGER.error(
+ "{} [Compaction][ExceptionHandler] Another exception occurs during
handling exception",
+ fullStorageGroupName,
+ e);
+ handleSuccess = false;
+ }
+ return handleSuccess;
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTask.java
index 95d4927..21b359f 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTask.java
@@ -22,6 +22,7 @@ import org.apache.iotdb.db.engine.compaction.TsFileIdentifier;
import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogAnalyzer;
import org.apache.iotdb.db.engine.compaction.task.AbstractCompactionTask;
+import org.apache.iotdb.db.engine.modification.ModificationFile;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
import org.apache.iotdb.tsfile.write.writer.RestorableTsFileIOWriter;
@@ -113,6 +114,7 @@ public class SizeTieredCompactionRecoverTask extends
SizeTieredCompactionTask {
File resourceFile = new File(targetFile.getPath() + ".resource");
RestorableTsFileIOWriter writer = new
RestorableTsFileIOWriter(targetFile, false);
+ writer.close();
if (writer.hasCrashed()) {
LOGGER.info(
"{}-{} [Compaction][Recover] target file {} crash, start to
delete it",
@@ -120,7 +122,6 @@ public class SizeTieredCompactionRecoverTask extends
SizeTieredCompactionTask {
virtualStorageGroup,
targetFile);
// the target tsfile is crashed, it is not completed
- writer.close();
if (!targetFile.delete()) {
LOGGER.error(
"{}-{} [Compaction][Recover] fail to delete target file {},
this may cause data incorrectness",
@@ -151,10 +152,17 @@ public class SizeTieredCompactionRecoverTask extends
SizeTieredCompactionTask {
sourceTsFileResources.add(new TsFileResource(sourceFile));
}
}
+ ModificationFile modificationFileForTargetFile =
+ ModificationFile.getNormalMods(targetResource);
+ if (!modificationFileForTargetFile.exists()) {
+ InnerSpaceCompactionUtils.combineModsInCompaction(
+ sourceTsFileResources, targetResource);
+ }
InnerSpaceCompactionUtils.deleteTsFilesInDisk(
sourceTsFileResources, fullStorageGroupName);
- combineModsInCompaction(sourceTsFileResources, targetResource);
+ InnerSpaceCompactionUtils.deleteModificationForSourceFile(
+ sourceTsFileResources, logicalStorageGroupName + "-" +
virtualStorageGroup);
}
}
} catch (IOException e) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTask.java
index 0661da8..e2ab204 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionTask.java
@@ -19,11 +19,10 @@
package org.apache.iotdb.db.engine.compaction.inner.sizetiered;
import
org.apache.iotdb.db.engine.compaction.inner.AbstractInnerSpaceCompactionTask;
+import
org.apache.iotdb.db.engine.compaction.inner.InnerSpaceCompactionExceptionHandler;
import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger;
import org.apache.iotdb.db.engine.compaction.task.AbstractCompactionTask;
-import org.apache.iotdb.db.engine.modification.Modification;
-import org.apache.iotdb.db.engine.modification.ModificationFile;
import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
import org.apache.iotdb.db.engine.storagegroup.TsFileNameGenerator;
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
@@ -36,8 +35,6 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.IOException;
-import java.util.ArrayList;
-import java.util.Collection;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
@@ -72,34 +69,6 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
this.tsFileManager = tsFileManager;
}
- public static void combineModsInCompaction(
- Collection<TsFileResource> mergeTsFiles, TsFileResource targetTsFile)
throws IOException {
- List<Modification> modifications = new ArrayList<>();
- for (TsFileResource mergeTsFile : mergeTsFiles) {
- try (ModificationFile sourceCompactionModificationFile =
- ModificationFile.getCompactionMods(mergeTsFile)) {
-
modifications.addAll(sourceCompactionModificationFile.getModifications());
- if (sourceCompactionModificationFile.exists()) {
- sourceCompactionModificationFile.remove();
- }
- }
- ModificationFile sourceModificationFile =
ModificationFile.getNormalMods(mergeTsFile);
- if (sourceModificationFile.exists()) {
- sourceModificationFile.remove();
- }
- }
- if (!modifications.isEmpty()) {
- try (ModificationFile modificationFile =
ModificationFile.getNormalMods(targetTsFile)) {
- for (Modification modification : modifications) {
- // we have to set modification offset to MAX_VALUE, as the offset of
source chunk may
- // change after compaction
- modification.setFileOffset(Long.MAX_VALUE);
- modificationFile.write(modification);
- }
- }
- }
- }
-
@Override
protected void doCompaction() throws Exception {
long startTime = System.currentTimeMillis();
@@ -115,6 +84,21 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
fullStorageGroupName,
selectedTsFileResourceList.size());
File logFile = null;
+ SizeTieredCompactionLogger sizeTieredCompactionLogger = null;
+ // to mark if we got the write lock or read lock of the selected tsfile
+ boolean[] isHoldingReadLock = new
boolean[selectedTsFileResourceList.size()];
+ boolean[] isHoldingWriteLock = new
boolean[selectedTsFileResourceList.size()];
+ for (int i = 0; i < selectedTsFileResourceList.size(); ++i) {
+ isHoldingReadLock[i] = false;
+ isHoldingWriteLock[i] = false;
+ }
+ LOGGER.info(
+ "{} [Compaction] Try to get the read lock of all selected files",
fullStorageGroupName);
+ for (int i = 0; i < selectedTsFileResourceList.size(); ++i) {
+ selectedTsFileResourceList.get(i).readLock();
+ isHoldingReadLock[i] = true;
+ }
+
try {
logFile =
new File(
@@ -122,8 +106,8 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
+ File.separator
+ targetFileName
+ SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
- SizeTieredCompactionLogger sizeTieredCompactionLogger =
- new SizeTieredCompactionLogger(logFile.getPath());
+ sizeTieredCompactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+
for (TsFileResource resource : selectedTsFileResourceList) {
sizeTieredCompactionLogger.logFileInfo(SOURCE_INFO,
resource.getTsFile());
}
@@ -131,6 +115,7 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
sizeTieredCompactionLogger.logFileInfo(TARGET_INFO,
targetTsFileResource.getTsFile());
LOGGER.info(
"{} [Compaction] compaction with {}", fullStorageGroupName,
selectedTsFileResourceList);
+
// carry out the compaction
InnerSpaceCompactionUtils.compact(
targetTsFileResource, selectedTsFileResourceList,
fullStorageGroupName, true);
@@ -144,25 +129,50 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
throw new InterruptedException(
String.format("%s [Compaction] abort", fullStorageGroupName));
}
+ LOGGER.info(
+ "{} [Compaction] Compacted target files, try to get the write lock
of source files",
+ fullStorageGroupName);
+
+ // release the read lock of all source files, and get the write lock of
them to delete them
+ for (int i = 0; i < selectedTsFileResourceList.size(); ++i) {
+ selectedTsFileResourceList.get(i).readUnlock();
+ isHoldingReadLock[i] = false;
+ selectedTsFileResourceList.get(i).writeLock();
+ isHoldingWriteLock[i] = true;
+ }
+ LOGGER.info(
+ "{} [Compaction] Get the write lock of files, try to get the write
lock of TsFileResourceList",
+ fullStorageGroupName);
+
// get write lock for TsFileResource list with timeout
try {
tsFileManager.writeLockWithTimeout("size-tired compaction", 60_000);
} catch (WriteLockFailedException e) {
- // if current compaction thread couldn't get writelock
+ // if current compaction thread couldn't get write lock
// a WriteLockFailException will be thrown, then terminate the thread
itself
LOGGER.warn(
"{} [SizeTiredCompactionTask] failed to get write lock, abort the
task and delete the target file {}",
fullStorageGroupName,
targetTsFileResource.getTsFile(),
e);
- targetTsFileResource.getTsFile().delete();
- logFile.delete();
throw new InterruptedException(
String.format(
"%s [Compaction] compaction abort because cannot acquire write
lock",
fullStorageGroupName));
}
+
try {
+ LOGGER.info(
+ "{} [SizeTiredCompactionTask] old file deleted, start to rename
mods file",
+ fullStorageGroupName);
+ InnerSpaceCompactionUtils.combineModsInCompaction(
+ selectedTsFileResourceList, targetTsFileResource);
+
+ // delete the old files
+ InnerSpaceCompactionUtils.deleteTsFilesInDisk(
+ selectedTsFileResourceList, fullStorageGroupName);
+ InnerSpaceCompactionUtils.deleteModificationForSourceFile(
+ selectedTsFileResourceList, fullStorageGroupName);
// replace the old files with new file, the new is in same position as
the old
for (TsFileResource resource : selectedTsFileResourceList) {
TsFileResourceManager.getInstance().removeTsFileResource(resource);
@@ -175,13 +185,7 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
} finally {
tsFileManager.writeUnlock();
}
- // delete the old files
- InnerSpaceCompactionUtils.deleteTsFilesInDisk(
- selectedTsFileResourceList, fullStorageGroupName);
- LOGGER.info(
- "{} [SizeTiredCompactionTask] old file deleted, start to rename mods
file",
- fullStorageGroupName);
- combineModsInCompaction(selectedTsFileResourceList,
targetTsFileResource);
+
long costTime = System.currentTimeMillis() - startTime;
LOGGER.info(
"{} [SizeTiredCompactionTask] all compaction task finish, target
file is {},"
@@ -192,9 +196,31 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
if (logFile.exists()) {
logFile.delete();
}
+ } catch (Throwable throwable) {
+ LOGGER.error(
+ "{} [Compaction] Throwable is caught during execution of
SizeTieredCompaction, {}",
+ fullStorageGroupName,
+ throwable);
+ LOGGER.warn("{} [Compaction] Start to handle exception",
fullStorageGroupName);
+ if (sizeTieredCompactionLogger != null) {
+ sizeTieredCompactionLogger.close();
+ }
+ InnerSpaceCompactionExceptionHandler.handleException(
+ fullStorageGroupName,
+ logFile,
+ targetTsFileResource,
+ selectedTsFileResourceList,
+ tsFileManager,
+ tsFileResourceList);
} finally {
- for (TsFileResource resource : selectedTsFileResourceList) {
- resource.setMerging(false);
+ for (int i = 0; i < selectedTsFileResourceList.size(); ++i) {
+ if (isHoldingReadLock[i]) {
+ selectedTsFileResourceList.get(i).readUnlock();
+ }
+ if (isHoldingWriteLock[i]) {
+ selectedTsFileResourceList.get(i).writeUnlock();
+ }
+ selectedTsFileResourceList.get(i).setMerging(false);
}
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/utils/InnerSpaceCompactionUtils.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/utils/InnerSpaceCompactionUtils.java
index 8c8f9bf..af2a114 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/utils/InnerSpaceCompactionUtils.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/utils/InnerSpaceCompactionUtils.java
@@ -54,6 +54,7 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.io.IOException;
import java.math.BigInteger;
+import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
@@ -541,7 +542,7 @@ public class InnerSpaceCompactionUtils {
public static void deleteTsFilesInDisk(
Collection<TsFileResource> mergeTsFiles, String storageGroupName) {
- logger.info("{} [compaction] merge starts to delete real file ",
storageGroupName);
+ logger.info("{} [Compaction] Compaction starts to delete real file ",
storageGroupName);
for (TsFileResource mergeTsFile : mergeTsFiles) {
deleteTsFile(mergeTsFile);
logger.info(
@@ -549,16 +550,56 @@ public class InnerSpaceCompactionUtils {
}
}
+ /** Delete all modification files for source files */
+ public static void deleteModificationForSourceFile(
+ Collection<TsFileResource> sourceFiles, String storageGroupName) throws
IOException {
+ logger.info("{} [Compaction] Start to delete modifications of source
files", storageGroupName);
+ for (TsFileResource tsFileResource : sourceFiles) {
+ ModificationFile compactionModificationFile =
+ ModificationFile.getCompactionMods(tsFileResource);
+ if (compactionModificationFile.exists()) {
+ compactionModificationFile.remove();
+ }
+
+ ModificationFile normalModification =
ModificationFile.getNormalMods(tsFileResource);
+ if (normalModification.exists()) {
+ normalModification.remove();
+ }
+ }
+ }
+
+ /**
+ * Collect all the compaction modification files of source files, and
combines them as the
+ * modification file of target file.
+ */
+ public static void combineModsInCompaction(
+ Collection<TsFileResource> mergeTsFiles, TsFileResource targetTsFile)
throws IOException {
+ List<Modification> modifications = new ArrayList<>();
+ for (TsFileResource mergeTsFile : mergeTsFiles) {
+ try (ModificationFile sourceCompactionModificationFile =
+ ModificationFile.getCompactionMods(mergeTsFile)) {
+
modifications.addAll(sourceCompactionModificationFile.getModifications());
+ }
+ }
+ if (!modifications.isEmpty()) {
+ try (ModificationFile modificationFile =
ModificationFile.getNormalMods(targetTsFile)) {
+ for (Modification modification : modifications) {
+ // we have to set modification offset to MAX_VALUE, as the offset of
source chunk may
+ // change after compaction
+ modification.setFileOffset(Long.MAX_VALUE);
+ modificationFile.write(modification);
+ }
+ }
+ }
+ }
+
public static void deleteTsFile(TsFileResource seqFile) {
- seqFile.writeLock();
try {
FileReaderManager.getInstance().closeFileAndRemoveReader(seqFile.getTsFilePath());
seqFile.setDeleted(true);
seqFile.delete();
} catch (IOException e) {
logger.error(e.getMessage(), e);
- } finally {
- seqFile.writeUnlock();
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileManager.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileManager.java
index 73632c3..c72311e 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileManager.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileManager.java
@@ -57,6 +57,8 @@ public class TsFileManager {
private List<TsFileResource> sequenceRecoverTsFileResources = new
ArrayList<>();
private List<TsFileResource> unsequenceRecoverTsFileResources = new
ArrayList<>();
+ private boolean allowCompaction = true;
+
public TsFileManager(
String storageGroupName, String virtualStorageGroup, String
storageGroupDir) {
this.storageGroupName = storageGroupName;
@@ -276,6 +278,14 @@ public class TsFileManager {
}
}
+ public boolean isAllowCompaction() {
+ return allowCompaction;
+ }
+
+ public void setAllowCompaction(boolean allowCompaction) {
+ this.allowCompaction = allowCompaction;
+ }
+
public String getVirtualStorageGroup() {
return virtualStorageGroup;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
index d39534f..ab2f0e1 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
@@ -347,6 +347,14 @@ public class TsFileResource {
return modFile;
}
+ public void resetModFile() {
+ if (modFile != null) {
+ synchronized (this) {
+ modFile = null;
+ }
+ }
+ }
+
public void setFile(File file) {
this.file = file;
}
@@ -449,26 +457,33 @@ public class TsFileResource {
}
/** Remove the data file, its resource file, and its modification file
physically. */
- public void remove() {
+ public boolean remove() {
try {
fsFactory.deleteIfExists(file);
} catch (IOException e) {
logger.error("TsFile {} cannot be deleted: {}", file, e.getMessage());
+ return false;
+ }
+ if (!removeResourceFile()) {
+ return false;
}
- removeResourceFile();
try {
fsFactory.deleteIfExists(fsFactory.getFile(file.getPath() +
ModificationFile.FILE_SUFFIX));
} catch (IOException e) {
logger.error("ModificationFile {} cannot be deleted: {}", file,
e.getMessage());
+ return false;
}
+ return true;
}
- public void removeResourceFile() {
+ public boolean removeResourceFile() {
try {
fsFactory.deleteIfExists(fsFactory.getFile(file.getPath() +
RESOURCE_SUFFIX));
} catch (IOException e) {
logger.error("TsFileResource {} cannot be deleted: {}", file,
e.getMessage());
+ return false;
}
+ return true;
}
void moveTo(File targetDir) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceList.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceList.java
index 34a7b7f..5ba13fc 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceList.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceList.java
@@ -26,6 +26,7 @@ import
org.apache.iotdb.tsfile.exception.NotImplementedException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.io.IOException;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Iterator;
@@ -189,6 +190,64 @@ public class TsFileResourceList implements
List<TsFileResource> {
}
/**
+ * Insert a tsfile resource to the list, the tsfile will be inserted before
the first tsfile whose
+ * timestamp is greater than its. If there is no tsfile whose timestamp is
greater than the new
+ * node's, the new node will be inserted to the tail of the list.
+ */
+ public boolean keepOrderInsert(TsFileResource newNode) throws IOException {
+ writeLock();
+ try {
+ if (newNode.prev != null || newNode.next != null) {
+ // this node already in a list
+ return false;
+ }
+ if (tail == null) {
+ header = newNode;
+ tail = newNode;
+ count++;
+ } else {
+ // find the position to insert of this node
+ // the list should be ordered by file timestamp
+ long timeOfNewNode =
+
TsFileNameGenerator.getTsFileName(newNode.getTsFile().getName()).getTime();
+
+ if
(TsFileNameGenerator.getTsFileName(header.getTsFile().getName()).getTime()
+ > timeOfNewNode) {
+ // the timestamp of head node is greater than the new node
+ // insert it before the head
+ insertBefore(header, newNode);
+ } else if
(TsFileNameGenerator.getTsFileName(tail.getTsFile().getName()).getTime()
+ < timeOfNewNode) {
+ // the timestamp of new node is greater than the tail node
+ // insert it after the tail
+ insertAfter(tail, newNode);
+ } else {
+ // the timestamp of new node is between the timestamp of head and
tail node
+ // find the first node whose timestamp is greater than new node
+ // and insert the new node before this node
+ TsFileResource currNode = header;
+ while (currNode.next != null) {
+ if
(TsFileNameGenerator.getTsFileName(currNode.getTsFile().getName()).getTime()
+ > timeOfNewNode) {
+ break;
+ }
+ currNode = currNode.next;
+ }
+ if
(TsFileNameGenerator.getTsFileName(currNode.getTsFile().getName()).getTime()
+ < timeOfNewNode) {
+ LOGGER.error("Cannot find an appropriate place to insert {}",
newNode);
+ } else {
+ insertBefore(currNode, newNode);
+ }
+ }
+ }
+ return true;
+ } finally {
+ writeUnlock();
+ }
+ }
+
+ /**
* The tsFileResourceListNode to be removed must be in the list, otherwise
may cause unknown
* behavior
*/
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTest.java
new file mode 100644
index 0000000..d382a71
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/AbstractInnerSpaceCompactionTest.java
@@ -0,0 +1,295 @@
+/*
+ * 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.engine.compaction.inner;
+
+import org.apache.iotdb.db.conf.IoTDBConstant;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.constant.TestConstant;
+import org.apache.iotdb.db.engine.cache.ChunkCache;
+import org.apache.iotdb.db.engine.cache.TimeSeriesMetadataCache;
+import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionRecoverTest;
+import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.db.exception.metadata.MetadataException;
+import org.apache.iotdb.db.metadata.path.PartialPath;
+import org.apache.iotdb.db.query.control.FileReaderManager;
+import org.apache.iotdb.db.service.IoTDB;
+import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
+import org.apache.iotdb.tsfile.read.common.Path;
+import org.apache.iotdb.tsfile.write.TsFileWriter;
+import org.apache.iotdb.tsfile.write.record.TSRecord;
+import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
+import org.apache.iotdb.tsfile.write.schema.UnaryMeasurementSchema;
+
+import org.apache.commons.io.FileUtils;
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+
+import static org.apache.iotdb.db.conf.IoTDBConstant.PATH_SEPARATOR;
+
+public abstract class AbstractInnerSpaceCompactionTest {
+ protected static final Logger logger =
+ LoggerFactory.getLogger(SizeTieredCompactionRecoverTest.class);
+
+ protected File tempSGDir;
+ protected static final String COMPACTION_TEST_SG = "root.compactionTest";
+ protected TsFileManager tsFileManager;
+
+ protected int seqFileNum = 5;
+ protected int unseqFileNum = 0;
+ protected int measurementNum = 10;
+ protected int deviceNum = 10;
+ protected long ptNum = 100;
+ protected long flushInterval = 20;
+ protected TSEncoding encoding = TSEncoding.PLAIN;
+
+ protected static String SEQ_DIRS =
+ TestConstant.BASE_OUTPUT_PATH
+ + "data"
+ + File.separator
+ + "sequence"
+ + File.separator
+ + "root.compactionTest"
+ + File.separator
+ + "0"
+ + File.separator
+ + "0";
+ protected static String UNSEQ_DIRS =
+ TestConstant.BASE_OUTPUT_PATH
+ + "data"
+ + File.separator
+ + "unsequence"
+ + File.separator
+ + "root.compactionTest"
+ + File.separator
+ + "0"
+ + File.separator
+ + "0";
+
+ protected String[] deviceIds;
+ protected UnaryMeasurementSchema[] measurementSchemas;
+
+ protected List<TsFileResource> seqResources = new ArrayList<>();
+ protected List<TsFileResource> unseqResources = new ArrayList<>();
+
+ private int prevMergeChunkThreshold;
+
+ @Before
+ public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ tempSGDir =
+ new File(
+ TestConstant.BASE_OUTPUT_PATH
+ + "data"
+ + File.separator
+ + "sequence"
+ + File.separator
+ + "root.compactionTest"
+ + File.separator
+ + "0"
+ + File.separator
+ + "0");
+ if (!tempSGDir.exists()) {
+ Assert.assertTrue(tempSGDir.mkdirs());
+ }
+ if (!new File(SEQ_DIRS).exists()) {
+ Assert.assertTrue(new File(SEQ_DIRS).mkdirs());
+ }
+ if (!new File(UNSEQ_DIRS).exists()) {
+ Assert.assertTrue(new File(UNSEQ_DIRS).mkdirs());
+ }
+
+ EnvironmentUtils.envSetUp();
+ IoTDB.metaManager.init();
+ prevMergeChunkThreshold =
+
IoTDBDescriptor.getInstance().getConfig().getMergeChunkPointNumberThreshold();
+
IoTDBDescriptor.getInstance().getConfig().setMergeChunkPointNumberThreshold(-1);
+ prepareSeries();
+ prepareFiles(seqFileNum, unseqFileNum);
+ tsFileManager = new TsFileManager(COMPACTION_TEST_SG, "0",
tempSGDir.getAbsolutePath());
+ }
+
+ void prepareSeries() throws MetadataException {
+ measurementSchemas = new UnaryMeasurementSchema[measurementNum];
+ for (int i = 0; i < measurementNum; i++) {
+ measurementSchemas[i] =
+ new UnaryMeasurementSchema(
+ "sensor" + i, TSDataType.DOUBLE, encoding,
CompressionType.UNCOMPRESSED);
+ }
+ deviceIds = new String[deviceNum];
+ for (int i = 0; i < deviceNum; i++) {
+ deviceIds[i] = COMPACTION_TEST_SG + PATH_SEPARATOR + "device" + i;
+ }
+ IoTDB.metaManager.setStorageGroup(new PartialPath(COMPACTION_TEST_SG));
+ for (String device : deviceIds) {
+ for (UnaryMeasurementSchema measurementSchema : measurementSchemas) {
+ PartialPath devicePath = new PartialPath(device);
+ IoTDB.metaManager.createTimeseries(
+ devicePath.concatNode(measurementSchema.getMeasurementId()),
+ measurementSchema.getType(),
+ measurementSchema.getEncodingType(),
+ measurementSchema.getCompressor(),
+ Collections.emptyMap());
+ }
+ }
+ }
+
+ void prepareFiles(int seqFileNum, int unseqFileNum) throws IOException,
WriteProcessException {
+ for (int i = 0; i < seqFileNum; i++) {
+ File file =
+ new File(
+ SEQ_DIRS
+ + File.separator.concat(
+ i
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + i
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + ".tsfile"));
+ TsFileResource tsFileResource = new TsFileResource(file);
+ tsFileResource.setClosed(true);
+ tsFileResource.updatePlanIndexes((long) i);
+ seqResources.add(tsFileResource);
+ prepareFile(tsFileResource, i * ptNum, ptNum, 0);
+ }
+ for (int i = 0; i < unseqFileNum; i++) {
+ File file =
+ new File(
+ UNSEQ_DIRS
+ + File.separator.concat(
+ (10000 + i)
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + (10000 + i)
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + ".tsfile"));
+ TsFileResource tsFileResource = new TsFileResource(file);
+ tsFileResource.setClosed(true);
+ tsFileResource.updatePlanIndexes(i + seqFileNum);
+ unseqResources.add(tsFileResource);
+ prepareFile(tsFileResource, i * ptNum, ptNum * (i + 1) / unseqFileNum,
10000);
+ }
+
+ File file =
+ new File(
+ UNSEQ_DIRS
+ + File.separator.concat(
+ unseqFileNum
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + unseqFileNum
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + IoTDBConstant.FILE_NAME_SEPARATOR
+ + 0
+ + ".tsfile"));
+ TsFileResource tsFileResource = new TsFileResource(file);
+ tsFileResource.setClosed(true);
+ tsFileResource.updatePlanIndexes(seqFileNum + unseqFileNum);
+ unseqResources.add(tsFileResource);
+ prepareFile(tsFileResource, 0, ptNum * unseqFileNum, 20000);
+ }
+
+ void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum,
long valueOffset)
+ throws IOException, WriteProcessException {
+ TsFileWriter fileWriter = new TsFileWriter(tsFileResource.getTsFile());
+ for (String deviceId : deviceIds) {
+ for (UnaryMeasurementSchema measurementSchema : measurementSchemas) {
+ fileWriter.registerTimeseries(new Path(deviceId), measurementSchema);
+ }
+ }
+ for (long i = timeOffset; i < timeOffset + ptNum; i++) {
+ for (int j = 0; j < deviceNum; j++) {
+ TSRecord record = new TSRecord(i, deviceIds[j]);
+ for (int k = 0; k < measurementNum; k++) {
+ record.addTuple(
+ DataPoint.getDataPoint(
+ measurementSchemas[k].getType(),
+ measurementSchemas[k].getMeasurementId(),
+ String.valueOf(i + valueOffset)));
+ }
+ fileWriter.write(record);
+ tsFileResource.updateStartTime(deviceIds[j], i);
+ tsFileResource.updateEndTime(deviceIds[j], i);
+ }
+ if ((i + 1) % flushInterval == 0) {
+ fileWriter.flushAllChunkGroups();
+ }
+ }
+ fileWriter.close();
+ }
+
+ @After
+ public void tearDown() throws IOException, StorageEngineException {
+ removeFiles();
+ seqResources.clear();
+ unseqResources.clear();
+ IoTDBDescriptor.getInstance()
+ .getConfig()
+ .setMergeChunkPointNumberThreshold(prevMergeChunkThreshold);
+ ChunkCache.getInstance().clear();
+ TimeSeriesMetadataCache.getInstance().clear();
+ IoTDB.metaManager.clear();
+ EnvironmentUtils.cleanEnv();
+ if (tempSGDir.exists()) {
+ FileUtils.deleteDirectory(tempSGDir);
+ }
+ }
+
+ private void removeFiles() throws IOException {
+ for (TsFileResource tsFileResource : seqResources) {
+ if (tsFileResource.getTsFile().exists()) {
+ tsFileResource.remove();
+ }
+ }
+ for (TsFileResource tsFileResource : unseqResources) {
+ if (tsFileResource.getTsFile().exists()) {
+ tsFileResource.remove();
+ }
+ }
+ File[] files =
FSFactoryProducer.getFSFactory().listFilesBySuffix("target", ".tsfile");
+ for (File file : files) {
+ file.delete();
+ }
+ File[] resourceFiles =
+ FSFactoryProducer.getFSFactory().listFilesBySuffix("target",
".resource");
+ for (File resourceFile : resourceFiles) {
+ resourceFile.delete();
+ }
+ FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders();
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSeqCompactionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSeqCompactionTest.java
index 69fef0b..15966cd 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSeqCompactionTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSeqCompactionTest.java
@@ -22,7 +22,6 @@ package org.apache.iotdb.db.engine.compaction.inner;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.cache.ChunkCache;
import org.apache.iotdb.db.engine.cache.TimeSeriesMetadataCache;
-import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionTask;
import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger;
import org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils;
@@ -55,7 +54,9 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
-import static
org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils.*;
+import static
org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils.putChunk;
+import static
org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils.putOnePageChunk;
+import static
org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils.putOnePageChunks;
public class InnerSeqCompactionTest {
static final String COMPACTION_TEST_SG = "root.compactionTest";
@@ -224,7 +225,8 @@ public class InnerSeqCompactionTest {
new SizeTieredCompactionLogger("target", COMPACTION_TEST_SG);
InnerSpaceCompactionUtils.compact(
targetTsFileResource, sourceResources, COMPACTION_TEST_SG,
true);
- SizeTieredCompactionTask.combineModsInCompaction(sourceResources,
targetTsFileResource);
+ InnerSpaceCompactionUtils.combineModsInCompaction(
+ sourceResources, targetTsFileResource);
List<TsFileResource> targetTsFileResources = new ArrayList<>();
targetTsFileResources.add(targetTsFileResource);
// check data
@@ -450,7 +452,7 @@ public class InnerSeqCompactionTest {
new SizeTieredCompactionLogger("target", COMPACTION_TEST_SG);
InnerSpaceCompactionUtils.compact(
targetTsFileResource, toMergeResources, COMPACTION_TEST_SG,
true);
- SizeTieredCompactionTask.combineModsInCompaction(
+ InnerSpaceCompactionUtils.combineModsInCompaction(
toMergeResources, targetTsFileResource);
List<TsFileResource> targetTsFileResources = new ArrayList<>();
targetTsFileResources.add(targetTsFileResource);
@@ -727,7 +729,7 @@ public class InnerSeqCompactionTest {
new SizeTieredCompactionLogger("target", COMPACTION_TEST_SG);
InnerSpaceCompactionUtils.compact(
targetTsFileResource, toMergeResources, COMPACTION_TEST_SG,
true);
- SizeTieredCompactionTask.combineModsInCompaction(
+ InnerSpaceCompactionUtils.combineModsInCompaction(
toMergeResources, targetTsFileResource);
List<TsFileResource> targetTsFileResources = new ArrayList<>();
targetTsFileResources.add(targetTsFileResource);
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionTest.java
new file mode 100644
index 0000000..e343ae3
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionExceptionTest.java
@@ -0,0 +1,427 @@
+/*
+ * 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.engine.compaction.inner;
+
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
+import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger;
+import
org.apache.iotdb.db.engine.compaction.utils.CompactionFileGeneratorUtils;
+import org.apache.iotdb.db.engine.modification.Modification;
+import org.apache.iotdb.db.engine.modification.ModificationFile;
+import org.apache.iotdb.db.engine.storagegroup.TsFileNameGenerator;
+import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.tsfile.utils.Pair;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.io.File;
+import java.io.FileOutputStream;
+import java.nio.channels.FileChannel;
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.Map;
+
+public class InnerSpaceCompactionExceptionTest extends
AbstractInnerSpaceCompactionTest {
+
+ /**
+ * Test when all source files exist, and target file is not complete. System
should delete target
+ * file and its resource at this time.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testWhenAllSourceExistsAndTargetNotComplete() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ try (FileOutputStream os = new
FileOutputStream(targetResource.getTsFile(), true);
+ FileChannel channel = os.getChannel()) {
+ channel.truncate(targetResource.getTsFileSize() - 10);
+ }
+ compactionLogger.close();
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertFalse(targetResource.getTsFile().exists());
+ Assert.assertFalse(targetResource.resourceFileExists());
+
+ for (TsFileResource resource : seqResources) {
+ Assert.assertTrue(resource.resourceFileExists());
+ Assert.assertTrue(resource.getTsFile().exists());
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ }
+
+ /**
+ * Test when all source files exist, and target file is complete. System
should delete target file
+ * and its resource at this time.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testWhenAllSourceExistsAndTargetComplete() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ compactionLogger.close();
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertFalse(targetResource.getTsFile().exists());
+ Assert.assertFalse(targetResource.resourceFileExists());
+
+ for (TsFileResource resource : seqResources) {
+ Assert.assertTrue(resource.resourceFileExists());
+ Assert.assertTrue(resource.getTsFile().exists());
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ }
+
+ /**
+ * Test some source files lost and target file is complete. System should
delete source files and
+ * add target file to list.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testWhenSomeSourceLostAndTargetComplete() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ seqResources.get(0).remove();
+ compactionLogger.close();
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertTrue(targetResource.getTsFile().exists());
+ Assert.assertTrue(targetResource.resourceFileExists());
+
+ seqResources.remove(0);
+ for (TsFileResource resource : seqResources) {
+ Assert.assertFalse(resource.resourceFileExists());
+ Assert.assertFalse(resource.getTsFile().exists());
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ Assert.assertEquals(1,
tsFileManager.getSequenceListByTimePartition(0).size());
+ Assert.assertEquals(targetResource,
tsFileManager.getSequenceListByTimePartition(0).get(0));
+ }
+
+ /**
+ * Test some source files are lost and target file is not complete. System
should be set to read
+ * only at this time.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testWhenSomeSourceLostAndTargetNotComplete() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ seqResources.get(0).remove();
+ try (FileOutputStream os = new
FileOutputStream(targetResource.getTsFile(), true);
+ FileChannel channel = os.getChannel()) {
+ channel.truncate(targetResource.getTsFileSize() - 10);
+ }
+ compactionLogger.close();
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertTrue(targetResource.getTsFile().exists());
+ Assert.assertTrue(targetResource.resourceFileExists());
+
+ seqResources.remove(0);
+ for (TsFileResource resource : seqResources) {
+ Assert.assertTrue(resource.resourceFileExists());
+ Assert.assertTrue(resource.getTsFile().exists());
+ }
+
+ Assert.assertFalse(tsFileManager.isAllowCompaction());
+ Assert.assertTrue(IoTDBDescriptor.getInstance().getConfig().isReadOnly());
+ IoTDBDescriptor.getInstance().getConfig().setReadOnly(false);
+ }
+
+ /**
+ * Test some source files are lost, and there are some compaction mods for
source files. System
+ * should collect these mods together and write them to mods file of target
file.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testHandleWithCompactionMods() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ for (int i = 0; i < seqResources.size(); i++) {
+ Map<String, Pair<Long, Long>> deleteMap = new HashMap<>();
+ deleteMap.put(
+ deviceIds[0] + "." + measurementSchemas[0].getMeasurementId(),
+ new Pair<>(i * ptNum, i * ptNum + 10));
+ CompactionFileGeneratorUtils.generateMods(deleteMap,
seqResources.get(i), true);
+ }
+
+ seqResources.get(0).remove();
+ compactionLogger.close();
+
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertTrue(targetResource.getTsFile().exists());
+ Assert.assertTrue(targetResource.resourceFileExists());
+ Assert.assertTrue(targetResource.getModFile().exists());
+ Collection<Modification> modifications =
targetResource.getModFile().getModifications();
+ Assert.assertEquals(seqResources.size(), modifications.size());
+ for (Modification modification : modifications) {
+ Assert.assertEquals(deviceIds[0], modification.getDevice());
+ Assert.assertEquals(measurementSchemas[0].getMeasurementId(),
modification.getMeasurement());
+ Assert.assertEquals(Long.MAX_VALUE, modification.getFileOffset());
+ }
+
+ seqResources.remove(0);
+ for (TsFileResource resource : seqResources) {
+ Assert.assertFalse(resource.resourceFileExists());
+ Assert.assertFalse(resource.getTsFile().exists());
+ Assert.assertFalse(resource.getModFile().exists());
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ }
+
+ /**
+ * Test some source files are lost, and there are some mods file for source
files. System should
+ * collect them and generate a new mods file for target file.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testHandleWithNormalMods() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+
+ for (int i = 0; i < seqResources.size(); i++) {
+ Map<String, Pair<Long, Long>> deleteMap = new HashMap<>();
+ deleteMap.put(
+ deviceIds[0] + "." + measurementSchemas[0].getMeasurementId(),
+ new Pair<>(i * ptNum, i * ptNum + 10));
+ CompactionFileGeneratorUtils.generateMods(deleteMap,
seqResources.get(i), false);
+ }
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+
+ seqResources.get(0).remove();
+ compactionLogger.close();
+
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertTrue(targetResource.getTsFile().exists());
+ Assert.assertTrue(targetResource.resourceFileExists());
+ Assert.assertFalse(targetResource.getModFile().exists());
+
+ seqResources.remove(0);
+ for (TsFileResource resource : seqResources) {
+ Assert.assertFalse(resource.resourceFileExists());
+ Assert.assertFalse(resource.getTsFile().exists());
+ Assert.assertFalse(resource.getModFile().exists());
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ }
+
+ /**
+ * Test source files exists, target file is not complete, and there are mods
and compaction mods
+ * for source files. System should remove target file, and combine the
compaction mods and normal
+ * mods together for each source file.
+ *
+ * @throws Exception
+ */
+ @Test
+ public void testHandleWithCompactionModsAndNormalMods() throws Exception {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ TsFileResource targetResource =
+ new TsFileResource(new
File(seqResources.get(0).getTsFile().getParent(), targetFileName));
+ File logFile =
+ new File(
+ targetResource.getTsFile().getPath() +
SizeTieredCompactionLogger.COMPACTION_LOG_NAME);
+ SizeTieredCompactionLogger compactionLogger = new
SizeTieredCompactionLogger(logFile.getPath());
+ for (TsFileResource source : seqResources) {
+ compactionLogger.logFileInfo(SizeTieredCompactionLogger.SOURCE_INFO,
source.getTsFile());
+ }
+ compactionLogger.logSequence(true);
+ compactionLogger.logFileInfo(
+ SizeTieredCompactionLogger.TARGET_INFO, targetResource.getTsFile());
+ for (int i = 0; i < seqResources.size(); i++) {
+ Map<String, Pair<Long, Long>> deleteMap = new HashMap<>();
+ deleteMap.put(
+ deviceIds[0] + "." + measurementSchemas[0].getMeasurementId(),
+ new Pair<>(i * ptNum, i * ptNum + 5));
+ CompactionFileGeneratorUtils.generateMods(deleteMap,
seqResources.get(i), false);
+ }
+ InnerSpaceCompactionUtils.compact(targetResource, seqResources,
COMPACTION_TEST_SG, true);
+ for (int i = 0; i < seqResources.size(); i++) {
+ Map<String, Pair<Long, Long>> deleteMap = new HashMap<>();
+ deleteMap.put(
+ deviceIds[0] + "." + measurementSchemas[0].getMeasurementId(),
+ new Pair<>(i * ptNum + 10, i * ptNum + 15));
+ CompactionFileGeneratorUtils.generateMods(deleteMap,
seqResources.get(i), true);
+ }
+ compactionLogger.close();
+
+ InnerSpaceCompactionExceptionHandler.handleException(
+ COMPACTION_TEST_SG,
+ logFile,
+ targetResource,
+ seqResources,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0));
+ Assert.assertFalse(targetResource.getTsFile().exists());
+ Assert.assertFalse(targetResource.resourceFileExists());
+ Assert.assertFalse(targetResource.getModFile().exists());
+
+ for (TsFileResource resource : seqResources) {
+ Assert.assertTrue(resource.resourceFileExists());
+ Assert.assertTrue(resource.getTsFile().exists());
+ Assert.assertTrue(resource.getModFile().exists());
+
Assert.assertFalse(ModificationFile.getCompactionMods(resource).exists());
+ Collection<Modification> modifications =
resource.getModFile().getModifications();
+ Assert.assertEquals(2, modifications.size());
+ for (Modification modification : modifications) {
+ Assert.assertEquals(deviceIds[0], modification.getDevice());
+ Assert.assertEquals(
+ measurementSchemas[0].getMeasurementId(),
modification.getMeasurement());
+ Assert.assertEquals(Long.MAX_VALUE, modification.getFileOffset());
+ }
+ }
+
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerUnseqCompactionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerUnseqCompactionTest.java
index ff67dde..cbe519c 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerUnseqCompactionTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerUnseqCompactionTest.java
@@ -21,7 +21,6 @@ package org.apache.iotdb.db.engine.compaction.inner;
import org.apache.iotdb.db.engine.cache.ChunkCache;
import org.apache.iotdb.db.engine.cache.TimeSeriesMetadataCache;
-import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionTask;
import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger;
import org.apache.iotdb.db.engine.compaction.utils.CompactionCheckerUtils;
@@ -354,7 +353,7 @@ public class InnerUnseqCompactionTest {
new SizeTieredCompactionLogger("target", COMPACTION_TEST_SG);
InnerSpaceCompactionUtils.compact(
targetTsFileResource, toMergeResources, COMPACTION_TEST_SG,
false);
- SizeTieredCompactionTask.combineModsInCompaction(
+ InnerSpaceCompactionUtils.combineModsInCompaction(
toMergeResources, targetTsFileResource);
List<TsFileResource> targetTsFileResources = new ArrayList<>();
targetTsFileResources.add(targetTsFileResource);
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionHandleExceptionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionHandleExceptionTest.java
new file mode 100644
index 0000000..0371c9d
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionHandleExceptionTest.java
@@ -0,0 +1,197 @@
+/*
+ * 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.engine.compaction.inner.sizetiered;
+
+import
org.apache.iotdb.db.engine.compaction.inner.AbstractInnerSpaceCompactionTest;
+import org.apache.iotdb.db.engine.storagegroup.TsFileNameGenerator;
+import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.db.exception.metadata.MetadataException;
+import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.io.File;
+import java.io.FileOutputStream;
+import java.io.IOException;
+import java.nio.channels.FileChannel;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class SizeTieredCompactionHandleExceptionTest extends
AbstractInnerSpaceCompactionTest {
+ @Before
+ public void setUp() throws IOException, MetadataException,
WriteProcessException {
+ this.seqFileNum = 10;
+ super.setUp();
+ }
+
+ @After
+ public void tearDown() throws StorageEngineException, IOException {
+ super.tearDown();
+ }
+
+ @Test
+ public void testHandleExceptionTargetCompleteAndSourceExists() {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ SizeTieredCompactionTask task =
+ new SizeTieredCompactionTask(
+ COMPACTION_TEST_SG,
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ try {
+ new Thread(
+ () -> {
+ try {
+ task.call();
+ } catch (Exception e) {
+
+ }
+ })
+ .start();
+ Thread.sleep(65_000);
+ } catch (Exception e) {
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ Assert.assertEquals(10, tsFileManager.getTsFileList(true).size());
+ }
+
+ @Test
+ public void testHandleExceptionTargetNotCompleteAndSourceNotExists() {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ SizeTieredCompactionTask task =
+ new SizeTieredCompactionTask(
+ COMPACTION_TEST_SG,
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ try {
+ seqResources.get(seqResources.size() - 1).remove();
+ new Thread(
+ () -> {
+ try {
+ task.call();
+ } catch (Exception e) {
+
+ }
+ })
+ .start();
+ Thread.sleep(5_000);
+ } catch (Exception e) {
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ Assert.assertFalse(tsFileManager.isAllowCompaction());
+ Assert.assertEquals(10, tsFileManager.getTsFileList(true).size());
+ }
+
+ @Test
+ public void testHandleExceptionTargetCompleteAndSourceNotExists() {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ SizeTieredCompactionTask task =
+ new SizeTieredCompactionTask(
+ COMPACTION_TEST_SG,
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ try {
+ new Thread(
+ () -> {
+ try {
+ task.call();
+ } catch (Exception e) {
+
+ }
+ })
+ .start();
+ Thread.sleep(10_000);
+ seqResources.get(0).remove();
+ tsFileManager.getTsFileList(true).remove(seqResources.get(0));
+ Thread.sleep(60_000);
+ } catch (Exception e) {
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ Assert.assertEquals(1, tsFileManager.getTsFileList(true).size());
+ }
+
+ @Test
+ public void testHandleExceptionTargetNotCompleteAndSourceExists() {
+ tsFileManager.addAll(seqResources, true);
+ tsFileManager.addAll(unseqResources, false);
+ SizeTieredCompactionTask task =
+ new SizeTieredCompactionTask(
+ COMPACTION_TEST_SG,
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ try {
+ new Thread(
+ () -> {
+ try {
+ task.call();
+ } catch (Exception e) {
+
+ }
+ })
+ .start();
+ Thread.sleep(10_000);
+ String targetFileName =
+ TsFileNameGenerator.getInnerCompactionFileName(seqResources,
true).getName();
+ File targetFile = new File(seqResources.get(0).getTsFile().getParent(),
targetFileName);
+ FileChannel channel = new FileOutputStream(targetFile,
true).getChannel();
+ channel.truncate(10);
+ channel.close();
+ Thread.sleep(60_000);
+ } catch (Exception e) {
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ Assert.assertTrue(tsFileManager.isAllowCompaction());
+ Assert.assertEquals(10, tsFileManager.getTsFileList(true).size());
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java
index 1862654..58c9935 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionRecoverTest.java
@@ -20,11 +20,8 @@
package org.apache.iotdb.db.engine.compaction.inner.sizetiered;
import org.apache.iotdb.db.conf.IoTDBConstant;
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import org.apache.iotdb.db.constant.TestConstant;
-import org.apache.iotdb.db.engine.cache.ChunkCache;
-import org.apache.iotdb.db.engine.cache.TimeSeriesMetadataCache;
import org.apache.iotdb.db.engine.compaction.CompactionTaskManager;
+import
org.apache.iotdb.db.engine.compaction.inner.AbstractInnerSpaceCompactionTest;
import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUtils;
import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger;
import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
@@ -33,34 +30,19 @@ import
org.apache.iotdb.db.engine.storagegroup.TsFileResource;
import org.apache.iotdb.db.exception.StorageEngineException;
import org.apache.iotdb.db.exception.metadata.MetadataException;
import org.apache.iotdb.db.metadata.path.MeasurementPath;
-import org.apache.iotdb.db.metadata.path.PartialPath;
-import org.apache.iotdb.db.query.control.FileReaderManager;
import org.apache.iotdb.db.query.reader.series.SeriesRawDataBatchReader;
-import org.apache.iotdb.db.service.IoTDB;
import org.apache.iotdb.db.utils.EnvironmentUtils;
import org.apache.iotdb.db.utils.SchemaTestUtils;
import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
-import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
import org.apache.iotdb.tsfile.read.common.BatchData;
-import org.apache.iotdb.tsfile.read.common.Path;
import org.apache.iotdb.tsfile.read.reader.IBatchReader;
-import org.apache.iotdb.tsfile.write.TsFileWriter;
-import org.apache.iotdb.tsfile.write.record.TSRecord;
-import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
-import org.apache.iotdb.tsfile.write.schema.UnaryMeasurementSchema;
import org.apache.iotdb.tsfile.write.writer.TsFileOutput;
-import org.apache.commons.io.FileUtils;
import org.junit.After;
-import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import java.io.BufferedReader;
import java.io.BufferedWriter;
@@ -69,247 +51,23 @@ import java.io.FileReader;
import java.io.FileWriter;
import java.io.IOException;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
-import static org.apache.iotdb.db.conf.IoTDBConstant.PATH_SEPARATOR;
import static
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger.COMPACTION_LOG_NAME;
import static
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger.SOURCE_INFO;
import static
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLogger.TARGET_INFO;
import static org.junit.Assert.assertEquals;
-public class SizeTieredCompactionRecoverTest {
- private static final Logger logger =
- LoggerFactory.getLogger(SizeTieredCompactionRecoverTest.class);
-
- File tempSGDir;
- protected static final String COMPACTION_TEST_SG = "root.compactionTest";
- protected TsFileManager tsFileManager;
-
- protected int seqFileNum = 5;
- protected int unseqFileNum = 0;
- protected int measurementNum = 10;
- protected int deviceNum = 10;
- protected long ptNum = 100;
- protected long flushInterval = 20;
- protected TSEncoding encoding = TSEncoding.PLAIN;
-
- static String SEQ_DIRS =
- TestConstant.BASE_OUTPUT_PATH
- + "data"
- + File.separator
- + "sequence"
- + File.separator
- + "root.compactionTest"
- + File.separator
- + "0"
- + File.separator
- + "0";
- static String UNSEQ_DIRS =
- TestConstant.BASE_OUTPUT_PATH
- + "data"
- + File.separator
- + "unsequence"
- + File.separator
- + "root.compactionTest"
- + File.separator
- + "0"
- + File.separator
- + "0";
-
- protected String[] deviceIds;
- protected UnaryMeasurementSchema[] measurementSchemas;
-
- protected List<TsFileResource> seqResources = new ArrayList<>();
- protected List<TsFileResource> unseqResources = new ArrayList<>();
-
- private int prevMergeChunkThreshold;
+public class SizeTieredCompactionRecoverTest extends
AbstractInnerSpaceCompactionTest {
@Before
public void setUp() throws IOException, WriteProcessException,
MetadataException {
- tempSGDir =
- new File(
- TestConstant.BASE_OUTPUT_PATH
- + "data"
- + File.separator
- + "sequence"
- + File.separator
- + "root.compactionTest"
- + File.separator
- + "0"
- + File.separator
- + "0");
- if (!tempSGDir.exists()) {
- Assert.assertTrue(tempSGDir.mkdirs());
- }
- if (!new File(SEQ_DIRS).exists()) {
- Assert.assertTrue(new File(SEQ_DIRS).mkdirs());
- }
- if (!new File(UNSEQ_DIRS).exists()) {
- Assert.assertTrue(new File(UNSEQ_DIRS).mkdirs());
- }
-
- EnvironmentUtils.envSetUp();
- IoTDB.metaManager.init();
- prevMergeChunkThreshold =
-
IoTDBDescriptor.getInstance().getConfig().getMergeChunkPointNumberThreshold();
-
IoTDBDescriptor.getInstance().getConfig().setMergeChunkPointNumberThreshold(-1);
- prepareSeries();
- prepareFiles(seqFileNum, unseqFileNum);
- tsFileManager = new TsFileManager(COMPACTION_TEST_SG, "0",
tempSGDir.getAbsolutePath());
- }
-
- void prepareSeries() throws MetadataException {
- measurementSchemas = new UnaryMeasurementSchema[measurementNum];
- for (int i = 0; i < measurementNum; i++) {
- measurementSchemas[i] =
- new UnaryMeasurementSchema(
- "sensor" + i, TSDataType.DOUBLE, encoding,
CompressionType.UNCOMPRESSED);
- }
- deviceIds = new String[deviceNum];
- for (int i = 0; i < deviceNum; i++) {
- deviceIds[i] = COMPACTION_TEST_SG + PATH_SEPARATOR + "device" + i;
- }
- IoTDB.metaManager.setStorageGroup(new PartialPath(COMPACTION_TEST_SG));
- for (String device : deviceIds) {
- for (UnaryMeasurementSchema measurementSchema : measurementSchemas) {
- PartialPath devicePath = new PartialPath(device);
- IoTDB.metaManager.createTimeseries(
- devicePath.concatNode(measurementSchema.getMeasurementId()),
- measurementSchema.getType(),
- measurementSchema.getEncodingType(),
- measurementSchema.getCompressor(),
- Collections.emptyMap());
- }
- }
- }
-
- void prepareFiles(int seqFileNum, int unseqFileNum) throws IOException,
WriteProcessException {
- for (int i = 0; i < seqFileNum; i++) {
- File file =
- new File(
- SEQ_DIRS
- + File.separator.concat(
- i
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + i
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + ".tsfile"));
- TsFileResource tsFileResource = new TsFileResource(file);
- tsFileResource.setClosed(true);
- tsFileResource.updatePlanIndexes((long) i);
- seqResources.add(tsFileResource);
- prepareFile(tsFileResource, i * ptNum, ptNum, 0);
- }
- for (int i = 0; i < unseqFileNum; i++) {
- File file =
- new File(
- UNSEQ_DIRS
- + File.separator.concat(
- (10000 + i)
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + (10000 + i)
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + ".tsfile"));
- TsFileResource tsFileResource = new TsFileResource(file);
- tsFileResource.setClosed(true);
- tsFileResource.updatePlanIndexes(i + seqFileNum);
- unseqResources.add(tsFileResource);
- prepareFile(tsFileResource, i * ptNum, ptNum * (i + 1) / unseqFileNum,
10000);
- }
-
- File file =
- new File(
- UNSEQ_DIRS
- + File.separator.concat(
- unseqFileNum
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + unseqFileNum
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + IoTDBConstant.FILE_NAME_SEPARATOR
- + 0
- + ".tsfile"));
- TsFileResource tsFileResource = new TsFileResource(file);
- tsFileResource.setClosed(true);
- tsFileResource.updatePlanIndexes(seqFileNum + unseqFileNum);
- unseqResources.add(tsFileResource);
- prepareFile(tsFileResource, 0, ptNum * unseqFileNum, 20000);
- }
-
- void prepareFile(TsFileResource tsFileResource, long timeOffset, long ptNum,
long valueOffset)
- throws IOException, WriteProcessException {
- TsFileWriter fileWriter = new TsFileWriter(tsFileResource.getTsFile());
- for (String deviceId : deviceIds) {
- for (UnaryMeasurementSchema measurementSchema : measurementSchemas) {
- fileWriter.registerTimeseries(new Path(deviceId), measurementSchema);
- }
- }
- for (long i = timeOffset; i < timeOffset + ptNum; i++) {
- for (int j = 0; j < deviceNum; j++) {
- TSRecord record = new TSRecord(i, deviceIds[j]);
- for (int k = 0; k < measurementNum; k++) {
- record.addTuple(
- DataPoint.getDataPoint(
- measurementSchemas[k].getType(),
- measurementSchemas[k].getMeasurementId(),
- String.valueOf(i + valueOffset)));
- }
- fileWriter.write(record);
- tsFileResource.updateStartTime(deviceIds[j], i);
- tsFileResource.updateEndTime(deviceIds[j], i);
- }
- if ((i + 1) % flushInterval == 0) {
- fileWriter.flushAllChunkGroups();
- }
- }
- fileWriter.close();
+ super.setUp();
}
@After
public void tearDown() throws IOException, StorageEngineException {
- removeFiles();
- seqResources.clear();
- unseqResources.clear();
- IoTDBDescriptor.getInstance()
- .getConfig()
- .setMergeChunkPointNumberThreshold(prevMergeChunkThreshold);
- ChunkCache.getInstance().clear();
- TimeSeriesMetadataCache.getInstance().clear();
- IoTDB.metaManager.clear();
- EnvironmentUtils.cleanEnv();
- if (tempSGDir.exists()) {
- FileUtils.deleteDirectory(tempSGDir);
- }
- }
-
- private void removeFiles() throws IOException {
- for (TsFileResource tsFileResource : seqResources) {
- if (tsFileResource.getTsFile().exists()) {
- tsFileResource.remove();
- }
- }
- for (TsFileResource tsFileResource : unseqResources) {
- if (tsFileResource.getTsFile().exists()) {
- tsFileResource.remove();
- }
- }
- File[] files =
FSFactoryProducer.getFSFactory().listFilesBySuffix("target", ".tsfile");
- for (File file : files) {
- file.delete();
- }
- File[] resourceFiles =
- FSFactoryProducer.getFSFactory().listFilesBySuffix("target",
".resource");
- for (File resourceFile : resourceFiles) {
- resourceFile.delete();
- }
- FileReaderManager.getInstance().closeAndRemoveAllOpenedReaders();
+ super.tearDown();
}
/** Target file uncompleted, source files and log exists */
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceListTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceListTest.java
index a63137a..15edccc 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceListTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/TsFileResourceListTest.java
@@ -65,6 +65,36 @@ public class TsFileResourceListTest {
}
@Test
+ public void testKeepOrderInsert() throws Exception {
+ TsFileResourceList tsFileResourceList = new TsFileResourceList();
+ TsFileResource resource1 = generateTsFileResource(10);
+ // 10
+ tsFileResourceList.keepOrderInsert(resource1);
+ Assert.assertEquals(resource1, tsFileResourceList.get(0));
+ Assert.assertEquals(1, tsFileResourceList.size());
+ TsFileResource resource2 = generateTsFileResource(100);
+ // 10 100
+ tsFileResourceList.keepOrderInsert(resource2);
+ Assert.assertEquals(resource2, tsFileResourceList.get(1));
+ Assert.assertEquals(2, tsFileResourceList.size());
+ TsFileResource resource3 = generateTsFileResource(50);
+ // 10 50 100
+ tsFileResourceList.keepOrderInsert(resource3);
+ Assert.assertEquals(resource3, tsFileResourceList.get(1));
+ Assert.assertEquals(3, tsFileResourceList.size());
+ TsFileResource resource4 = generateTsFileResource(75);
+ // 10 50 75 100
+ tsFileResourceList.keepOrderInsert(resource4);
+ Assert.assertEquals(resource4, tsFileResourceList.get(2));
+ Assert.assertEquals(4, tsFileResourceList.size());
+ TsFileResource resource5 = generateTsFileResource(5);
+ // 5 10 50 75 100
+ tsFileResourceList.keepOrderInsert(resource5);
+ Assert.assertEquals(resource5, tsFileResourceList.get(0));
+ Assert.assertEquals(5, tsFileResourceList.size());
+ }
+
+ @Test
public void testRemove() {
TsFileResourceList tsFileResourceList = new TsFileResourceList();
List<TsFileResource> tsFileResources = new ArrayList<>();