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) {

Reply via email to