This is an automated email from the ASF dual-hosted git repository. qiaojialin pushed a commit to branch rel/0.12 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 9be4d2fdbc78fad1d82c93d6522443ba5ad9dbda 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; }
