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

Reply via email to