This is an automated email from the ASF dual-hosted git repository.
rong 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 4cf7a483ffb Pipe: avoid overriding wal entries when loading entries
from different wal nodes (#10130)
4cf7a483ffb is described below
commit 4cf7a483ffbbba5a241a8b659f78b41fea893d2e
Author: Alan Choo <[email protected]>
AuthorDate: Mon Jun 12 22:48:09 2023 +0800
Pipe: avoid overriding wal entries when loading entries from different wal
nodes (#10130)
---
.../java/org/apache/iotdb/db/wal/node/WALNode.java | 4 +
.../apache/iotdb/db/wal/utils/WALEntryHandler.java | 9 +-
.../iotdb/db/wal/utils/WALEntryPosition.java | 15 +-
.../iotdb/db/wal/utils/WALInsertNodeCache.java | 6 +-
.../iotdb/db/wal/buffer/WALBufferCommonTest.java | 1 +
.../db/wal/checkpoint/CheckpointManagerTest.java | 1 +
.../iotdb/db/wal/node/ConsensusReqReaderTest.java | 2 +
.../iotdb/db/wal/node/WALEntryHandlerTest.java | 156 +++++++++++++++------
.../org/apache/iotdb/db/wal/node/WALNodeTest.java | 2 +
.../db/wal/recover/WALRecoverManagerTest.java | 2 +
10 files changed, 152 insertions(+), 46 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
b/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
index 3eacc838cec..98f985da322 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/node/WALNode.java
@@ -790,6 +790,10 @@ public class WALNode implements IWALNode {
checkpointManager.close();
}
+ public String getIdentifier() {
+ return identifier;
+ }
+
public File getLogDirectory() {
return logDirectory;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryHandler.java
b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryHandler.java
index e8849dc6536..52324b512dd 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryHandler.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryHandler.java
@@ -94,11 +94,18 @@ public class WALEntryHandler {
}
}
// read from the wal file
+ InsertNode node = null;
try {
- return walEntryPosition.readInsertNodeViaCache();
+ node = walEntryPosition.readInsertNodeViaCache();
} catch (Exception e) {
throw new WALPipeException("Fail to get value because the file content
isn't correct.", e);
}
+
+ if (node == null) {
+ throw new WALPipeException(
+ String.format("Fail to get the wal value of the position %s.",
walEntryPosition));
+ }
+ return node;
}
public long getMemTableId() {
diff --git
a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryPosition.java
b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryPosition.java
index adab38dc5b3..acb66a393df 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryPosition.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALEntryPosition.java
@@ -34,6 +34,7 @@ import java.util.Objects;
*/
public class WALEntryPosition {
private static final WALInsertNodeCache CACHE =
WALInsertNodeCache.getInstance();
+ private volatile String identifier = "";
private volatile long walFileVersionId = -1;
private volatile long position;
private volatile int size;
@@ -44,7 +45,8 @@ public class WALEntryPosition {
public WALEntryPosition() {}
- public WALEntryPosition(long walFileVersionId, long position, int size) {
+ public WALEntryPosition(String identifier, long walFileVersionId, long
position, int size) {
+ this.identifier = identifier;
this.walFileVersionId = walFileVersionId;
this.position = position;
this.size = size;
@@ -111,6 +113,7 @@ public class WALEntryPosition {
public void setWalNode(WALNode walNode) {
this.walNode = walNode;
+ this.identifier = walNode.getIdentifier();
}
public void setEntryPosition(long walFileVersionId, long position) {
@@ -118,6 +121,10 @@ public class WALEntryPosition {
this.walFileVersionId = walFileVersionId;
}
+ public String getIdentifier() {
+ return identifier;
+ }
+
public long getWalFileVersionId() {
return walFileVersionId;
}
@@ -140,7 +147,7 @@ public class WALEntryPosition {
@Override
public int hashCode() {
- return Objects.hash(walFileVersionId, position);
+ return Objects.hash(identifier, walFileVersionId, position);
}
@Override
@@ -152,6 +159,8 @@ public class WALEntryPosition {
return false;
}
WALEntryPosition that = (WALEntryPosition) o;
- return walFileVersionId == that.walFileVersionId && position ==
that.position;
+ return identifier.equals(that.identifier)
+ && walFileVersionId == that.walFileVersionId
+ && position == that.position;
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALInsertNodeCache.java
b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALInsertNodeCache.java
index 736009f0212..2a3810268bb 100644
--- a/server/src/main/java/org/apache/iotdb/db/wal/utils/WALInsertNodeCache.java
+++ b/server/src/main/java/org/apache/iotdb/db/wal/utils/WALInsertNodeCache.java
@@ -52,7 +52,7 @@ public class WALInsertNodeCache {
private final LoadingCache<WALEntryPosition, InsertNode> lruCache;
/** ids of all pinned memTables */
- private final Set<Long> memTablesNeedSearch = ConcurrentHashMap.newKeySet();;
+ private final Set<Long> memTablesNeedSearch = ConcurrentHashMap.newKeySet();
private WALInsertNodeCache() {
lruCache =
@@ -142,7 +142,9 @@ public class WALInsertNodeCache {
buffer.clear();
InsertNode node = parse(buffer);
if (node != null) {
- res.put(new WALEntryPosition(walFileVersionId, position,
size), node);
+ res.put(
+ new WALEntryPosition(pos.getIdentifier(),
walFileVersionId, position, size),
+ node);
}
}
position += size;
diff --git
a/server/src/test/java/org/apache/iotdb/db/wal/buffer/WALBufferCommonTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/buffer/WALBufferCommonTest.java
index 4054ca796b7..05f01f2dabc 100644
---
a/server/src/test/java/org/apache/iotdb/db/wal/buffer/WALBufferCommonTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/wal/buffer/WALBufferCommonTest.java
@@ -95,6 +95,7 @@ public abstract class WALBufferCommonTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// wait a moment
while (!walBuffer.isAllWALEntriesConsumed()) {
Thread.sleep(1_000);
diff --git
a/server/src/test/java/org/apache/iotdb/db/wal/checkpoint/CheckpointManagerTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/checkpoint/CheckpointManagerTest.java
index 1ee9d175e3f..4fb39df9ff1 100644
---
a/server/src/test/java/org/apache/iotdb/db/wal/checkpoint/CheckpointManagerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/wal/checkpoint/CheckpointManagerTest.java
@@ -114,6 +114,7 @@ public class CheckpointManagerTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// check first valid version id
assertEquals(memTablesNum / 2,
checkpointManager.getFirstValidWALVersionId());
// recover info from checkpoint file
diff --git
a/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java
index 508338a6e07..e57191e7133 100644
---
a/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java
@@ -227,6 +227,7 @@ public class ConsensusReqReaderTest {
InsertRowNode insertRowNode = getInsertRowNode(devicePath);
walNode.log(0, insertRowNode); // put -1 after 6
Assert.assertTrue(future.get());
+ checkThread.shutdown();
}
@Test
@@ -272,6 +273,7 @@ public class ConsensusReqReaderTest {
insertRowNode = getInsertRowNode(devicePath);
walNode.log(0, insertRowNode); // put -1 after 6
Assert.assertTrue(future.get());
+ checkThread.shutdown();
}
@Test
diff --git
a/server/src/test/java/org/apache/iotdb/db/wal/node/WALEntryHandlerTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/node/WALEntryHandlerTest.java
index d34ca54cae8..5fe00536ccd 100644
--- a/server/src/test/java/org/apache/iotdb/db/wal/node/WALEntryHandlerTest.java
+++ b/server/src/test/java/org/apache/iotdb/db/wal/node/WALEntryHandlerTest.java
@@ -34,6 +34,7 @@ import org.apache.iotdb.db.wal.checkpoint.CheckpointManager;
import org.apache.iotdb.db.wal.checkpoint.MemTableInfo;
import org.apache.iotdb.db.wal.exception.MemTablePinException;
import org.apache.iotdb.db.wal.utils.WALEntryHandler;
+import org.apache.iotdb.db.wal.utils.WALInsertNodeCache;
import org.apache.iotdb.db.wal.utils.WALMode;
import org.apache.iotdb.db.wal.utils.listener.WALFlushListener;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
@@ -44,46 +45,65 @@ import org.junit.After;
import org.junit.Before;
import org.junit.Test;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.Callable;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
public class WALEntryHandlerTest {
private static final IoTDBConfig config =
IoTDBDescriptor.getInstance().getConfig();
- private static final String identifier = String.valueOf(Integer.MAX_VALUE);
- private static final String logDirectory =
TestConstant.BASE_OUTPUT_PATH.concat("wal-test");
+ private static final String identifier1 = String.valueOf(Integer.MAX_VALUE);
+ private static final String identifier2 = String.valueOf(Integer.MAX_VALUE -
1);
+ private static final String logDirectory1 =
+ TestConstant.BASE_OUTPUT_PATH.concat("wal-test" + identifier1);
+ private static final String logDirectory2 =
+ TestConstant.BASE_OUTPUT_PATH.concat("wal-test" + identifier2);
+
private static final String devicePath = "root.test_sg.test_d";
private WALMode prevMode;
private boolean prevIsClusterMode;
- private WALNode walNode;
+ private WALNode walNode1;
+ private WALNode walNode2;
@Before
public void setUp() throws Exception {
- EnvironmentUtils.cleanDir(logDirectory);
+ EnvironmentUtils.cleanDir(logDirectory1);
+ EnvironmentUtils.cleanDir(logDirectory2);
prevMode = config.getWalMode();
prevIsClusterMode = config.isClusterMode();
config.setWalMode(WALMode.SYNC);
config.setClusterMode(true);
- walNode = new WALNode(identifier, logDirectory);
+ walNode1 = new WALNode(identifier1, logDirectory1);
+ walNode2 = new WALNode(identifier2, logDirectory2);
}
@After
public void tearDown() throws Exception {
- walNode.close();
+ walNode1.close();
+ walNode2.close();
config.setWalMode(prevMode);
config.setClusterMode(prevIsClusterMode);
- EnvironmentUtils.cleanDir(logDirectory);
+ EnvironmentUtils.cleanDir(logDirectory1);
+ EnvironmentUtils.cleanDir(logDirectory2);
+ WALInsertNodeCache.getInstance().clear();
}
@Test(expected = MemTablePinException.class)
public void pinDeletedMemTable() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
WALFlushListener flushListener =
- walNode.log(
+ walNode1.log(
memTable.getMemTableId(), getInsertRowNode(devicePath,
System.currentTimeMillis()));
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// pin flushed memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.pinMemTable();
@@ -92,27 +112,27 @@ public class WALEntryHandlerTest {
@Test
public void pinMemTable() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
InsertRowNode node1 = getInsertRowNode(devicePath,
System.currentTimeMillis());
node1.setSearchIndex(1);
- WALFlushListener flushListener = walNode.log(memTable.getMemTableId(),
node1);
+ WALFlushListener flushListener = walNode1.log(memTable.getMemTableId(),
node1);
// pin memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.pinMemTable();
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// roll wal file
- walNode.rollWALFile();
- walNode.rollWALFile();
+ walNode1.rollWALFile();
+ walNode1.rollWALFile();
// find node1
- ConsensusReqReader.ReqIterator itr = walNode.getReqIterator(1);
+ ConsensusReqReader.ReqIterator itr = walNode1.getReqIterator(1);
assertTrue(itr.hasNext());
assertEquals(
node1,
WALEntry.deserializeForConsensus(itr.next().getRequests().get(0).serializeToByteBuffer()));
// try to delete flushed but pinned memTable
- walNode.deleteOutdatedFiles();
+ walNode1.deleteOutdatedFiles();
// try to find node1
- itr = walNode.getReqIterator(1);
+ itr = walNode1.getReqIterator(1);
assertTrue(itr.hasNext());
assertEquals(
node1,
@@ -122,11 +142,11 @@ public class WALEntryHandlerTest {
@Test(expected = MemTablePinException.class)
public void unpinDeletedMemTable() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
WALFlushListener flushListener =
- walNode.log(
+ walNode1.log(
memTable.getMemTableId(), getInsertRowNode(devicePath,
System.currentTimeMillis()));
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// pin flushed memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.unpinMemTable();
@@ -135,17 +155,17 @@ public class WALEntryHandlerTest {
@Test
public void unpinFlushedMemTable() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
WALFlushListener flushListener =
- walNode.log(
+ walNode1.log(
memTable.getMemTableId(), getInsertRowNode(devicePath,
System.currentTimeMillis()));
WALEntryHandler handler = flushListener.getWalEntryHandler();
// pin twice
handler.pinMemTable();
handler.pinMemTable();
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// unpin 1
- CheckpointManager checkpointManager = walNode.getCheckpointManager();
+ CheckpointManager checkpointManager = walNode1.getCheckpointManager();
handler.unpinMemTable();
MemTableInfo oldestMemTableInfo =
checkpointManager.getOldestMemTableInfo();
assertEquals(memTable.getMemTableId(), oldestMemTableInfo.getMemTableId());
@@ -159,19 +179,19 @@ public class WALEntryHandlerTest {
@Test
public void unpinMemTable() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
InsertRowNode node1 = getInsertRowNode(devicePath,
System.currentTimeMillis());
node1.setSearchIndex(1);
- WALFlushListener flushListener = walNode.log(memTable.getMemTableId(),
node1);
+ WALFlushListener flushListener = walNode1.log(memTable.getMemTableId(),
node1);
// pin memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.pinMemTable();
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// roll wal file
- walNode.rollWALFile();
- walNode.rollWALFile();
+ walNode1.rollWALFile();
+ walNode1.rollWALFile();
// find node1
- ConsensusReqReader.ReqIterator itr = walNode.getReqIterator(1);
+ ConsensusReqReader.ReqIterator itr = walNode1.getReqIterator(1);
assertTrue(itr.hasNext());
assertEquals(
node1,
@@ -179,44 +199,100 @@ public class WALEntryHandlerTest {
// unpin flushed memTable
handler.unpinMemTable();
// try to delete flushed but pinned memTable
- walNode.deleteOutdatedFiles();
+ walNode1.deleteOutdatedFiles();
// try to find node1
- itr = walNode.getReqIterator(1);
+ itr = walNode1.getReqIterator(1);
assertFalse(itr.hasNext());
}
@Test
public void getUnFlushedValue() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
InsertRowNode node1 = getInsertRowNode(devicePath,
System.currentTimeMillis());
node1.setSearchIndex(1);
- WALFlushListener flushListener = walNode.log(memTable.getMemTableId(),
node1);
+ WALFlushListener flushListener = walNode1.log(memTable.getMemTableId(),
node1);
// pin memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.pinMemTable();
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
assertEquals(node1, handler.getValue());
}
@Test
public void getFlushedValue() throws Exception {
IMemTable memTable = new PrimitiveMemTable();
- walNode.onMemTableCreated(memTable, logDirectory + "/" + "fake.tsfile");
+ walNode1.onMemTableCreated(memTable, logDirectory1 + "/" + "fake.tsfile");
InsertRowNode node1 = getInsertRowNode(devicePath,
System.currentTimeMillis());
node1.setSearchIndex(1);
- WALFlushListener flushListener = walNode.log(memTable.getMemTableId(),
node1);
+ WALFlushListener flushListener = walNode1.log(memTable.getMemTableId(),
node1);
// pin memTable
WALEntryHandler handler = flushListener.getWalEntryHandler();
handler.pinMemTable();
- walNode.onMemTableFlushed(memTable);
+ walNode1.onMemTableFlushed(memTable);
// wait until wal flushed
- while (!walNode.isAllWALEntriesConsumed()) {
+ while (!walNode1.isAllWALEntriesConsumed()) {
Thread.sleep(50);
}
assertEquals(node1, handler.getValue());
}
+ @Test
+ public void testConcurrentGetValue() throws Exception {
+ int threadsNum = 10;
+ ExecutorService executorService = Executors.newFixedThreadPool(threadsNum);
+ List<Future<Void>> futures = new ArrayList<>();
+ for (int i = 0; i < threadsNum; ++i) {
+ WALNode walNode = i % 2 == 0 ? walNode1 : walNode2;
+ String logDirectory = i % 2 == 0 ? logDirectory1 : logDirectory2;
+ Callable<Void> writeTask =
+ () -> {
+ IMemTable memTable = new PrimitiveMemTable();
+ walNode.onMemTableCreated(memTable, logDirectory + "/" +
"fake.tsfile");
+
+ List<WALFlushListener> walFlushListeners = new ArrayList<>();
+ List<InsertRowNode> expectedInsertRowNodes = new ArrayList<>();
+ try {
+ for (int j = 0; j < 1_000; ++j) {
+ long memTableId = memTable.getMemTableId();
+ InsertRowNode node =
+ getInsertRowNode(devicePath + memTableId,
System.currentTimeMillis());
+ expectedInsertRowNodes.add(node);
+ WALFlushListener walFlushListener = walNode.log(memTableId,
node);
+ walFlushListeners.add(walFlushListener);
+ }
+ } catch (IllegalPathException e) {
+ fail();
+ }
+
+ // wait until wal flushed
+ while (!walNode1.isAllWALEntriesConsumed() &&
!walNode2.isAllWALEntriesConsumed()) {
+ Thread.sleep(50);
+ }
+
+ walFlushListeners.get(0).getWalEntryHandler().pinMemTable();
+ walNode.onMemTableFlushed(memTable);
+
+ for (int j = 0; j < expectedInsertRowNodes.size(); ++j) {
+ InsertRowNode expect = expectedInsertRowNodes.get(j);
+ InsertRowNode actual =
+ (InsertRowNode)
walFlushListeners.get(j).getWalEntryHandler().getValue();
+ assertEquals(expect, actual);
+ }
+
+ walFlushListeners.get(0).getWalEntryHandler().unpinMemTable();
+ return null;
+ };
+ Future<Void> future = executorService.submit(writeTask);
+ futures.add(future);
+ }
+ // wait until all write tasks are done
+ for (Future<Void> future : futures) {
+ future.get();
+ }
+ executorService.shutdown();
+ }
+
private InsertRowNode getInsertRowNode(String devicePath, long time) throws
IllegalPathException {
TSDataType[] dataTypes =
new TSDataType[] {
diff --git a/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java
index d2820fbde69..59d2bb1b86f 100644
--- a/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java
+++ b/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java
@@ -117,6 +117,7 @@ public class WALNodeTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// wait a moment
while (!walNode.isAllWALEntriesConsumed()) {
Thread.sleep(1_000);
@@ -247,6 +248,7 @@ public class WALNodeTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// recover info from checkpoint file
Map<Long, MemTableInfo> actualMemTableId2Info =
CheckpointRecoverUtils.recoverMemTableInfo(new
File(logDirectory)).getMemTableId2Info();
diff --git
a/server/src/test/java/org/apache/iotdb/db/wal/recover/WALRecoverManagerTest.java
b/server/src/test/java/org/apache/iotdb/db/wal/recover/WALRecoverManagerTest.java
index 38f35bcd3f4..5f97aa9211e 100644
---
a/server/src/test/java/org/apache/iotdb/db/wal/recover/WALRecoverManagerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/wal/recover/WALRecoverManagerTest.java
@@ -177,6 +177,7 @@ public class WALRecoverManagerTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// wait a moment
while (!walBuffer.isAllWALEntriesConsumed()) {
Thread.sleep(1_000);
@@ -235,6 +236,7 @@ public class WALRecoverManagerTest {
for (Future<Void> future : futures) {
future.get();
}
+ executorService.shutdown();
// wait a moment
while (!walBuffer.isAllWALEntriesConsumed()) {
Thread.sleep(1_000);