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

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


The following commit(s) were added to refs/heads/master by this push:
     new 463e23271c2 Revert "Optimized wal file deletion algorithm (#11682)" 
(#11852)
463e23271c2 is described below

commit 463e23271c25bd8179f54225b51f8e3f91fdb3b7
Author: Haonan <[email protected]>
AuthorDate: Fri Jan 5 09:58:25 2024 +0800

    Revert "Optimized wal file deletion algorithm (#11682)" (#11852)
    
    This reverts commit 9fdeb95b6f6d7bc456213a690ec536cf4d0959a8.
---
 .../dataregion/wal/buffer/AbstractWALBuffer.java   |   1 -
 .../dataregion/wal/buffer/WALBuffer.java           |  24 -
 .../wal/checkpoint/CheckpointManager.java          |  28 +-
 .../storageengine/dataregion/wal/node/WALNode.java | 182 +++----
 .../dataregion/wal/node/WALEntryHandlerTest.java   |  13 +-
 .../wal/node/WalDeleteOutdatedNewTest.java         | 585 ---------------------
 6 files changed, 105 insertions(+), 728 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java
index a1c3dcd6672..08bca5f6649 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java
@@ -96,7 +96,6 @@ public abstract class AbstractWALBuffer implements IWALBuffer 
{
    * @throws IOException If failing to close or open the log writer
    */
   protected File rollLogWriter(long searchIndex, WALFileStatus fileStatus) 
throws IOException {
-
     // close file
     currentWALFileWriter.close();
     addDiskUsage(currentWALFileWriter.size());
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
index a93ac366748..cdae8baa1c7 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java
@@ -33,7 +33,6 @@ import 
org.apache.iotdb.db.storageengine.dataregion.wal.checkpoint.CheckpointMan
 import 
org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALNodeClosedException;
 import org.apache.iotdb.db.storageengine.dataregion.wal.io.WALMetaData;
 import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileStatus;
-import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileUtils;
 import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALMode;
 import 
org.apache.iotdb.db.storageengine.dataregion.wal.utils.listener.WALFlushListener;
 import org.apache.iotdb.db.utils.MmapUtil;
@@ -46,13 +45,9 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.nio.MappedByteBuffer;
 import java.util.ArrayList;
-import java.util.HashSet;
 import java.util.List;
-import java.util.Map;
-import java.util.Set;
 import java.util.concurrent.ArrayBlockingQueue;
 import java.util.concurrent.BlockingQueue;
-import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.locks.Condition;
@@ -107,9 +102,6 @@ public class WALBuffer extends AbstractWALBuffer {
   // single thread to sync syncingBuffer to disk
   private final ExecutorService syncBufferThread;
 
-  // manage wal files which have MemTableIds
-  private final Map<Long, Set<Long>> memTableIdsOfWal = new 
ConcurrentHashMap<>();
-
   public WALBuffer(String identifier, String logDirectory) throws 
FileNotFoundException {
     this(identifier, logDirectory, new CheckpointManager(identifier, 
logDirectory), 0, 0L);
   }
@@ -194,8 +186,6 @@ public class WALBuffer extends AbstractWALBuffer {
     final List<Checkpoint> checkpoints = new ArrayList<>();
     final List<WALFlushListener> fsyncListeners = new ArrayList<>();
     WALFlushListener rollWALFileWriterListener = null;
-
-    final Set<Long> memTableIds = new HashSet<>();
   }
 
   /** This task serializes WALEntry to workingBuffer and will call fsync at 
last. */
@@ -319,7 +309,6 @@ public class WALBuffer extends AbstractWALBuffer {
       info.metaData.add(size, searchIndex);
       walEntry.getWalFlushListener().getWalEntryHandler().setSize(size);
       info.fsyncListeners.add(walEntry.getWalFlushListener());
-      info.memTableIds.add(walEntry.getMemTableId());
     }
 
     /**
@@ -521,11 +510,6 @@ public class WALBuffer extends AbstractWALBuffer {
       } finally {
         switchSyncingBufferToIdle();
       }
-      long walFileVersion =
-          
WALFileUtils.parseVersionId(currentWALFileWriter.getLogFile().getName());
-      memTableIdsOfWal
-          .computeIfAbsent(walFileVersion, memTableIds -> new HashSet<>())
-          .addAll(info.memTableIds);
 
       boolean forceSuccess = false;
       // try to roll log writer
@@ -700,12 +684,4 @@ public class WALBuffer extends AbstractWALBuffer {
   public CheckpointManager getCheckpointManager() {
     return checkpointManager;
   }
-
-  public Map<Long, Set<Long>> getMemTableIdsOfWal() {
-    return memTableIdsOfWal;
-  }
-
-  public void removeMemTableIdsOfWal(Long walVersionId) {
-    this.memTableIdsOfWal.remove(walVersionId);
-  }
 }
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 6c8b81300a6..50ac8ca6155 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
@@ -40,7 +40,6 @@ import java.nio.ByteBuffer;
 import java.nio.file.Files;
 import java.util.ArrayList;
 import java.util.Collections;
-import java.util.Comparator;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -89,7 +88,7 @@ public class CheckpointManager implements AutoCloseable {
     logHeader();
   }
 
-  public List<MemTableInfo> activeOrPinnedMemTables() {
+  public List<MemTableInfo> snapshotMemTableInfos() {
     infoLock.lock();
     try {
       return new ArrayList<>(memTableId2Info.values());
@@ -122,7 +121,7 @@ public class CheckpointManager implements AutoCloseable {
    */
   private void makeGlobalInfoCP() {
     long start = System.nanoTime();
-    List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
+    List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
     memTableInfos.removeIf(MemTableInfo::isFlushed);
     Checkpoint checkpoint = new 
Checkpoint(CheckpointType.GLOBAL_MEMORY_TABLE_INFO, memTableInfos);
     logByCachedByteBuffer(checkpoint);
@@ -316,13 +315,20 @@ public class CheckpointManager implements AutoCloseable {
   }
   // endregion
 
-  /** Get MemTableInfo of oldest unpinned MemTable, whose first version id is 
smallest. */
-  public MemTableInfo getOldestUnpinnedMemTableInfo() {
+  /** Get MemTableInfo of oldest MemTable, whose first version id is smallest. 
*/
+  public MemTableInfo getOldestMemTableInfo() {
     // find oldest memTable
-    return activeOrPinnedMemTables().stream()
-        .filter(memTableInfo -> !memTableInfo.isPinned())
-        .min(Comparator.comparingLong(MemTableInfo::getMemTableId))
-        .orElse(null);
+    List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
+    if (memTableInfos.isEmpty()) {
+      return null;
+    }
+    MemTableInfo oldestMemTableInfo = memTableInfos.get(0);
+    for (MemTableInfo memTableInfo : memTableInfos) {
+      if (oldestMemTableInfo.getFirstFileVersionId() > 
memTableInfo.getFirstFileVersionId()) {
+        oldestMemTableInfo = memTableInfo;
+      }
+    }
+    return oldestMemTableInfo;
   }
 
   /**
@@ -331,7 +337,7 @@ public class CheckpointManager implements AutoCloseable {
    * @return Return {@link Long#MIN_VALUE} if no file is valid
    */
   public long getFirstValidWALVersionId() {
-    List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
+    List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
     long firstValidVersionId = memTableInfos.isEmpty() ? Long.MIN_VALUE : 
Long.MAX_VALUE;
     for (MemTableInfo memTableInfo : memTableInfos) {
       firstValidVersionId = Math.min(firstValidVersionId, 
memTableInfo.getFirstFileVersionId());
@@ -341,7 +347,7 @@ public class CheckpointManager implements AutoCloseable {
 
   /** Get total cost of active memTables. */
   public long getTotalCostOfActiveMemTables() {
-    List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
+    List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
     long totalCost = 0;
     for (MemTableInfo memTableInfo : memTableInfos) {
       // flushed memTables are not active
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 a2a6c9ea342..9501834f475 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
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.db.storageengine.dataregion.wal.node;
 
+import org.apache.iotdb.commons.conf.IoTDBConstant;
 import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.file.SystemFileFactory;
 import org.apache.iotdb.commons.utils.TestOnly;
@@ -65,19 +66,18 @@ import java.io.File;
 import java.io.FileNotFoundException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
-import java.util.Arrays;
 import java.util.Collections;
 import java.util.LinkedList;
 import java.util.List;
 import java.util.ListIterator;
 import java.util.Map;
 import java.util.NoSuchElementException;
-import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicLong;
-import java.util.stream.Collectors;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
 
 /**
  * This class encapsulates {@link IWALBuffer} and {@link CheckpointManager}. 
If search is enabled,
@@ -150,7 +150,6 @@ public class WALNode implements IWALNode {
   }
 
   private WALFlushListener log(WALEntry walEntry) {
-
     buffer.write(walEntry);
     // set handler for pipe
     walEntry.getWalFlushListener().getWalEntryHandler().setWalNode(this, 
walEntry.getMemTableId());
@@ -232,18 +231,16 @@ public class WALNode implements IWALNode {
   }
 
   private class DeleteOutdatedFileTask implements Runnable {
-    private File[] sortedWalFilesExcludingLast;
-
-    private List<MemTableInfo> activeOrPinnedMemTables;
-
-    private Map<Long, Set<Long>> memTableIdsOfWalMap = new 
ConcurrentHashMap<>();
     private static final int MAX_RECURSION_TIME = 5;
-
+    // .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 = 0d;
 
     private List<Long> pinnedMemTableIds;
 
+    private File[] filesShouldDelete;
+
     private int fileIndexAfterFilterSafelyDeleteIndex = Integer.MAX_VALUE;
     private List<Long> successfullyDeleted;
     private long deleteFileSize;
@@ -254,46 +251,21 @@ public class WALNode implements IWALNode {
       // Do nothing
     }
 
-    private boolean initAndCheckIfNeedContinue() {
-      rollWalFileIfHaveNoActiveMemTable();
-      File[] allWalFilesOfOneNode = WALFileUtils.listAllWALFiles(logDirectory);
-      if (allWalFilesOfOneNode == null || allWalFilesOfOneNode.length <= 1) {
-        if (logger.isDebugEnabled()) {
-          logger.debug(
-              "wal node-{}:no wal file or wal file number less than or equal 
to one was found",
-              identifier);
-        }
-        return false;
+    private void init() {
+      this.firstValidVersionId = initFirstValidWALVersionId();
+      this.filesShouldDelete = 
logDirectory.listFiles(this::filterFilesToDelete);
+      if (filesShouldDelete == null) {
+        filesShouldDelete = new File[0];
       }
-      WALFileUtils.ascSortByVersionId(allWalFilesOfOneNode);
-      this.sortedWalFilesExcludingLast =
-          Arrays.copyOfRange(allWalFilesOfOneNode, 0, 
allWalFilesOfOneNode.length - 1);
-      this.activeOrPinnedMemTables = 
checkpointManager.activeOrPinnedMemTables();
-      this.memTableIdsOfWalMap = buffer.getMemTableIdsOfWal();
-
       this.pinnedMemTableIds = initPinnedMemTableIds();
+      WALFileUtils.ascSortByVersionId(filesShouldDelete);
       this.fileIndexAfterFilterSafelyDeleteIndex = 
initFileIndexAfterFilterSafelyDeleteIndex();
       this.successfullyDeleted = new ArrayList<>();
       this.deleteFileSize = 0;
-      return true;
-    }
-
-    /**
-     * 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
-     */
-    public void rollWalFileIfHaveNoActiveMemTable() {
-      long firstVersionId = checkpointManager.getFirstValidWALVersionId();
-      if (firstVersionId == Long.MIN_VALUE) {
-        // roll wal log writer to delete current wal file
-        if (buffer.getCurrentWALFileSize() > 0) {
-          rollWALFile();
-        }
-      }
     }
 
     private List<Long> initPinnedMemTableIds() {
-      List<MemTableInfo> memTableInfos = 
checkpointManager.activeOrPinnedMemTables();
+      List<MemTableInfo> memTableInfos = 
checkpointManager.snapshotMemTableInfos();
       if (memTableInfos.isEmpty()) {
         return new ArrayList<>();
       }
@@ -311,11 +283,8 @@ public class WALNode implements IWALNode {
       // 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.
       while (recursionTime < MAX_RECURSION_TIME) {
-        // init delete outdated file task fields, if the number of wal files 
is less than one, the
-        // subsequent logic is not executed
-        if (!initAndCheckIfNeedContinue()) {
-          break;
-        }
+        // init delete outdated file task fields
+        init();
 
         // delete outdated WAL files and record which delete successfully and 
which delete failed.
         deleteOutdatedFilesAndUpdateMetric();
@@ -354,18 +323,27 @@ public class WALNode implements IWALNode {
     }
 
     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;
+      }
+
       if (!pinnedMemTableIds.isEmpty()
-          || fileIndexAfterFilterSafelyDeleteIndex < 
sortedWalFilesExcludingLast.length) {
+          || fileIndexAfterFilterSafelyDeleteIndex < filesShouldDelete.length) 
{
         if (logger.isDebugEnabled()) {
           StringBuilder summary =
               new StringBuilder(
                   String.format(
-                      "wal node-%s delete outdated files summary:the range is: 
[%d,%d], delete successful is [%s], safely delete file index is: [%s].The 
following reasons influenced the result: %s",
+                      "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(sortedWalFilesExcludingLast[0].getName()),
+                      
WALFileUtils.parseVersionId(filesShouldDelete[0].getName()),
                       WALFileUtils.parseVersionId(
-                          
sortedWalFilesExcludingLast[sortedWalFilesExcludingLast.length - 1]
-                              .getName()),
+                          filesShouldDelete[filesShouldDelete.length - 
1].getName()),
                       StringUtils.join(successfullyDeleted, ","),
                       fileIndexAfterFilterSafelyDeleteIndex,
                       System.getProperty("line.separator")));
@@ -377,7 +355,7 @@ public class WALNode implements IWALNode {
                 .append(".")
                 .append(System.getProperty("line.separator"));
           }
-          if (fileIndexAfterFilterSafelyDeleteIndex < 
sortedWalFilesExcludingLast.length) {
+          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",
@@ -389,32 +367,33 @@ public class WALNode implements IWALNode {
 
       } else {
         logger.debug(
-            "Successfully delete {} outdated wal files for wal node-{}",
+            "Successfully delete {} outdated wal files for wal node-{},first 
valid version id is {}",
             successfullyDeleted.size(),
-            identifier);
+            identifier,
+            firstValidVersionId);
       }
     }
 
     /** Delete obsolete wal files while recording which succeeded or failed */
     private void deleteOutdatedFilesAndUpdateMetric() {
-      for (File currentWal : sortedWalFilesExcludingLast) {
-        long searchIndex = 
WALFileUtils.parseStartSearchIndex(currentWal.getName());
-        WALFileStatus walFileStatus = 
WALFileUtils.parseStatusCode(currentWal.getName());
-        long versionId = WALFileUtils.parseVersionId(currentWal.getName());
-        if (canDeleteFile(searchIndex, walFileStatus, versionId)) {
-          long fileSize = currentWal.length();
-          if (currentWal.delete()) {
-            deleteFileSize += fileSize;
-            Long memTableRamCostSum = 
walFileVersionId2MemTablesTotalCost.remove(versionId);
-            if (memTableRamCostSum != null) {
-              totalCostOfFlushedMemTables.addAndGet(-memTableRamCostSum);
-            }
-            buffer.removeMemTableIdsOfWal(versionId);
-            successfullyDeleted.add(versionId);
-          } else {
-            logger.info(
-                "Fail to delete outdated wal file {} of wal node-{}.", 
currentWal, identifier);
+      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);
@@ -424,15 +403,15 @@ public class WALNode implements IWALNode {
     private int initFileIndexAfterFilterSafelyDeleteIndex() {
       int endFileIndex =
           safelyDeletedSearchIndex == DEFAULT_SAFELY_DELETED_SEARCH_INDEX
-              ? sortedWalFilesExcludingLast.length
+              ? filesShouldDelete.length
               : WALFileUtils.binarySearchFileBySearchIndex(
-                  sortedWalFilesExcludingLast, safelyDeletedSearchIndex + 1);
+                  filesShouldDelete, safelyDeletedSearchIndex + 1);
       // delete files whose file status is CONTAINS_NONE_SEARCH_INDEX
       if (endFileIndex == -1) {
         endFileIndex = 0;
       }
-      while (endFileIndex < sortedWalFilesExcludingLast.length) {
-        if 
(WALFileUtils.parseStatusCode(sortedWalFilesExcludingLast[endFileIndex].getName())
+      while (endFileIndex < filesShouldDelete.length) {
+        if 
(WALFileUtils.parseStatusCode(filesShouldDelete[endFileIndex].getName())
             == WALFileStatus.CONTAINS_SEARCH_INDEX) {
           break;
         }
@@ -441,6 +420,17 @@ public class WALNode implements IWALNode {
       return endFileIndex;
     }
 
+    private boolean filterFilesToDelete(File dir, String name) {
+      Pattern pattern = WALFileUtils.WAL_FILE_NAME_PATTERN;
+      Matcher matcher = pattern.matcher(name);
+      boolean toDelete = false;
+      if (matcher.find()) {
+        long versionId = 
Long.parseLong(matcher.group(IoTDBConstant.WAL_VERSION_ID));
+        toDelete = versionId < firstValidVersionId;
+      }
+      return toDelete;
+    }
+
     /** Return true iff effective information ratio is too small or disk usage 
is too large. */
     private boolean shouldSnapshotOrFlush() {
       return effectiveInfoRatio < config.getWalMinEffectiveInfoRatio()
@@ -457,7 +447,7 @@ public class WALNode implements IWALNode {
         return false;
       }
       // find oldest memTable
-      MemTableInfo oldestMemTableInfo = 
checkpointManager.getOldestUnpinnedMemTableInfo();
+      MemTableInfo oldestMemTableInfo = 
checkpointManager.getOldestMemTableInfo();
       if (oldestMemTableInfo == null) {
         return false;
       }
@@ -594,25 +584,22 @@ public class WALNode implements IWALNode {
       }
     }
 
-    public boolean isContainsActiveOrPinnedMemTable(Long versionId) {
-      Set<Long> memTableIdsOfCurrentWal = memTableIdsOfWalMap.get(versionId);
-      // If this set is empty, there is a case where WalEntry has been logged 
but not persisted,
-      // because WalEntry is persisted asynchronously. In this case, the file 
cannot be deleted
-      // directly, so it is considered active
-      if (memTableIdsOfCurrentWal == null || 
memTableIdsOfCurrentWal.isEmpty()) {
-        return true;
+    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 !Collections.disjoint(
-          activeOrPinnedMemTables.stream()
-              .map(MemTableInfo::getMemTableId)
-              .collect(Collectors.toSet()),
-          memTableIdsOfCurrentWal);
-    }
-
-    private boolean canDeleteFile(long searchIndex, WALFileStatus 
walFileStatus, long versionId) {
-      return (searchIndex < safelyDeletedSearchIndex
-              || walFileStatus == WALFileStatus.CONTAINS_NONE_SEARCH_INDEX)
-          && !isContainsActiveOrPinnedMemTable(versionId);
+      return firstVersionId;
     }
   }
 
@@ -991,9 +978,4 @@ public class WALNode implements IWALNode {
   public void setBufferSize(int size) {
     buffer.setBufferSize(size);
   }
-
-  @TestOnly
-  public WALBuffer getWALBuffer() {
-    return buffer;
-  }
 }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALEntryHandlerTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALEntryHandlerTest.java
index 58e9aefec32..08feee58228 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALEntryHandlerTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALEntryHandlerTest.java
@@ -124,12 +124,9 @@ public class WALEntryHandlerTest {
     // pin memTable
     WALEntryHandler handler = flushListener.getWalEntryHandler();
     handler.pinMemTable();
+    walNode1.onMemTableFlushed(memTable);
     // roll wal file
     walNode1.rollWALFile();
-    InsertRowNode node2 = getInsertRowNode(devicePath, 
System.currentTimeMillis());
-    node2.setSearchIndex(2);
-    walNode1.log(memTable.getMemTableId(), node2);
-    walNode1.onMemTableFlushed(memTable);
     walNode1.rollWALFile();
     // find node1
     ConsensusReqReader.ReqIterator itr = walNode1.getReqIterator(1);
@@ -176,11 +173,13 @@ public class WALEntryHandlerTest {
     // unpin 1
     CheckpointManager checkpointManager = walNode1.getCheckpointManager();
     handler.unpinMemTable();
-    MemTableInfo oldestMemTableInfo = 
checkpointManager.getOldestUnpinnedMemTableInfo();
-    assertNull(oldestMemTableInfo);
+    MemTableInfo oldestMemTableInfo = 
checkpointManager.getOldestMemTableInfo();
+    assertEquals(memTable.getMemTableId(), oldestMemTableInfo.getMemTableId());
+    assertNull(oldestMemTableInfo.getMemTable());
+    assertTrue(oldestMemTableInfo.isPinned());
     // unpin 2
     handler.unpinMemTable();
-    assertNull(checkpointManager.getOldestUnpinnedMemTableInfo());
+    assertNull(checkpointManager.getOldestMemTableInfo());
   }
 
   @Test
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WalDeleteOutdatedNewTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WalDeleteOutdatedNewTest.java
deleted file mode 100644
index c528965b58f..00000000000
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WalDeleteOutdatedNewTest.java
+++ /dev/null
@@ -1,585 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied.  See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-package org.apache.iotdb.db.storageengine.dataregion.wal.node;
-
-import org.apache.iotdb.commons.exception.IllegalPathException;
-import org.apache.iotdb.commons.path.PartialPath;
-import org.apache.iotdb.consensus.iot.log.ConsensusReqReader;
-import org.apache.iotdb.db.conf.IoTDBConfig;
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId;
-import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
-import org.apache.iotdb.db.storageengine.dataregion.memtable.IMemTable;
-import org.apache.iotdb.db.storageengine.dataregion.memtable.PrimitiveMemTable;
-import 
org.apache.iotdb.db.storageengine.dataregion.wal.exception.MemTablePinException;
-import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALEntryHandler;
-import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileUtils;
-import 
org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALInsertNodeCache;
-import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALMode;
-import 
org.apache.iotdb.db.storageengine.dataregion.wal.utils.listener.WALFlushListener;
-import org.apache.iotdb.db.utils.EnvironmentUtils;
-import org.apache.iotdb.db.utils.constant.TestConstant;
-import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
-import org.apache.iotdb.tsfile.utils.Binary;
-import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
-
-import org.awaitility.Awaitility;
-import org.junit.After;
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
-
-import java.io.File;
-import java.util.Map;
-import java.util.Set;
-
-public class WalDeleteOutdatedNewTest {
-  private static final IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
-  private static final String identifier1 = String.valueOf(Integer.MAX_VALUE);
-  private static final String logDirectory1 = 
TestConstant.BASE_OUTPUT_PATH.concat("1/2910/");
-  private static final String databasePath = "root.test_sg";
-  private static final String devicePath = databasePath + ".test_d";
-  private static final String dataRegionId = "1";
-  private WALMode prevMode;
-  private boolean prevIsClusterMode;
-  private WALNode walNode1;
-
-  @Before
-  public void setUp() throws Exception {
-    EnvironmentUtils.cleanDir(logDirectory1);
-    prevMode = config.getWalMode();
-    prevIsClusterMode = config.isClusterMode();
-    config.setWalMode(WALMode.SYNC);
-    config.setClusterMode(true);
-    walNode1 = new WALNode(identifier1, logDirectory1);
-  }
-
-  @After
-  public void tearDown() throws Exception {
-    walNode1.close();
-    config.setWalMode(prevMode);
-    config.setClusterMode(prevIsClusterMode);
-    EnvironmentUtils.cleanDir(logDirectory1);
-
-    WALInsertNodeCache.getInstance(1).clear();
-  }
-
-  /**
-   * The simulation here is to write the last file, because serialization and 
disk flushing
-   * operations are asynchronous, so you have to wait until all the entries 
are processed to get the
-   * correct result, when WalEntry is not consumed to get memTableIdsOfWal, 
the result is not
-   * accurate, so when the actual deletion of expired wal files, Don't read 
the last wal file.
-   */
-  @Test
-  public void test01() throws IllegalPathException {
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 4));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 5));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 6));
-
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(1, memTableIdsOfWal.size());
-    Assert.assertEquals(2, memTableIdsOfWal.get(0L).size());
-    Assert.assertEquals(1, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-    walNode1.close();
-  }
-
-  /**
-   * Ensure that the memtableIds maintained by each wal file are accurate:<br>
-   * 1. _0-1-1.wal:memTable0、memTable1 <br>
-   * 2. roll wal file <br>
-   * 3. _1-6-1.wal: memTable1 <br>
-   * 4. wait until all walEntry consumed
-   */
-  @Test
-  public void test02() throws IllegalPathException {
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 4));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 5));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 6));
-
-    walNode1.rollWALFile();
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 7));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 8));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 9));
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(2, memTableIdsOfWal.size());
-    Assert.assertEquals(1, memTableIdsOfWal.get(1L).size());
-    Assert.assertEquals(2, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-  }
-
-  /**
-   * Ensure that files that can be cleaned can be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 、 memTable1 <br>
-   * 2. roll wal file <br>
-   * 3. _1-1-1.wal: memTable1 <br>
-   * 4. wait until all walEntry consumed <br>
-   * 5. memTable0 flush, memTable1 flush <br>
-   * 6. delete outdated wal files
-   */
-  @Test
-  public void test03() throws IllegalPathException {
-
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    walNode1.onMemTableFlushed(memTable0);
-    walNode1.onMemTableFlushed(memTable1);
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-    // before deleted
-    Assert.assertEquals(2, memTableIdsOfWal.size());
-    Assert.assertEquals(2, memTableIdsOfWal.get(0L).size());
-    File[] files = WALFileUtils.listAllWALFiles(new File(logDirectory1));
-    Assert.assertEquals(2, files.length);
-
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-
-    // after deleted
-    Assert.assertEquals(0, memTableIdsOfWalAfter.size());
-    File[] filesAfter = WALFileUtils.listAllWALFiles(new File(logDirectory1));
-    Assert.assertEquals(1, filesAfter.length);
-  }
-
-  /**
-   * Ensure that files that can be cleaned can be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. roll wal file <br>
-   * 3. _1-1-0.wal: memTable1 <br>
-   * 4. roll wal file <br>
-   * 5. _2-1-1.wal: memTable1 <br>
-   * 6. wait until all walEntry consumed <br>
-   * 7. memTable0 flush, memTable1 flush, memTable2 flush <br>
-   * 6. delete outdated wal files
-   */
-  @Test
-  public void test04() throws IllegalPathException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.rollWALFile();
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-    walNode1.onMemTableFlushed(memTable2);
-    walNode1.onMemTableFlushed(memTable0);
-    walNode1.onMemTableFlushed(memTable1);
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(3, memTableIdsOfWal.size());
-    Assert.assertEquals(3, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(0, memTableIdsOfWalAfter.size());
-    Assert.assertEquals(1, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-  }
-
-  /**
-   * Ensure that wal pinned to memtable cannot be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. pin memTable0 <br>
-   * 3. memTable0 flush <br>
-   * 4. roll wal file <br>
-   * 5. _1-1-1.wal: memTable0、memTable1 <br>
-   * 6. roll wal file <br>
-   * 7. _2-1-1.wal: memTable1 <br>
-   * 8. roll wal file <br>
-   * 9. _2-1-1.wal: memTable1 <br>
-   * 10. wait until all walEntry consumed <br>
-   * 11. memTable0 flush, memTable1 flush <br>
-   * 12. delete outdated wal files
-   */
-  @Test
-  public void test05() throws IllegalPathException, MemTablePinException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    WALFlushListener listener =
-        walNode1.log(
-            memTable0.getMemTableId(),
-            generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-
-    // pin memTable
-    WALEntryHandler handler = listener.getWalEntryHandler();
-    handler.pinMemTable();
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-    walNode1.rollWALFile();
-
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 4));
-    walNode1.rollWALFile();
-
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 5));
-    walNode1.onMemTableFlushed(memTable0);
-    walNode1.onMemTableFlushed(memTable1);
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(4, memTableIdsOfWal.size());
-    Assert.assertEquals(4, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(3, memTableIdsOfWalAfter.size());
-    Assert.assertEquals(3, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-  }
-
-  /**
-   * Ensure that the flushed wal related to memtable cannot be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. roll wal file <br>
-   * 3. _1-1-1.wal: memTable0 <br>
-   * 4. roll wal file <br>
-   * 5. _2-1-1.wal: memTable0 <br>
-   * 6. roll wal file <br>
-   * 7. _2-1-1.wal: memTable0 <br>
-   * 8. wait until all walEntry consumed <br>
-   * 9. delete outdated wal files
-   */
-  @Test
-  public void test06() throws IllegalPathException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 4));
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(4, memTableIdsOfWal.size());
-    Assert.assertEquals(4, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(4, memTableIdsOfWalAfter.size());
-    Assert.assertEquals(4, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-  }
-
-  /**
-   * Ensure that files that can be cleaned can be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. roll wal file <br>
-   * 3. _1-1-0.wal: memTable1、memTable2 <br>
-   * 4. roll wal file <br>
-   * 5. _2-1-0.wal: memTable2 <br>
-   * 6. roll wal file <br>
-   * 7. _3-1-0.wal: memTable3 <br>
-   * 8. roll wal file <br>
-   * 9. _4-1-0.wal: memTable3 <br>
-   * 10. wait until all walEntry consumed <br>
-   * 11. memTable1 flush, memTable2 flush, memTable3 flush <br>
-   * 12. delete outdated wal files
-   */
-  @Test
-  public void test07() throws IllegalPathException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.rollWALFile();
-    IMemTable memTable3 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable3, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable3.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.rollWALFile();
-    walNode1.log(
-        memTable3.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.onMemTableFlushed(memTable1);
-    walNode1.onMemTableFlushed(memTable2);
-    walNode1.onMemTableFlushed(memTable3);
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    Map<Long, Set<Long>> memTableIdsOfWal = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(5, memTableIdsOfWal.size());
-    Assert.assertEquals(5, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(2, memTableIdsOfWalAfter.size());
-    Assert.assertEquals(2, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-
-    walNode1.onMemTableFlushed(memTable0);
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-    walNode1.deleteOutdatedFiles();
-    Map<Long, Set<Long>> memTableIdsOfWalAfterAfter = 
walNode1.getWALBuffer().getMemTableIdsOfWal();
-    Assert.assertEquals(0, memTableIdsOfWalAfterAfter.size());
-    Assert.assertEquals(1, WALFileUtils.listAllWALFiles(new 
File(logDirectory1)).length);
-  }
-
-  /**
-   * Ensure that files that can be cleaned can be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. roll wal file <br>
-   * 3. _1-1-0.wal: memTable1<br>
-   * 4. memTable1 flush <br>
-   * 5. roll wal file <br>
-   * 6. _2-1-0.wal: memTable2 <br>
-   * 7. wait until all walEntry consumed <br>
-   * 8. delete outdated wal files
-   */
-  @Test
-  public void test08() throws IllegalPathException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-
-    ConsensusReqReader.ReqIterator itr1 = walNode1.getReqIterator(1);
-    Assert.assertFalse(itr1.hasNext());
-
-    IMemTable memTable1 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable1, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable1.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), -1));
-    walNode1.onMemTableFlushed(memTable1);
-    walNode1.rollWALFile();
-
-    ConsensusReqReader.ReqIterator itr2 = walNode1.getReqIterator(1);
-    Assert.assertTrue(itr2.hasNext());
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    ConsensusReqReader.ReqIterator itr3 = walNode1.getReqIterator(1);
-    Assert.assertTrue(itr3.hasNext());
-    walNode1.deleteOutdatedFiles();
-
-    ConsensusReqReader.ReqIterator itr4 = walNode1.getReqIterator(1);
-    Assert.assertFalse(itr4.hasNext());
-    walNode1.rollWALFile();
-    Assert.assertTrue(itr4.hasNext());
-  }
-
-  /**
-   * Ensure that files that can be cleaned can be deleted: <br>
-   * 1. _0-0-1.wal: memTable0 <br>
-   * 2. roll wal file <br>
-   * 3. _2-1-0.wal: memTable2 <br>
-   * 4. wait until all walEntry consumed <br>
-   * 5. delete outdated wal files
-   */
-  @Test
-  public void test09() throws IllegalPathException {
-    IMemTable memTable0 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable0, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable0.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 1));
-    walNode1.rollWALFile();
-
-    IMemTable memTable2 = new PrimitiveMemTable(databasePath, dataRegionId);
-    walNode1.onMemTableCreated(memTable2, logDirectory1 + "/" + "fake.tsfile");
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 2));
-    walNode1.log(
-        memTable2.getMemTableId(),
-        generateInsertRowNode(devicePath, System.currentTimeMillis(), 3));
-    Awaitility.await().until(() -> walNode1.isAllWALEntriesConsumed());
-
-    ConsensusReqReader.ReqIterator itr3 = walNode1.getReqIterator(1);
-    Assert.assertFalse(itr3.hasNext());
-  }
-
-  public static InsertRowNode generateInsertRowNode(String devicePath, long 
time, long searchIndex)
-      throws IllegalPathException {
-    TSDataType[] dataTypes =
-        new TSDataType[] {
-          TSDataType.DOUBLE,
-          TSDataType.FLOAT,
-          TSDataType.INT64,
-          TSDataType.INT32,
-          TSDataType.BOOLEAN,
-          TSDataType.TEXT
-        };
-
-    Object[] columns = new Object[6];
-    columns[0] = 1.0d;
-    columns[1] = 2f;
-    columns[2] = 10000L;
-    columns[3] = 100;
-    columns[4] = false;
-    columns[5] = new Binary("hh" + 0, TSFileConfig.STRING_CHARSET);
-
-    InsertRowNode node =
-        new InsertRowNode(
-            new PlanNodeId(""),
-            new PartialPath(devicePath),
-            false,
-            new String[] {"s1", "s2", "s3", "s4", "s5", "s6"},
-            dataTypes,
-            time,
-            columns,
-            false);
-    MeasurementSchema[] schemas = new MeasurementSchema[6];
-    for (int i = 0; i < 6; i++) {
-      schemas[i] = new MeasurementSchema("s" + (i + 1), dataTypes[i]);
-    }
-    node.setMeasurementSchemas(schemas);
-    node.setSearchIndex(searchIndex);
-    return node;
-  }
-}

Reply via email to