This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch rel/1.2
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/1.2 by this push:
new f174e9ff255 [To rel/1.2] Refactoring DeleteOutdatedFileTask in WalNode
(#10992)
f174e9ff255 is described below
commit f174e9ff255915f2dbfe60867676beb8b62a644b
Author: Zhijia Cao <[email protected]>
AuthorDate: Thu Aug 31 10:59:49 2023 +0800
[To rel/1.2] Refactoring DeleteOutdatedFileTask in WalNode (#10992)
---
.../wal/checkpoint/CheckpointManager.java | 2 +-
.../storageengine/dataregion/wal/node/WALNode.java | 251 ++++++++++++++-------
2 files changed, 171 insertions(+), 82 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/checkpoint/CheckpointManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/checkpoint/CheckpointManager.java
index 367f4be6328..3899b07a3f1 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/checkpoint/CheckpointManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/checkpoint/CheckpointManager.java
@@ -88,7 +88,7 @@ public class CheckpointManager implements AutoCloseable {
logHeader();
}
- private List<MemTableInfo> snapshotMemTableInfos() {
+ public List<MemTableInfo> snapshotMemTableInfos() {
infoLock.lock();
try {
return new ArrayList<>(memTableId2Info.values());
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
index 8ba4b3149a8..8ee08e09156 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNode.java
@@ -55,6 +55,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.wal.utils.listener.WALFlushL
import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
import org.apache.iotdb.tsfile.utils.TsFileUtils;
+import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -184,6 +185,7 @@ public class WALNode implements IWALNode {
}
// region methods for pipe
+
/**
* Pin the wal files of the given memory table. Notice: cannot pin one
memTable too long,
* otherwise the wal disk usage may too large.
@@ -205,6 +207,7 @@ public class WALNode implements IWALNode {
// endregion
// region Task to delete outdated .wal files
+
/** Delete outdated .wal files. */
public void deleteOutdatedFiles() {
try {
@@ -219,34 +222,79 @@ public class WALNode implements IWALNode {
// .wal files whose version ids are less than first valid version id
should be deleted
private long firstValidVersionId;
// the effective information ratio
- private double effectiveInfoRatio;
- // recursion time of calling deletion
+ private double effectiveInfoRatio = 0d;
+
+ private List<Long> pinnedMemTableIds;
+
+ private File[] filesShouldDelete;
+
+ private int fileIndexAfterFilterSafelyDeleteIndex = Integer.MAX_VALUE;
+ private List<Long> successfullyDeleted;
+ private long deleteFileSize;
+
private int recursionTime = 0;
+ public DeleteOutdatedFileTask() {}
+
+ private void init() {
+ this.firstValidVersionId = initFirstValidWALVersionId();
+ this.filesShouldDelete =
logDirectory.listFiles(this::filterFilesToDelete);
+ if (filesShouldDelete == null) {
+ filesShouldDelete = new File[0];
+ }
+ this.pinnedMemTableIds = initPinnedMemTableIds();
+ WALFileUtils.ascSortByVersionId(filesShouldDelete);
+ this.fileIndexAfterFilterSafelyDeleteIndex =
initFileIndexAfterFilterSafelyDeleteIndex();
+ this.successfullyDeleted = new ArrayList<>();
+ this.deleteFileSize = 0;
+ }
+
+ private List<Long> initPinnedMemTableIds() {
+ List<MemTableInfo> memTableInfos =
checkpointManager.snapshotMemTableInfos();
+ if (memTableInfos.isEmpty()) {
+ return new ArrayList<>();
+ }
+ List<Long> pinnedIds = new ArrayList<>();
+ for (MemTableInfo memTableInfo : memTableInfos) {
+ if (memTableInfo.isFlushed() && memTableInfo.isPinned()) {
+ pinnedIds.add(memTableInfo.getMemTableId());
+ }
+ }
+ return pinnedIds;
+ }
+
@Override
public void run() {
- // init firstValidVersionId
- firstValidVersionId = checkpointManager.getFirstValidWALVersionId();
- if (firstValidVersionId == Long.MIN_VALUE) {
- // roll wal log writer to delete current wal file
- if (buffer.getCurrentWALFileSize() > 0) {
- rollWALFile();
- }
- // update firstValidVersionId
- firstValidVersionId = checkpointManager.getFirstValidWALVersionId();
- if (firstValidVersionId == Long.MIN_VALUE) {
- firstValidVersionId = buffer.getCurrentWALFileVersion();
+ // The intent of the loop execution here is to try to get as many
memTable flush or snapshot
+ // as possible when the valid information ratio is less than the
configured value.
+ // In addition, if the disk space used by wal exceeds the limit
threshold, resulting in a
+ // write rejection, the task will continue to attempt to delete expired
files until the
+ // threshold is no longer exceeded
+ while (recursionTime < MAX_RECURSION_TIME ||
WALManager.getInstance().shouldThrottle()) {
+ // init delete outdated file task fields
+ init();
+
+ // delete outdated WAL files and record which delete successfully and
which delete failed.
+ deleteOutdatedFilesAndUpdateMetric();
+
+ // summary the execution result and output a log
+ summarizeExecuteResult();
+
+ // update current effective info ration
+ updateEffectiveInfoRationAndUpdateMetric();
+
+ // decide whether to snapshot or flush based on the effective info
ration and throttle
+ // threshold
+ if (trySnapshotOrFlushMemTable()
+ && safelyDeletedSearchIndex != DEFAULT_SAFELY_DELETED_SEARCH_INDEX
+ && !WALManager.getInstance().shouldThrottle()) {
+ return;
}
+ recursionTime++;
}
+ }
- // delete outdated files
- logger.debug(
- "Start deleting outdated wal files for wal node-{}, the first valid
version id is {}, and the safely deleted search index is {}.",
- identifier,
- firstValidVersionId,
- safelyDeletedSearchIndex);
- deleteOutdatedFiles();
-
+ private void updateEffectiveInfoRationAndUpdateMetric() {
// calculate effective information ratio
long costOfActiveMemTables =
checkpointManager.getTotalCostOfActiveMemTables();
long costOfFlushedMemTables = totalCostOfFlushedMemTables.get();
@@ -255,91 +303,110 @@ public class WALNode implements IWALNode {
return;
}
effectiveInfoRatio = (double) costOfActiveMemTables / totalCost;
- WRITING_METRICS.recordWALNodeEffectiveInfoRatio(identifier,
effectiveInfoRatio);
logger.debug(
"Effective information ratio is {}, active memTables cost is {},
flushed memTables cost is {}",
effectiveInfoRatio,
costOfActiveMemTables,
costOfFlushedMemTables);
+ WRITING_METRICS.recordWALNodeEffectiveInfoRatio(identifier,
effectiveInfoRatio);
+ }
- // try updating first valid version id by snapshotting or flushing
memTable,
- // then delete old .wal files again
- if (!shouldSnapshotOrFlush()) {
+ private void summarizeExecuteResult() {
+ if (filesShouldDelete.length == 0) {
+ if (logger.isDebugEnabled()) {
+ logger.debug(
+ "wal node-{}:no wal file was found that should be deleted,
current first valid version id is {}",
+ identifier,
+ firstValidVersionId);
+ }
return;
}
- logger.debug(
- "Effective information ratio {} (active memTables cost is {},
flushed memTables cost is {}) of wal node-{} is below wal min effective info
ratio {}, some memTables will be snapshot or flushed.",
- effectiveInfoRatio,
- costOfActiveMemTables,
- costOfFlushedMemTables,
- identifier,
- config.getWalMinEffectiveInfoRatio());
- boolean isSuccess = snapshotOrFlushMemTable();
- if (isSuccess && recursionTime < MAX_RECURSION_TIME) {
- // wal is used to search, cannot optimize files deletion
- if (safelyDeletedSearchIndex != DEFAULT_SAFELY_DELETED_SEARCH_INDEX) {
- return;
+
+ if (!pinnedMemTableIds.isEmpty()
+ || fileIndexAfterFilterSafelyDeleteIndex < filesShouldDelete.length)
{
+ if (logger.isDebugEnabled()) {
+ StringBuilder summary =
+ new StringBuilder(
+ String.format(
+ "wal node-%s delete outdated files summary:the range
that should be removed is: [%d,%d], delete successful is [%s], end file index
is: [%s].The following reasons influenced the result: %s",
+ identifier,
+
WALFileUtils.parseVersionId(filesShouldDelete[0].getName()),
+ WALFileUtils.parseVersionId(
+ filesShouldDelete[filesShouldDelete.length -
1].getName()),
+ StringUtils.join(successfullyDeleted, ","),
+ fileIndexAfterFilterSafelyDeleteIndex,
+ System.getProperty("line.separator")));
+
+ if (!pinnedMemTableIds.isEmpty()) {
+ summary
+ .append("- MemTable has been flushed but pinned by PIPE, the
MemTableId list is : ")
+ .append(StringUtils.join(pinnedMemTableIds, ","))
+ .append(".")
+ .append(System.getProperty("line.separator"));
+ }
+ if (fileIndexAfterFilterSafelyDeleteIndex <
filesShouldDelete.length) {
+ summary.append(
+ String.format(
+ "- The data in the wal file was not consumed by the
consensus group,current search index is %d, safely delete index is %d",
+ getCurrentSearchIndex(), safelyDeletedSearchIndex));
+ }
+ String summaryLog = summary.toString();
+ logger.debug(summaryLog);
}
- recursionTime++;
- run();
+
+ } else {
+ logger.debug(
+ "Successfully delete {} outdated wal files for wal node-{},first
valid version id is {}",
+ successfullyDeleted.size(),
+ identifier,
+ firstValidVersionId);
}
}
- /** Return true iff cannot delete all outdated files because of
IoTConsensus. */
- private boolean deleteOutdatedFiles() {
- // find all files to delete
- // delete files whose version < firstValidVersionId
- File[] filesToDelete = logDirectory.listFiles(this::filterFilesToDelete);
- if (filesToDelete == null || filesToDelete.length == 0) {
- return false;
+ /** Delete obsolete wal files while recording which succeeded or failed */
+ private void deleteOutdatedFilesAndUpdateMetric() {
+ if (filesShouldDelete.length == 0) {
+ return;
}
+ for (int i = 0; i < fileIndexAfterFilterSafelyDeleteIndex; ++i) {
+ long fileSize = filesShouldDelete[i].length();
+ long versionId =
WALFileUtils.parseVersionId(filesShouldDelete[i].getName());
+ if (filesShouldDelete[i].delete()) {
+ deleteFileSize += fileSize;
+ Long memTableRamCostSum =
walFileVersionId2MemTablesTotalCost.remove(versionId);
+ if (memTableRamCostSum != null) {
+ totalCostOfFlushedMemTables.addAndGet(-memTableRamCostSum);
+ }
+ successfullyDeleted.add(versionId);
+ } else {
+ logger.info(
+ "Fail to delete outdated wal file {} of wal node-{}.",
+ filesShouldDelete[i],
+ identifier);
+ }
+ }
+ buffer.subtractDiskUsage(deleteFileSize);
+ buffer.subtractFileNum(successfullyDeleted.size());
+ }
- // delete files whose content's search index are all <=
safelyDeletedSearchIndex
- WALFileUtils.ascSortByVersionId(filesToDelete);
- // judge DEFAULT_SAFELY_DELETED_SEARCH_INDEX for standalone,
Long.MIN_VALUE for iot
+ private int initFileIndexAfterFilterSafelyDeleteIndex() {
int endFileIndex =
safelyDeletedSearchIndex == DEFAULT_SAFELY_DELETED_SEARCH_INDEX
- ? filesToDelete.length
+ ? filesShouldDelete.length
: WALFileUtils.binarySearchFileBySearchIndex(
- filesToDelete, safelyDeletedSearchIndex + 1);
+ filesShouldDelete, safelyDeletedSearchIndex + 1);
// delete files whose file status is CONTAINS_NONE_SEARCH_INDEX
if (endFileIndex == -1) {
endFileIndex = 0;
}
- while (endFileIndex < filesToDelete.length) {
- if (WALFileUtils.parseStatusCode(filesToDelete[endFileIndex].getName())
+ while (endFileIndex < filesShouldDelete.length) {
+ if
(WALFileUtils.parseStatusCode(filesShouldDelete[endFileIndex].getName())
== WALFileStatus.CONTAINS_SEARCH_INDEX) {
break;
}
endFileIndex++;
}
-
- // delete files
- int deletedFilesNum = 0;
- long deletedFilesSize = 0;
- for (int i = 0; i < endFileIndex; ++i) {
- long fileSize = filesToDelete[i].length();
- if (filesToDelete[i].delete()) {
- deletedFilesNum++;
- deletedFilesSize += fileSize;
- } else {
- logger.info(
- "Fail to delete outdated wal file {} of wal node-{}.",
filesToDelete[i], identifier);
- }
- // update totalRamCostOfFlushedMemTables
- long versionId =
WALFileUtils.parseVersionId(filesToDelete[i].getName());
- Long memTableRamCostSum =
walFileVersionId2MemTablesTotalCost.remove(versionId);
- if (memTableRamCostSum != null) {
- totalCostOfFlushedMemTables.addAndGet(-memTableRamCostSum);
- }
- }
- buffer.subtractDiskUsage(deletedFilesSize);
- buffer.subtractFileNum(deletedFilesNum);
- logger.debug(
- "Successfully delete {} outdated wal files for wal node-{}.",
- deletedFilesNum,
- identifier);
- return endFileIndex < filesToDelete.length;
+ return endFileIndex;
}
private boolean filterFilesToDelete(File dir, String name) {
@@ -364,7 +431,10 @@ public class WALNode implements IWALNode {
*
* @return true if snapshot or flush is executed successfully
*/
- private boolean snapshotOrFlushMemTable() {
+ private boolean trySnapshotOrFlushMemTable() {
+ if (!shouldSnapshotOrFlush()) {
+ return false;
+ }
// find oldest memTable
MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestMemTableInfo();
if (oldestMemTableInfo == null) {
@@ -498,7 +568,26 @@ public class WALNode implements IWALNode {
dataRegion.writeUnlock();
}
}
+
+ public long initFirstValidWALVersionId() {
+ long firstVersionId = checkpointManager.getFirstValidWALVersionId();
+ // This means that the relevant memTable in the file has been
successfully flushed, so we
+ // should scroll through a new wal file so that the current file can be
deleted
+ if (firstVersionId == Long.MIN_VALUE) {
+ // roll wal log writer to delete current wal file
+ if (buffer.getCurrentWALFileSize() > 0) {
+ rollWALFile();
+ }
+ // update firstValidVersionId
+ firstVersionId = checkpointManager.getFirstValidWALVersionId();
+ if (firstVersionId == Long.MIN_VALUE) {
+ firstVersionId = buffer.getCurrentWALFileVersion();
+ }
+ }
+ return firstVersionId;
+ }
}
+
// endregion
// region Search interfaces for consensus group