This is an automated email from the ASF dual-hosted git repository.

qiaojialin pushed a commit to branch autoai
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/autoai by this push:
     new 0ada293  replace synchronized with write lock in cross compaction 
selection (#3633)
0ada293 is described below

commit 0ada293dff1c4a94374b7a9251ba4e8dc2712900
Author: liuxuxin <[email protected]>
AuthorDate: Tue Jul 27 19:54:23 2021 +0800

    replace synchronized with write lock in cross compaction selection (#3633)
---
 .../db/engine/compaction/TsFileManagement.java     | 148 +++++++++++----------
 .../db/engine/merge/manage/MergeResource.java      |   3 +-
 .../merge/selector/MaxFileMergeFileSelector.java   |  15 ++-
 .../db/engine/storagegroup/TsFileResource.java     |   2 +-
 4 files changed, 87 insertions(+), 81 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
index d997e04..dde54a7 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/compaction/TsFileManagement.java
@@ -193,7 +193,7 @@ public abstract class TsFileManagement {
     }
   }
 
-  public synchronized boolean merge(
+  public boolean merge(
       boolean fullMerge,
       List<TsFileResource> seqMergeList,
       List<TsFileResource> unSeqMergeList,
@@ -209,89 +209,93 @@ public abstract class TsFileManagement {
       }
     }
     isUnseqMerging = true;
-
-    if (seqMergeList.isEmpty()) {
-      logger.info("{} no seq files to be merged", storageGroupName);
-      isUnseqMerging = false;
-      return false;
-    }
-
-    if (unSeqMergeList.isEmpty()) {
-      logger.info("{} no unseq files to be merged", storageGroupName);
-      isUnseqMerging = false;
-      return false;
-    }
-
-    if (unSeqMergeList.size() > maxOpenFileNumInEachUnseqCompaction) {
-      logger.info(
-          "{} too much unseq files to be merged, reduce it to {}",
-          storageGroupName,
-          maxOpenFileNumInEachUnseqCompaction);
-      unSeqMergeList = unSeqMergeList.subList(0, 
maxOpenFileNumInEachUnseqCompaction);
-    }
-
-    long budget = 
IoTDBDescriptor.getInstance().getConfig().getMergeMemoryBudget();
-    long timeLowerBound = System.currentTimeMillis() - dataTTL;
-    MergeResource mergeResource = new MergeResource(seqMergeList, 
unSeqMergeList, timeLowerBound);
-
-    IMergeFileSelector fileSelector = getMergeFileSelector(budget, 
mergeResource);
+    writeLock();
     try {
-      List[] mergeFiles = fileSelector.select();
-      if (mergeFiles.length == 0) {
-        logger.info(
-            "{} cannot select merge candidates under the budget {}", 
storageGroupName, budget);
+      if (seqMergeList.isEmpty()) {
+        logger.info("{} no seq files to be merged", storageGroupName);
         isUnseqMerging = false;
         return false;
       }
-      // avoid pending tasks holds the metadata and streams
-      mergeResource.clear();
-      String taskName = storageGroupName + "-" + System.currentTimeMillis();
-      // do not cache metadata until true candidates are chosen, or too much 
metadata will be
-      // cached during selection
-      mergeResource.setCacheDeviceMeta(true);
-
-      for (TsFileResource tsFileResource : mergeResource.getSeqFiles()) {
-        tsFileResource.setMerging(true);
-      }
-      for (TsFileResource tsFileResource : mergeResource.getUnseqFiles()) {
-        tsFileResource.setMerging(true);
+
+      if (unSeqMergeList.isEmpty()) {
+        logger.info("{} no unseq files to be merged", storageGroupName);
+        isUnseqMerging = false;
+        return false;
       }
 
-      mergeStartTime = System.currentTimeMillis();
-      MergeTask mergeTask =
-          new MergeTask(
-              mergeResource,
-              storageGroupDir,
-              this::mergeEndAction,
-              taskName,
-              fullMerge,
-              fileSelector.getConcurrentMergeNum(),
-              storageGroupName);
-      mergingModification =
-          new ModificationFile(storageGroupDir + File.separator + 
MERGING_MODIFICATION_FILE_NAME);
-      MergeManager.getINSTANCE().submitMainTask(mergeTask);
-      if (logger.isInfoEnabled()) {
+      if (unSeqMergeList.size() > maxOpenFileNumInEachUnseqCompaction) {
         logger.info(
-            "{} submits a merge task {}, merging {} seqFiles, {} unseqFiles",
+            "{} too much unseq files to be merged, reduce it to {}",
             storageGroupName,
-            taskName,
-            mergeFiles[0].size(),
-            mergeFiles[1].size());
+            maxOpenFileNumInEachUnseqCompaction);
+        unSeqMergeList = unSeqMergeList.subList(0, 
maxOpenFileNumInEachUnseqCompaction);
       }
-      // wait until unseq merge has finished
-      while (isUnseqMerging) {
-        try {
-          Thread.sleep(200);
-        } catch (InterruptedException e) {
-          logger.error("{} [Compaction] shutdown", storageGroupName, e);
-          Thread.currentThread().interrupt();
+
+      long budget = 
IoTDBDescriptor.getInstance().getConfig().getMergeMemoryBudget();
+      long timeLowerBound = System.currentTimeMillis() - dataTTL;
+      MergeResource mergeResource = new MergeResource(seqMergeList, 
unSeqMergeList, timeLowerBound);
+
+      IMergeFileSelector fileSelector = getMergeFileSelector(budget, 
mergeResource);
+      try {
+        List[] mergeFiles = fileSelector.select();
+        if (mergeFiles.length == 0) {
+          logger.info(
+              "{} cannot select merge candidates under the budget {}", 
storageGroupName, budget);
+          isUnseqMerging = false;
           return false;
         }
+        // avoid pending tasks holds the metadata and streams
+        mergeResource.clear();
+        String taskName = storageGroupName + "-" + System.currentTimeMillis();
+        // do not cache metadata until true candidates are chosen, or too much 
metadata will be
+        // cached during selection
+        mergeResource.setCacheDeviceMeta(true);
+
+        for (TsFileResource tsFileResource : mergeResource.getSeqFiles()) {
+          tsFileResource.setMerging(true);
+        }
+        for (TsFileResource tsFileResource : mergeResource.getUnseqFiles()) {
+          tsFileResource.setMerging(true);
+        }
+
+        mergeStartTime = System.currentTimeMillis();
+        MergeTask mergeTask =
+            new MergeTask(
+                mergeResource,
+                storageGroupDir,
+                this::mergeEndAction,
+                taskName,
+                fullMerge,
+                fileSelector.getConcurrentMergeNum(),
+                storageGroupName);
+        mergingModification =
+            new ModificationFile(storageGroupDir + File.separator + 
MERGING_MODIFICATION_FILE_NAME);
+        MergeManager.getINSTANCE().submitMainTask(mergeTask);
+        if (logger.isInfoEnabled()) {
+          logger.info(
+              "{} submits a merge task {}, merging {} seqFiles, {} unseqFiles",
+              storageGroupName,
+              taskName,
+              mergeFiles[0].size(),
+              mergeFiles[1].size());
+        }
+        // wait until unseq merge has finished
+        while (isUnseqMerging) {
+          try {
+            Thread.sleep(200);
+          } catch (InterruptedException e) {
+            logger.error("{} [Compaction] shutdown", storageGroupName, e);
+            Thread.currentThread().interrupt();
+            return false;
+          }
+        }
+        return true;
+      } catch (MergeException | IOException e) {
+        logger.error("{} cannot select file for merge", storageGroupName, e);
+        return false;
       }
-      return true;
-    } catch (MergeException | IOException e) {
-      logger.error("{} cannot select file for merge", storageGroupName, e);
-      return false;
+    } finally {
+      writeUnlock();
     }
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
index df1add1..bb00989 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/manage/MergeResource.java
@@ -81,7 +81,8 @@ public class MergeResource {
   private boolean filterResource(TsFileResource res) {
     return res.getTsFile().exists()
         && !res.isDeleted()
-        && (!res.isClosed() || res.stillLives(timeLowerBound));
+        && (!res.isClosed() || res.stillLives(timeLowerBound))
+        && !res.isMerging();
   }
 
   public MergeResource(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
index 3deb763..d55eaab 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/merge/selector/MaxFileMergeFileSelector.java
@@ -155,7 +155,7 @@ public class MaxFileMergeFileSelector implements 
IMergeFileSelector {
       if (seqSelectedNum != resource.getSeqFiles().size()) {
         selectOverlappedSeqFiles(unseqFile);
       }
-      boolean isClosed = checkClosed(unseqFile);
+      boolean isClosed = checkChoosable(unseqFile);
       if (!isClosed) {
         tmpSelectedSeqFiles.clear();
         unseqIndex++;
@@ -206,18 +206,19 @@ public class MaxFileMergeFileSelector implements 
IMergeFileSelector {
     return false;
   }
 
-  private boolean checkClosed(TsFileResource unseqFile) {
-    boolean isClosed = unseqFile.isClosed();
-    if (!isClosed) {
+  private boolean checkChoosable(TsFileResource unseqFile) {
+    boolean choosable = unseqFile.isClosed() && !unseqFile.isMerging();
+    if (!choosable) {
       return false;
     }
     for (Integer seqIdx : tmpSelectedSeqFiles) {
-      if (!resource.getSeqFiles().get(seqIdx).isClosed()) {
-        isClosed = false;
+      if (!resource.getSeqFiles().get(seqIdx).isClosed()
+          || resource.getSeqFiles().get(seqIdx).isMerging()) {
+        choosable = false;
         break;
       }
     }
-    return isClosed;
+    return choosable;
   }
 
   private void selectOverlappedSeqFiles(TsFileResource unseqFile) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
index 16cdb54..6ba652c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileResource.java
@@ -522,7 +522,7 @@ public class TsFileResource {
     this.deleted = deleted;
   }
 
-  boolean isMerging() {
+  public boolean isMerging() {
     return isMerging;
   }
 

Reply via email to