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 9fdeb95b6f6 Optimized wal file deletion algorithm (#11682)
9fdeb95b6f6 is described below
commit 9fdeb95b6f6d7bc456213a690ec536cf4d0959a8
Author: Zhijia Cao <[email protected]>
AuthorDate: Thu Jan 4 14:21:47 2024 +0800
Optimized wal file deletion algorithm (#11682)
---
.../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, 728 insertions(+), 105 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 08bca5f6649..a1c3dcd6672 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,6 +96,7 @@ 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 cdae8baa1c7..a93ac366748 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,6 +33,7 @@ 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;
@@ -45,9 +46,13 @@ 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;
@@ -102,6 +107,9 @@ 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);
}
@@ -186,6 +194,8 @@ 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. */
@@ -309,6 +319,7 @@ 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());
}
/**
@@ -510,6 +521,11 @@ 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
@@ -684,4 +700,12 @@ 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 50ac8ca6155..6c8b81300a6 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,6 +40,7 @@ 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;
@@ -88,7 +89,7 @@ public class CheckpointManager implements AutoCloseable {
logHeader();
}
- public List<MemTableInfo> snapshotMemTableInfos() {
+ public List<MemTableInfo> activeOrPinnedMemTables() {
infoLock.lock();
try {
return new ArrayList<>(memTableId2Info.values());
@@ -121,7 +122,7 @@ public class CheckpointManager implements AutoCloseable {
*/
private void makeGlobalInfoCP() {
long start = System.nanoTime();
- List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
+ List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
memTableInfos.removeIf(MemTableInfo::isFlushed);
Checkpoint checkpoint = new
Checkpoint(CheckpointType.GLOBAL_MEMORY_TABLE_INFO, memTableInfos);
logByCachedByteBuffer(checkpoint);
@@ -315,20 +316,13 @@ public class CheckpointManager implements AutoCloseable {
}
// endregion
- /** Get MemTableInfo of oldest MemTable, whose first version id is smallest.
*/
- public MemTableInfo getOldestMemTableInfo() {
+ /** Get MemTableInfo of oldest unpinned MemTable, whose first version id is
smallest. */
+ public MemTableInfo getOldestUnpinnedMemTableInfo() {
// find oldest memTable
- 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;
+ return activeOrPinnedMemTables().stream()
+ .filter(memTableInfo -> !memTableInfo.isPinned())
+ .min(Comparator.comparingLong(MemTableInfo::getMemTableId))
+ .orElse(null);
}
/**
@@ -337,7 +331,7 @@ public class CheckpointManager implements AutoCloseable {
* @return Return {@link Long#MIN_VALUE} if no file is valid
*/
public long getFirstValidWALVersionId() {
- List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
+ List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
long firstValidVersionId = memTableInfos.isEmpty() ? Long.MIN_VALUE :
Long.MAX_VALUE;
for (MemTableInfo memTableInfo : memTableInfos) {
firstValidVersionId = Math.min(firstValidVersionId,
memTableInfo.getFirstFileVersionId());
@@ -347,7 +341,7 @@ public class CheckpointManager implements AutoCloseable {
/** Get total cost of active memTables. */
public long getTotalCostOfActiveMemTables() {
- List<MemTableInfo> memTableInfos = snapshotMemTableInfos();
+ List<MemTableInfo> memTableInfos = activeOrPinnedMemTables();
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 9501834f475..a2a6c9ea342 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,7 +19,6 @@
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;
@@ -66,18 +65,19 @@ 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.regex.Matcher;
-import java.util.regex.Pattern;
+import java.util.stream.Collectors;
/**
* This class encapsulates {@link IWALBuffer} and {@link CheckpointManager}.
If search is enabled,
@@ -150,6 +150,7 @@ public class WALNode implements IWALNode {
}
private WALFlushListener log(WALEntry walEntry) {
+
buffer.write(walEntry);
// set handler for pipe
walEntry.getWalFlushListener().getWalEntryHandler().setWalNode(this,
walEntry.getMemTableId());
@@ -231,16 +232,18 @@ 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;
@@ -251,21 +254,46 @@ public class WALNode implements IWALNode {
// Do nothing
}
- private void init() {
- this.firstValidVersionId = initFirstValidWALVersionId();
- this.filesShouldDelete =
logDirectory.listFiles(this::filterFilesToDelete);
- if (filesShouldDelete == null) {
- filesShouldDelete = new File[0];
+ 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;
}
+ 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.snapshotMemTableInfos();
+ List<MemTableInfo> memTableInfos =
checkpointManager.activeOrPinnedMemTables();
if (memTableInfos.isEmpty()) {
return new ArrayList<>();
}
@@ -283,8 +311,11 @@ 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
- init();
+ // 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;
+ }
// delete outdated WAL files and record which delete successfully and
which delete failed.
deleteOutdatedFilesAndUpdateMetric();
@@ -323,27 +354,18 @@ 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 < filesShouldDelete.length)
{
+ || fileIndexAfterFilterSafelyDeleteIndex <
sortedWalFilesExcludingLast.length) {
if (logger.isDebugEnabled()) {
StringBuilder summary =
new StringBuilder(
String.format(
- "wal node-%s delete outdated files summary:the range
that should be removed is: [%d,%d], delete successful is [%s], end file index
is: [%s].The following reasons influenced the result: %s",
+ "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",
identifier,
-
WALFileUtils.parseVersionId(filesShouldDelete[0].getName()),
+
WALFileUtils.parseVersionId(sortedWalFilesExcludingLast[0].getName()),
WALFileUtils.parseVersionId(
- filesShouldDelete[filesShouldDelete.length -
1].getName()),
+
sortedWalFilesExcludingLast[sortedWalFilesExcludingLast.length - 1]
+ .getName()),
StringUtils.join(successfullyDeleted, ","),
fileIndexAfterFilterSafelyDeleteIndex,
System.getProperty("line.separator")));
@@ -355,7 +377,7 @@ public class WALNode implements IWALNode {
.append(".")
.append(System.getProperty("line.separator"));
}
- if (fileIndexAfterFilterSafelyDeleteIndex <
filesShouldDelete.length) {
+ if (fileIndexAfterFilterSafelyDeleteIndex <
sortedWalFilesExcludingLast.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",
@@ -367,33 +389,32 @@ public class WALNode implements IWALNode {
} else {
logger.debug(
- "Successfully delete {} outdated wal files for wal node-{},first
valid version id is {}",
+ "Successfully delete {} outdated wal files for wal node-{}",
successfullyDeleted.size(),
- identifier,
- firstValidVersionId);
+ identifier);
}
}
/** Delete obsolete wal files while recording which succeeded or failed */
private void deleteOutdatedFilesAndUpdateMetric() {
- if (filesShouldDelete.length == 0) {
- return;
- }
- for (int i = 0; i < fileIndexAfterFilterSafelyDeleteIndex; ++i) {
- long fileSize = filesShouldDelete[i].length();
- long versionId =
WALFileUtils.parseVersionId(filesShouldDelete[i].getName());
- if (filesShouldDelete[i].delete()) {
- deleteFileSize += fileSize;
- Long memTableRamCostSum =
walFileVersionId2MemTablesTotalCost.remove(versionId);
- if (memTableRamCostSum != null) {
- totalCostOfFlushedMemTables.addAndGet(-memTableRamCostSum);
+ 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);
}
- successfullyDeleted.add(versionId);
- } else {
- logger.info(
- "Fail to delete outdated wal file {} of wal node-{}.",
- filesShouldDelete[i],
- identifier);
}
}
buffer.subtractDiskUsage(deleteFileSize);
@@ -403,15 +424,15 @@ public class WALNode implements IWALNode {
private int initFileIndexAfterFilterSafelyDeleteIndex() {
int endFileIndex =
safelyDeletedSearchIndex == DEFAULT_SAFELY_DELETED_SEARCH_INDEX
- ? filesShouldDelete.length
+ ? sortedWalFilesExcludingLast.length
: WALFileUtils.binarySearchFileBySearchIndex(
- filesShouldDelete, safelyDeletedSearchIndex + 1);
+ sortedWalFilesExcludingLast, safelyDeletedSearchIndex + 1);
// delete files whose file status is CONTAINS_NONE_SEARCH_INDEX
if (endFileIndex == -1) {
endFileIndex = 0;
}
- while (endFileIndex < filesShouldDelete.length) {
- if
(WALFileUtils.parseStatusCode(filesShouldDelete[endFileIndex].getName())
+ while (endFileIndex < sortedWalFilesExcludingLast.length) {
+ if
(WALFileUtils.parseStatusCode(sortedWalFilesExcludingLast[endFileIndex].getName())
== WALFileStatus.CONTAINS_SEARCH_INDEX) {
break;
}
@@ -420,17 +441,6 @@ 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()
@@ -447,7 +457,7 @@ public class WALNode implements IWALNode {
return false;
}
// find oldest memTable
- MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestMemTableInfo();
+ MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestUnpinnedMemTableInfo();
if (oldestMemTableInfo == null) {
return false;
}
@@ -584,22 +594,25 @@ public class WALNode implements IWALNode {
}
}
- 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();
- }
+ 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;
}
- return firstVersionId;
+ 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);
}
}
@@ -978,4 +991,9 @@ 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 08feee58228..58e9aefec32 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,9 +124,12 @@ 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);
@@ -173,13 +176,11 @@ public class WALEntryHandlerTest {
// unpin 1
CheckpointManager checkpointManager = walNode1.getCheckpointManager();
handler.unpinMemTable();
- MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestMemTableInfo();
- assertEquals(memTable.getMemTableId(), oldestMemTableInfo.getMemTableId());
- assertNull(oldestMemTableInfo.getMemTable());
- assertTrue(oldestMemTableInfo.isPinned());
+ MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestUnpinnedMemTableInfo();
+ assertNull(oldestMemTableInfo);
// unpin 2
handler.unpinMemTable();
- assertNull(checkpointManager.getOldestMemTableInfo());
+ assertNull(checkpointManager.getOldestUnpinnedMemTableInfo());
}
@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
new file mode 100644
index 00000000000..c528965b58f
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WalDeleteOutdatedNewTest.java
@@ -0,0 +1,585 @@
+/*
+ * 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;
+ }
+}