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;
}