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 f5a7b99 [IOTDB-1899] Fix stream closed exception during compaction
(#4457)
f5a7b99 is described below
commit f5a7b990d902402b0d738ca18f0f587330db87e7
Author: liuxuxin <[email protected]>
AuthorDate: Sat Nov 27 09:59:38 2021 +0800
[IOTDB-1899] Fix stream closed exception during compaction (#4457)
---
.../db/engine/compaction/CompactionScheduler.java | 3 +-
.../engine/compaction/CompactionTaskManager.java | 36 ++-
.../inner/AbstractInnerSpaceCompactionTask.java | 2 +
.../sizetiered/SizeTieredCompactionSelector.java | 2 +-
.../inner/sizetiered/SizeTieredCompactionTask.java | 105 ++++-----
.../compaction/task/AbstractCompactionTask.java | 2 +
.../compaction/task/CompactionRecoverTask.java | 2 +-
.../engine/compaction/CompactionSchedulerTest.java | 18 +-
.../compaction/CompactionTaskManagerTest.java | 252 +++++++++++++++++++++
.../compaction/inner/InnerCompactionChunkTest.java | 4 +-
.../compaction/inner/InnerCompactionLogTest.java | 4 +-
.../inner/InnerCompactionMoreDataTest.java | 2 +-
.../inner/InnerCompactionSchedulerTest.java | 4 +-
.../compaction/inner/InnerCompactionTest.java | 2 +-
.../inner/InnerSpaceCompactionUtilsTest.java | 4 +-
.../task/FakedInnerSpaceCompactionTask.java | 44 ++--
.../engine/storagegroup/FakedTsFileResource.java | 11 +
.../storagegroup/StorageGroupProcessorTest.java | 2 +-
.../db/integration/IoTDBNewTsFileCompactionIT.java | 2 +-
19 files changed, 394 insertions(+), 107 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 fb2ecf5..6bccf4d 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
@@ -106,7 +106,8 @@ public class CompactionScheduler {
boolean taskSubmitted = true;
int concurrentCompactionThread = config.getConcurrentCompactionThread();
while (taskSubmitted
- && CompactionTaskManager.getInstance().getTaskCount() <
concurrentCompactionThread) {
+ && CompactionTaskManager.getInstance().getExecutingTaskCount()
+ < concurrentCompactionThread) {
taskSubmitted =
tryToSubmitInnerSpaceCompactionTask(
logicalStorageGroupName,
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java
index 8f85d2c..959b9fc 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManager.java
@@ -32,9 +32,11 @@ import com.google.common.collect.MinMaxPriorityQueue;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
+import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Callable;
@@ -53,11 +55,12 @@ public class CompactionTaskManager implements IService {
private WrappedScheduledExecutorService taskExecutionPool;
public static volatile AtomicInteger currentTaskNum = new AtomicInteger(0);
// TODO: record the task in time partition
- private MinMaxPriorityQueue<AbstractCompactionTask> compactionTaskQueue =
+ private MinMaxPriorityQueue<AbstractCompactionTask>
candidateCompactionTaskQueue =
MinMaxPriorityQueue.orderedBy(new
CompactionTaskComparator()).maximumSize(1000).create();
private Map<String, Set<Future<Void>>> storageGroupTasks = new
ConcurrentHashMap<>();
private Map<String, Map<Long, Set<Future<Void>>>> compactionTaskFutures =
new ConcurrentHashMap<>();
+ private List<AbstractCompactionTask> runningCompactionTaskList = new
ArrayList<>();
private ScheduledExecutorService compactionTaskSubmissionThreadPool;
private final long TASK_SUBMIT_INTERVAL =
IoTDBDescriptor.getInstance().getConfig().getCompactionSubmissionInterval();
@@ -172,13 +175,9 @@ public class CompactionTaskManager implements IService {
* with last priority will be removed from the task.
*/
public synchronized boolean addTaskToWaitingQueue(AbstractCompactionTask
compactionTask) {
- if (!compactionTaskQueue.contains(compactionTask)) {
- logger.debug(
- "Add a compaction task {} to queue, current queue size is {},
current task num is {}",
- compactionTask,
- compactionTaskQueue.size(),
- currentTaskNum.get());
- compactionTaskQueue.add(compactionTask);
+ if (!candidateCompactionTaskQueue.contains(compactionTask)
+ && !runningCompactionTaskList.contains(compactionTask)) {
+ candidateCompactionTaskQueue.add(compactionTask);
return true;
}
return false;
@@ -191,14 +190,19 @@ public class CompactionTaskManager implements IService {
public synchronized void submitTaskFromTaskQueue() {
while (currentTaskNum.get()
<
IoTDBDescriptor.getInstance().getConfig().getConcurrentCompactionThread()
- && compactionTaskQueue.size() > 0) {
- AbstractCompactionTask task = compactionTaskQueue.poll();
- if (task.checkValidAndSetMerging()) {
+ && candidateCompactionTaskQueue.size() > 0) {
+ AbstractCompactionTask task = candidateCompactionTaskQueue.poll();
+ if (task != null && task.checkValidAndSetMerging()) {
submitTask(task.getFullStorageGroupName(), task.getTimePartition(),
task);
+ runningCompactionTaskList.add(task);
}
}
}
+ public synchronized void removeRunningTaskFromList(AbstractCompactionTask
task) {
+ runningCompactionTaskList.remove(task);
+ }
+
/**
* This method will directly submit a task to thread pool if there is
available thread.
*
@@ -240,10 +244,18 @@ public class CompactionTaskManager implements IService {
}
}
- public int getTaskCount() {
+ public int getExecutingTaskCount() {
return taskExecutionPool.getActiveCount() +
taskExecutionPool.getQueue().size();
}
+ public int getTotalTaskCount() {
+ return getExecutingTaskCount() + candidateCompactionTaskQueue.size();
+ }
+
+ public synchronized List<AbstractCompactionTask>
getRunningCompactionTaskList() {
+ return new ArrayList<>(runningCompactionTaskList);
+ }
+
public long getFinishTaskNum() {
return taskExecutionPool.getCompletedTaskCount();
}
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 38d2883..fa1e830 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
@@ -124,6 +124,8 @@ public abstract class AbstractInnerSpaceCompactionTask
extends AbstractCompactio
.append(timePartition)
.append(" task file num is ")
.append(selectedTsFileResourceList.size())
+ .append(", files is ")
+ .append(selectedTsFileResourceList)
.append(", total compaction count is ")
.append(sumOfCompactionCount)
.toString();
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionSelector.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionSelector.java
index c3c416a..db665bd 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionSelector.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/inner/sizetiered/SizeTieredCompactionSelector.java
@@ -82,7 +82,7 @@ public class SizeTieredCompactionSelector extends
AbstractInnerSpaceCompactionSe
IoTDBDescriptor.getInstance().getConfig().getTargetCompactionFileSize(),
IoTDBDescriptor.getInstance().getConfig().getMaxCompactionCandidateFileNum(),
CompactionTaskManager.currentTaskNum.get(),
- CompactionTaskManager.getInstance().getTaskCount(),
+ CompactionTaskManager.getInstance().getExecutingTaskCount(),
IoTDBDescriptor.getInstance().getConfig().getConcurrentCompactionThread());
tsFileResources.readLock();
PriorityQueue<Pair<List<TsFileResource>, Long>> taskPriorityQueue =
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 6f85f21..0661da8 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
@@ -137,62 +137,65 @@ public class SizeTieredCompactionTask extends
AbstractInnerSpaceCompactionTask {
LOGGER.info(
"{} [SizeTiredCompactionTask] compact finish, close the logger",
fullStorageGroupName);
sizeTieredCompactionLogger.close();
- } finally {
- for (TsFileResource resource : selectedTsFileResourceList) {
- resource.setMerging(false);
+
+ LOGGER.info(
+ "{} [Compaction] compaction finish, start to delete old files",
fullStorageGroupName);
+ if (Thread.currentThread().isInterrupted()) {
+ throw new InterruptedException(
+ String.format("%s [Compaction] abort", fullStorageGroupName));
}
- }
- LOGGER.info(
- "{} [Compaction] compaction finish, start to delete old files",
fullStorageGroupName);
- if (Thread.currentThread().isInterrupted()) {
- throw new InterruptedException(String.format("%s [Compaction] abort",
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
- // 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 {}",
+ // 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
+ // 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 {
+ // replace the old files with new file, the new is in same position as
the old
+ for (TsFileResource resource : selectedTsFileResourceList) {
+ TsFileResourceManager.getInstance().removeTsFileResource(resource);
+ }
+ tsFileResourceList.insertBefore(selectedTsFileResourceList.get(0),
targetTsFileResource);
+
TsFileResourceManager.getInstance().registerSealedTsFileResource(targetTsFileResource);
+ for (TsFileResource resource : selectedTsFileResourceList) {
+ tsFileResourceList.remove(resource);
+ }
+ } 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 {},"
+ + "time cost is {} s",
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 {
- // replace the old files with new file, the new is in same position as
the old
- for (TsFileResource resource : selectedTsFileResourceList) {
- TsFileResourceManager.getInstance().removeTsFileResource(resource);
+ targetFileName,
+ costTime / 1000);
+ if (logFile.exists()) {
+ logFile.delete();
}
- tsFileResourceList.insertBefore(selectedTsFileResourceList.get(0),
targetTsFileResource);
-
TsFileResourceManager.getInstance().registerSealedTsFileResource(targetTsFileResource);
+ } finally {
for (TsFileResource resource : selectedTsFileResourceList) {
- tsFileResourceList.remove(resource);
+ resource.setMerging(false);
}
- } 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 {},"
- + "time cost is {} s",
- fullStorageGroupName,
- targetFileName,
- costTime / 1000);
- if (logFile.exists()) {
- logFile.delete();
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/AbstractCompactionTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/AbstractCompactionTask.java
index b1d9128..4bf4d46 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/AbstractCompactionTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/AbstractCompactionTask.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.engine.compaction.task;
import org.apache.iotdb.db.engine.compaction.CompactionScheduler;
+import org.apache.iotdb.db.engine.compaction.CompactionTaskManager;
import
org.apache.iotdb.db.engine.compaction.cross.inplace.InplaceCompactionRecoverTask;
import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionRecoverTask;
@@ -61,6 +62,7 @@ public abstract class AbstractCompactionTask implements
Callable<Void> {
if (!(this instanceof InplaceCompactionRecoverTask)
&& !(this instanceof SizeTieredCompactionRecoverTask)) {
CompactionScheduler.decPartitionCompaction(fullStorageGroupName,
timePartition);
+ CompactionTaskManager.getInstance().removeRunningTaskFromList(this);
}
this.currentTaskNum.decrementAndGet();
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/CompactionRecoverTask.java
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/CompactionRecoverTask.java
index 948e80f..63526d1 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/CompactionRecoverTask.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/task/CompactionRecoverTask.java
@@ -66,7 +66,7 @@ public class CompactionRecoverTask implements Callable<Void> {
compactionRecoverCallBack.call();
logger.info(
"recover task finish, current compaction thread is {}",
- CompactionTaskManager.getInstance().getTaskCount());
+ CompactionTaskManager.getInstance().getExecutingTaskCount());
return null;
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionSchedulerTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionSchedulerTest.java
index b871f15..8082e19 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionSchedulerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionSchedulerTest.java
@@ -94,7 +94,7 @@ public class CompactionSchedulerTest {
}
MergeManager.getINSTANCE().start();
CompactionTaskManager.getInstance().start();
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
@@ -216,7 +216,7 @@ public class CompactionSchedulerTest {
e.printStackTrace();
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -338,14 +338,14 @@ public class CompactionSchedulerTest {
e.printStackTrace();
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -468,7 +468,7 @@ public class CompactionSchedulerTest {
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -566,7 +566,7 @@ public class CompactionSchedulerTest {
fullPath, chunkPagePointsNum, 100 * i + 50, tsFileResource);
tsFileManager.add(tsFileResource, false);
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -723,7 +723,7 @@ public class CompactionSchedulerTest {
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -968,7 +968,7 @@ public class CompactionSchedulerTest {
e.printStackTrace();
}
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
@@ -1064,7 +1064,7 @@ public class CompactionSchedulerTest {
tsFileManager.add(tsFileResource, false);
}
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManagerTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManagerTest.java
new file mode 100644
index 0000000..46098c4
--- /dev/null
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/CompactionTaskManagerTest.java
@@ -0,0 +1,252 @@
+/*
+ * 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;
+
+import org.apache.iotdb.db.constant.TestConstant;
+import org.apache.iotdb.db.engine.compaction.inner.InnerCompactionTest;
+import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionTask;
+import org.apache.iotdb.db.engine.compaction.task.AbstractCompactionTask;
+import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
+
+import org.apache.commons.io.FileUtils;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.File;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class CompactionTaskManagerTest extends InnerCompactionTest {
+ static final Logger logger =
LoggerFactory.getLogger(CompactionTaskManagerTest.class);
+ File tempSGDir;
+ final long MAX_WAITING_TIME = 120_000;
+
+ @Before
+ public void setUp() throws Exception {
+ tempSGDir = new File(TestConstant.getTestTsFileDir("root.compactionTest",
0, 0));
+ if (tempSGDir.exists()) {
+ FileUtils.deleteDirectory(tempSGDir);
+ }
+ Assert.assertTrue(tempSGDir.mkdirs());
+ super.setUp();
+ }
+
+ @Test
+ public void testRepeatedSubmitBeforeExecution() throws Exception {
+ logger.warn("testRepeatedSubmitBeforeExecution");
+ TsFileManager tsFileManager =
+ new TsFileManager("root.compactionTest", "0",
tempSGDir.getAbsolutePath());
+ tsFileManager.addAll(seqResources, true);
+ SizeTieredCompactionTask task1 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ SizeTieredCompactionTask task2 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ CompactionTaskManager manager = CompactionTaskManager.getInstance();
+ try {
+ Assert.assertTrue(manager.addTaskToWaitingQueue(task1));
+ Assert.assertEquals(manager.getTotalTaskCount(), 1);
+ // a same task should not be submitted compaction task manager
+ Assert.assertFalse(manager.addTaskToWaitingQueue(task2));
+ Assert.assertEquals(manager.getTotalTaskCount(), 1);
+ manager.submitTaskFromTaskQueue();
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ Thread.sleep(5000);
+ Assert.assertEquals(0, manager.getTotalTaskCount());
+ long waitingTime = 0;
+ while (manager.getRunningCompactionTaskList().size() > 0) {
+ Thread.sleep(100);
+ waitingTime += 100;
+ if (waitingTime % 10000 == 0) {
+ logger.warn("{}", manager.getRunningCompactionTaskList());
+ }
+ if (waitingTime > MAX_WAITING_TIME) {
+ Assert.fail();
+ }
+ }
+ }
+
+ @Test
+ public void testRepeatedSubmitWhenExecuting() throws Exception {
+ logger.warn("testRepeatedSubmitWhenExecuting");
+ TsFileManager tsFileManager =
+ new TsFileManager("root.compactionTest", "0",
tempSGDir.getAbsolutePath());
+ tsFileManager.addAll(seqResources, true);
+ SizeTieredCompactionTask task1 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ SizeTieredCompactionTask task2 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ tsFileManager.writeLock("test");
+ try {
+ CompactionTaskManager manager = CompactionTaskManager.getInstance();
+ manager.addTaskToWaitingQueue(task1);
+ manager.submitTaskFromTaskQueue();
+ Thread.sleep(2000);
+ // When a same compaction task is executing, the compaction task should
not be submitted!
+ Assert.assertEquals(manager.getExecutingTaskCount(), 1);
+ Assert.assertFalse(manager.addTaskToWaitingQueue(task2));
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ long waitingTime = 0;
+ while
(CompactionTaskManager.getInstance().getRunningCompactionTaskList().size() > 0)
{
+ Thread.sleep(100);
+ waitingTime += 100;
+ if (waitingTime % 10000 == 0) {
+ logger.warn("{}",
CompactionTaskManager.getInstance().getRunningCompactionTaskList());
+ }
+ if (waitingTime > MAX_WAITING_TIME) {
+ Assert.fail();
+ }
+ }
+ }
+
+ @Test
+ public void testRepeatedSubmitAfterExecution() throws Exception {
+ logger.warn("testRepeatedSubmitAfterExecution");
+ TsFileManager tsFileManager =
+ new TsFileManager("root.compactionTest", "0",
tempSGDir.getAbsolutePath());
+ tsFileManager.addAll(seqResources, true);
+ SizeTieredCompactionTask task1 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ SizeTieredCompactionTask task2 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ CompactionTaskManager manager = CompactionTaskManager.getInstance();
+ manager.addTaskToWaitingQueue(task1);
+ manager.submitTaskFromTaskQueue();
+ while (manager.getTotalTaskCount() > 0) {
+ Thread.sleep(10);
+ }
+ tsFileManager.writeLock("test");
+ // an invalid task can be submitted to waiting queue, but should not be
submitted to thread pool
+ Assert.assertTrue(manager.addTaskToWaitingQueue(task2));
+ manager.submitTaskFromTaskQueue();
+ Assert.assertEquals(manager.getExecutingTaskCount(), 0);
+ long waitingTime = 0;
+ while (manager.getRunningCompactionTaskList().size() > 0) {
+ Thread.sleep(100);
+ waitingTime += 100;
+ if (waitingTime % 10000 == 0) {
+ logger.warn("{}", manager.getRunningCompactionTaskList());
+ }
+ if (waitingTime > MAX_WAITING_TIME) {
+ Assert.fail();
+ }
+ }
+ }
+
+ @Test
+ public void testRemoveSelfFromRunningList() throws Exception {
+ logger.warn("testRemoveSelfFromRunningList");
+ TsFileManager tsFileManager =
+ new TsFileManager("root.compactionTest", "0",
tempSGDir.getAbsolutePath());
+ tsFileManager.addAll(seqResources, true);
+ SizeTieredCompactionTask task1 =
+ new SizeTieredCompactionTask(
+ "root.compactionTest",
+ "0",
+ 0,
+ tsFileManager,
+ tsFileManager.getSequenceListByTimePartition(0),
+ seqResources,
+ true,
+ new AtomicInteger(0));
+ CompactionTaskManager manager = CompactionTaskManager.getInstance();
+ tsFileManager.writeLock("test");
+ try {
+ manager.addTaskToWaitingQueue(task1);
+ manager.submitTaskFromTaskQueue();
+ Thread.sleep(5000);
+ List<AbstractCompactionTask> runningList =
manager.getRunningCompactionTaskList();
+ // compaction task should add itself to running list
+ Assert.assertEquals(1, runningList.size());
+ Assert.assertTrue(runningList.contains(task1));
+ } finally {
+ tsFileManager.writeUnlock();
+ }
+ // after execution, task should remove itself from running list
+ Thread.sleep(5000);
+ List<AbstractCompactionTask> runningList =
manager.getRunningCompactionTaskList();
+ Assert.assertEquals(0, runningList.size());
+ long waitingTime = 0;
+ while (manager.getRunningCompactionTaskList().size() > 0) {
+ Thread.sleep(100);
+ waitingTime += 100;
+ if (waitingTime % 10000 == 0) {
+ logger.warn("{}", manager.getRunningCompactionTaskList());
+ }
+ if (waitingTime > MAX_WAITING_TIME) {
+ Assert.fail();
+ }
+ }
+ }
+}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionChunkTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionChunkTest.java
index cad4215..3f29a44 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionChunkTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionChunkTest.java
@@ -26,8 +26,6 @@ import
org.apache.iotdb.db.engine.compaction.inner.utils.InnerSpaceCompactionUti
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
import org.apache.iotdb.db.exception.StorageEngineException;
import org.apache.iotdb.db.exception.metadata.IllegalPathException;
-import org.apache.iotdb.db.exception.metadata.MetadataException;
-import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
import org.apache.iotdb.tsfile.read.common.BatchData;
@@ -62,7 +60,7 @@ public class InnerCompactionChunkTest extends
InnerCompactionTest {
File tempSGDir;
@Before
- public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ public void setUp() throws Exception {
tempSGDir = new File(TestConstant.getTestTsFileDir("root.compactionTest",
0, 0));
if (tempSGDir.exists()) {
FileUtils.deleteDirectory(tempSGDir);
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionLogTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionLogTest.java
index 1f5bc24..35229e2 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionLogTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionLogTest.java
@@ -23,8 +23,6 @@ import org.apache.iotdb.db.constant.TestConstant;
import org.apache.iotdb.db.engine.compaction.CompactionScheduler;
import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
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.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
import org.apache.commons.io.FileUtils;
@@ -45,7 +43,7 @@ public class InnerCompactionLogTest extends
InnerCompactionTest {
@Override
@Before
- public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ public void setUp() throws Exception {
tempSGDir = new File(TestConstant.getTestTsFileDir("root.compactionTest",
0, 0));
if (!tempSGDir.exists()) {
Assert.assertTrue(tempSGDir.mkdirs());
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionMoreDataTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionMoreDataTest.java
index 572e7de..256327a 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionMoreDataTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionMoreDataTest.java
@@ -182,7 +182,7 @@ public class InnerCompactionMoreDataTest extends
InnerCompactionTest {
}
@Before
- public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ public void setUp() throws Exception {
tempSGDir = new File(TestConstant.getTestTsFileDir("root.compactionTest",
0, 0));
if (!tempSGDir.exists()) {
Assert.assertTrue(tempSGDir.mkdirs());
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionSchedulerTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionSchedulerTest.java
index f935720..e1544fa 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionSchedulerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionSchedulerTest.java
@@ -82,7 +82,7 @@ public class InnerCompactionSchedulerTest {
CompactionTaskManager.getInstance().submitTaskFromTaskQueue();
try {
- Thread.sleep(1000);
+ Thread.sleep(5000);
} catch (Exception e) {
}
@@ -113,7 +113,7 @@ public class InnerCompactionSchedulerTest {
CompactionTaskManager.getInstance().submitTaskFromTaskQueue();
long waitingTime = 0;
- while (CompactionTaskManager.getInstance().getTaskCount() != 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() != 0) {
try {
Thread.sleep(100);
waitingTime += 100;
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java
index c7aca33..df023c1 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerCompactionTest.java
@@ -76,7 +76,7 @@ public abstract class InnerCompactionTest {
private int prevMergeChunkThreshold;
@Before
- public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ public void setUp() throws IOException, WriteProcessException,
MetadataException, Exception {
EnvironmentUtils.envSetUp();
IoTDB.metaManager.init();
prevMergeChunkThreshold =
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionUtilsTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionUtilsTest.java
index eb5b69c..4191a40 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionUtilsTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/inner/InnerSpaceCompactionUtilsTest.java
@@ -26,8 +26,6 @@ import
org.apache.iotdb.db.engine.compaction.inner.utils.SizeTieredCompactionLog
import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
import org.apache.iotdb.db.exception.StorageEngineException;
import org.apache.iotdb.db.exception.metadata.IllegalPathException;
-import org.apache.iotdb.db.exception.metadata.MetadataException;
-import org.apache.iotdb.tsfile.exception.write.WriteProcessException;
import org.apache.iotdb.tsfile.read.TsFileReader;
import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
import org.apache.iotdb.tsfile.read.common.Path;
@@ -54,7 +52,7 @@ public class InnerSpaceCompactionUtilsTest extends
InnerCompactionTest {
@Override
@Before
- public void setUp() throws IOException, WriteProcessException,
MetadataException {
+ public void setUp() throws Exception {
tempSGDir = new File(TestConstant.getTestTsFileDir("root.compactionTest",
0, 0));
if (!tempSGDir.exists()) {
assertTrue(tempSGDir.mkdirs());
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/task/FakedInnerSpaceCompactionTask.java
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/task/FakedInnerSpaceCompactionTask.java
index 50af7a5..b43fcc1 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/compaction/task/FakedInnerSpaceCompactionTask.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/compaction/task/FakedInnerSpaceCompactionTask.java
@@ -18,6 +18,7 @@
*/
package org.apache.iotdb.db.engine.compaction.task;
+import org.apache.iotdb.db.engine.compaction.CompactionTaskManager;
import
org.apache.iotdb.db.engine.compaction.inner.sizetiered.SizeTieredCompactionTask;
import org.apache.iotdb.db.engine.storagegroup.FakedTsFileResource;
import org.apache.iotdb.db.engine.storagegroup.TsFileManager;
@@ -53,28 +54,37 @@ public class FakedInnerSpaceCompactionTask extends
SizeTieredCompactionTask {
@Override
protected void doCompaction() throws IOException {
- TsFileNameGenerator.TsFileName name =
-
TsFileNameGenerator.getTsFileName(selectedTsFileResourceList.get(0).getTsFile().getName());
- String newName =
- TsFileNameGenerator.generateNewTsFileName(
- name.getTime(),
- name.getVersion(),
- name.getInnerCompactionCnt() + 1,
- name.getCrossCompactionCnt());
- FakedTsFileResource targetTsFileResource = new FakedTsFileResource(0,
newName);
- long targetFileSize = 0;
- for (TsFileResource resource : selectedTsFileResourceList) {
- targetFileSize += resource.getTsFileSize();
- }
- targetTsFileResource.setTsFileSize(targetFileSize);
- this.tsFileResourceList.insertBefore(selectedTsFileResourceList.get(0),
targetTsFileResource);
- for (TsFileResource tsFileResource : selectedTsFileResourceList) {
- this.tsFileResourceList.remove(tsFileResource);
+ try {
+ TsFileNameGenerator.TsFileName name =
+ TsFileNameGenerator.getTsFileName(
+ selectedTsFileResourceList.get(0).getTsFile().getName());
+ String newName =
+ TsFileNameGenerator.generateNewTsFileName(
+ name.getTime(),
+ name.getVersion(),
+ name.getInnerCompactionCnt() + 1,
+ name.getCrossCompactionCnt());
+ FakedTsFileResource targetTsFileResource = new FakedTsFileResource(0,
newName);
+ long targetFileSize = 0;
+ for (TsFileResource resource : selectedTsFileResourceList) {
+ targetFileSize += resource.getTsFileSize();
+ }
+ targetTsFileResource.setTsFileSize(targetFileSize);
+ this.tsFileResourceList.insertBefore(selectedTsFileResourceList.get(0),
targetTsFileResource);
+ for (TsFileResource tsFileResource : selectedTsFileResourceList) {
+ this.tsFileResourceList.remove(tsFileResource);
+ }
+ } finally {
+ CompactionTaskManager.getInstance().removeRunningTaskFromList(this);
}
}
@Override
public boolean equalsOtherTask(AbstractCompactionTask otherTask) {
+ if (otherTask instanceof FakedInnerSpaceCompactionTask) {
+ FakedInnerSpaceCompactionTask fakedOtherTask =
(FakedInnerSpaceCompactionTask) otherTask;
+ return
this.selectedTsFileResourceList.equals(fakedOtherTask.selectedTsFileResourceList);
+ }
return false;
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/FakedTsFileResource.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/FakedTsFileResource.java
index 7e3dede..d1d8ee6 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/FakedTsFileResource.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/FakedTsFileResource.java
@@ -73,4 +73,15 @@ public class FakedTsFileResource extends TsFileResource {
public File getTsFile() {
return new File(fakeTsfileName);
}
+
+ @Override
+ public boolean equals(Object otherObject) {
+ if (otherObject instanceof FakedTsFileResource) {
+ FakedTsFileResource otherResource = (FakedTsFileResource) otherObject;
+ return this.fakeTsfileName.equals(otherResource.fakeTsfileName)
+ && this.tsFileSize == otherResource.tsFileSize;
+ }
+
+ return false;
+ }
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
index f5fec46..40e96fc 100644
---
a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessorTest.java
@@ -620,7 +620,7 @@ public class StorageGroupProcessorTest {
processor.syncCloseAllWorkingTsFileProcessors();
processor.merge(IoTDBDescriptor.getInstance().getConfig().isForceFullMerge());
long totalWaitingTime = 0;
- while (CompactionTaskManager.getInstance().getTaskCount() > 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() > 0) {
// wait
try {
Thread.sleep(100);
diff --git
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBNewTsFileCompactionIT.java
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBNewTsFileCompactionIT.java
index 46f5309..cdf6c27 100644
---
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBNewTsFileCompactionIT.java
+++
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBNewTsFileCompactionIT.java
@@ -1052,7 +1052,7 @@ public class IoTDBNewTsFileCompactionIT {
long startTime = System.nanoTime();
// get the size of level 1's tsfile list to judge whether merge is finished
- while (CompactionTaskManager.getInstance().getTaskCount() != 0) {
+ while (CompactionTaskManager.getInstance().getExecutingTaskCount() != 0) {
TimeUnit.MILLISECONDS.sleep(100);
// wait too long, just break
if ((System.nanoTime() - startTime) >= MAX_WAIT_TIME_FOR_MERGE) {