This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch revert-11682-f_wal_delete in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 56156eab7d441f38b52587b998c041854e007e07 Author: Haonan <[email protected]> AuthorDate: Fri Jan 5 09:35:58 2024 +0800 Revert "Optimized wal file deletion algorithm (#11682)" 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; - } -}
