This is an automated email from the ASF dual-hosted git repository. tanxinyu pushed a commit to branch master_performance in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 0224027c50ecca51ae3f40f2f9abc9710b1fe7ac 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) {
