From 563ac043adfad816a6cc50b03283e85cfa33a6db Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Mon, 20 Jul 2026 10:41:35 +0800 Subject: [PATCH] Fix flaky concurrent queue and WAL tests --- .../wal/node/WALNodeWaitForRollFileTest.java | 141 +++++------ .../ConcurrentIterableLinkedQueueTest.java | 238 +++++++++++------- 2 files changed, 203 insertions(+), 176 deletions(-) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNodeWaitForRollFileTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNodeWaitForRollFileTest.java index b24a8cd29cf50..1720a76cfbfd7 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNodeWaitForRollFileTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/node/WALNodeWaitForRollFileTest.java @@ -47,13 +47,12 @@ import java.io.File; import java.util.Collections; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@ -205,7 +204,7 @@ public void testLegacySeparatorStillWorksAfterRollFile() throws Exception { * Verifies that waitForNextReady wakes up when a WAL file roll is triggered concurrently. A * background thread rolls the WAL file while the main thread waits on the iterator. */ - @Test(timeout = 30000) + @Test public void testWaitForNextReadyWakesUpOnConcurrentRoll() throws Exception { IMemTable memTable = new PrimitiveMemTable(databasePath, dataRegionId); walNode.onMemTableCreated(memTable, logDirectory + File.separator + "test.tsfile"); @@ -223,38 +222,32 @@ public void testWaitForNextReadyWakesUpOnConcurrentRoll() throws Exception { ConsensusReqReader.ReqIterator iterator = walNode.getReqIterator(1); - AtomicBoolean found = new AtomicBoolean(false); - AtomicReference error = new AtomicReference<>(); ExecutorService executor = Executors.newSingleThreadExecutor(); - - // background: wait for data to become available via waitForNextReady - Future waitFuture = - executor.submit( - () -> { - try { + CountDownLatch waiterStarted = new CountDownLatch(1); + try { + // background: wait for data to become available via waitForNextReady + Future waitFuture = + executor.submit( + () -> { + waiterStarted.countDown(); iterator.waitForNextReady(15, TimeUnit.SECONDS); - if (iterator.hasNext()) { - found.set(true); - } - } catch (Exception e) { - error.set(e); - } - }); - - // give the waiter thread time to start blocking - Thread.sleep(500); - - // trigger WAL file roll — this should signal rollLogWriterCondition and wake up the iterator - walNode.rollWALFile(); - Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> walNode.isAllWALEntriesConsumed()); + return iterator.hasNext(); + }); - waitFuture.get(20, TimeUnit.SECONDS); - executor.shutdown(); + assertTrue(waiterStarted.await(10, TimeUnit.SECONDS)); - if (error.get() != null) { - throw error.get(); + // trigger WAL file roll — this should signal rollLogWriterCondition and wake up the iterator + walNode.rollWALFile(); + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> walNode.isAllWALEntriesConsumed()); + + assertTrue( + "Iterator should have found data after WAL file roll", + waitFuture.get(20, TimeUnit.SECONDS)); + } finally { + shutdownExecutor(executor); } - assertTrue("Iterator should have found data after WAL file roll", found.get()); } /** @@ -262,10 +255,12 @@ public void testWaitForNextReadyWakesUpOnConcurrentRoll() throws Exception { * trigger an automatic WAL file roll (file size exceeds threshold). Uses a small WAL file size * threshold to trigger the roll quickly. */ - @Test(timeout = 60000) + @Test public void testWaitForNextReadyWithAutoRollOnSizeThreshold() throws Exception { // use small WAL file size to trigger auto-roll config.setWalFileSizeThresholdInByte(1024); + ExecutorService executor = Executors.newSingleThreadExecutor(); + CountDownLatch waiterStarted = new CountDownLatch(1); try { IMemTable memTable = new PrimitiveMemTable(databasePath, dataRegionId); @@ -285,24 +280,15 @@ public void testWaitForNextReadyWithAutoRollOnSizeThreshold() throws Exception { ConsensusReqReader.ReqIterator iterator = walNode.getReqIterator(1); - AtomicBoolean found = new AtomicBoolean(false); - AtomicReference error = new AtomicReference<>(); - ExecutorService executor = Executors.newSingleThreadExecutor(); - - Future waitFuture = + Future waitFuture = executor.submit( () -> { - try { - iterator.waitForNextReady(30, TimeUnit.SECONDS); - if (iterator.hasNext()) { - found.set(true); - } - } catch (Exception e) { - error.set(e); - } + waiterStarted.countDown(); + iterator.waitForNextReady(30, TimeUnit.SECONDS); + return iterator.hasNext(); }); - Thread.sleep(500); + assertTrue(waiterStarted.await(10, TimeUnit.SECONDS)); // write more data to exceed the small threshold and trigger auto-roll for (int i = 2; i <= 50; i++) { @@ -314,15 +300,15 @@ public void testWaitForNextReadyWithAutoRollOnSizeThreshold() throws Exception { Collections.singletonList(new int[] {0, node.getRowCount()})); } - waitFuture.get(40, TimeUnit.SECONDS); - executor.shutdown(); - - if (error.get() != null) { - fail("waitForNextReady threw unexpected exception: " + error.get().getMessage()); - } - assertTrue("Iterator should have found data after auto WAL file roll", found.get()); + assertTrue( + "Iterator should have found data after auto WAL file roll", + waitFuture.get(40, TimeUnit.SECONDS)); } finally { - config.setWalFileSizeThresholdInByte(2 * 1024 * 1024); + try { + shutdownExecutor(executor); + } finally { + config.setWalFileSizeThresholdInByte(2 * 1024 * 1024); + } } } @@ -332,7 +318,7 @@ public void testWaitForNextReadyWithAutoRollOnSizeThreshold() throws Exception { * buffer → waitForRollFile(30s) times out → rollWALFile() called → data moves to closed file → * hasNext() returns true → method returns. */ - @Test(timeout = 120000) + @Test public void testWaitForNextReadyAutoTriggersRollOnTimeout() throws Exception { IMemTable memTable = new PrimitiveMemTable(databasePath, dataRegionId); walNode.onMemTableCreated(memTable, logDirectory + File.separator + "test.tsfile"); @@ -352,8 +338,6 @@ public void testWaitForNextReadyAutoTriggersRollOnTimeout() throws Exception { ConsensusReqReader.ReqIterator iterator = walNode.getReqIterator(1); assertFalse("Data should not be visible before WAL file roll", iterator.hasNext()); - AtomicBoolean found = new AtomicBoolean(false); - AtomicReference error = new AtomicReference<>(); ExecutorService executor = Executors.newSingleThreadExecutor(); long startTime = System.currentTimeMillis(); @@ -362,33 +346,32 @@ public void testWaitForNextReadyAutoTriggersRollOnTimeout() throws Exception { // 1) wait 30s for rollLogWriterCondition (timeout) // 2) auto-call rollWALFile() // 3) data becomes readable, hasNext() returns true, method returns - Future waitFuture = - executor.submit( - () -> { - try { + try { + Future waitFuture = + executor.submit( + () -> { iterator.waitForNextReady(); - if (iterator.hasNext()) { - found.set(true); - } - } catch (Exception e) { - error.set(e); - } - }); - - waitFuture.get(90, TimeUnit.SECONDS); - executor.shutdown(); + return iterator.hasNext(); + }); - long elapsed = System.currentTimeMillis() - startTime; + assertTrue( + "Iterator should have found data after auto-triggered WAL file roll", + waitFuture.get(90, TimeUnit.SECONDS)); - if (error.get() != null) { - fail("waitForNextReady() threw unexpected exception: " + error.get().getMessage()); + long elapsed = System.currentTimeMillis() - startTime; + assertTrue( + "Should have waited at least 30s for the timeout to trigger auto-roll, but only waited " + + elapsed + + "ms", + elapsed >= TimeUnit.SECONDS.toMillis(WALNode.WAIT_FOR_NEXT_WAL_ENTRY_TIMEOUT_IN_SEC - 1)); + } finally { + shutdownExecutor(executor); } - assertTrue("Iterator should have found data after auto-triggered WAL file roll", found.get()); - assertTrue( - "Should have waited at least 30s for the timeout to trigger auto-roll, but only waited " - + elapsed - + "ms", - elapsed >= TimeUnit.SECONDS.toMillis(WALNode.WAIT_FOR_NEXT_WAL_ENTRY_TIMEOUT_IN_SEC - 1)); + } + + private static void shutdownExecutor(final ExecutorService executor) throws InterruptedException { + executor.shutdownNow(); + assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); } private InsertTabletNode getInsertTabletNode(String devicePath, long[] times) diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/ConcurrentIterableLinkedQueueTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/ConcurrentIterableLinkedQueueTest.java index f2b0898b418de..20d27df536eb8 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/ConcurrentIterableLinkedQueueTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/datastructure/ConcurrentIterableLinkedQueueTest.java @@ -30,8 +30,10 @@ import java.util.ArrayList; import java.util.HashSet; import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -79,28 +81,34 @@ public void testInsertNull() { } @Test(timeout = 60000) - public void testConcurrentAddAndRemove() throws InterruptedException { + public void testConcurrentAddAndRemove() throws Exception { final int numberOfAdds = 500; final ExecutorService executor = Executors.newFixedThreadPool(2); - // Thread 1 adds elements to the queue - executor.submit( - () -> { - for (int i = 0; i < numberOfAdds; i++) { - queue.add(i); - } - }); - - // Thread 2 removes elements before a certain index - executor.submit( - () -> { - for (int i = 0; i < numberOfAdds; i += 10) { - queue.tryRemoveBefore(i); - } - }); - - executor.shutdown(); - Assert.assertTrue(executor.awaitTermination(1, TimeUnit.MINUTES)); + try { + // Thread 1 adds elements to the queue + final Future addFuture = + executor.submit( + () -> { + for (int i = 0; i < numberOfAdds; i++) { + queue.add(i); + } + }); + + // Thread 2 removes elements before a certain index + final Future removeFuture = + executor.submit( + () -> { + for (int i = 0; i < numberOfAdds; i += 10) { + queue.tryRemoveBefore(i); + } + }); + + addFuture.get(10, TimeUnit.SECONDS); + removeFuture.get(10, TimeUnit.SECONDS); + } finally { + shutdownExecutor(executor); + } // Validate the state of the queue // The actual state depends on the timing of add and remove operations @@ -109,34 +117,48 @@ public void testConcurrentAddAndRemove() throws InterruptedException { } @Test(timeout = 60000) - public void testIterateFromEmptyQueue() { + public void testIterateFromEmptyQueue() throws Exception { final ConcurrentIterableLinkedQueue.DynamicIterator itr = queue.iterateFrom(1); - - final AtomicInteger value = new AtomicInteger(-1); - new Thread(() -> value.set(itr.next())).start(); - queue.add(3); - Awaitility.await().untilAsserted(() -> Assert.assertEquals(3, value.get())); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + final Future valueFuture = + executor.submit( + () -> { + return itr.next(); + }); + queue.add(3); + Assert.assertEquals(Integer.valueOf(3), valueFuture.get(10, TimeUnit.SECONDS)); + } finally { + shutdownExecutor(executor); + } } @Test(timeout = 60000) - public void testContinuousEmptyNext() throws InterruptedException { + public void testContinuousEmptyNext() throws Exception { final ConcurrentIterableLinkedQueue.DynamicIterator itr = queue.iterateFrom(0); - final AtomicInteger consumedValue = new AtomicInteger(0); - new Thread( - () -> { - while (true) { - Integer value = itr.next(0); - if (value != null) { - consumedValue.set(value); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + final CountDownLatch emptyRead = new CountDownLatch(1); + try { + final Future consumedValueFuture = + executor.submit( + () -> { + while (!Thread.currentThread().isInterrupted()) { + final Integer value = itr.next(0); + if (value != null) { + return value; + } + emptyRead.countDown(); } - } - }) - .start(); - Thread.sleep(6000); - queue.add(1); - Awaitility.await() - .atMost(100, TimeUnit.SECONDS) - .untilAsserted(() -> Assert.assertEquals(1, consumedValue.get())); + return null; + }); + + Assert.assertTrue(emptyRead.await(10, TimeUnit.SECONDS)); + queue.add(1); + Assert.assertEquals(Integer.valueOf(1), consumedValueFuture.get(10, TimeUnit.SECONDS)); + } finally { + itr.close(); + shutdownExecutor(executor); + } } @Test(timeout = 60000) @@ -180,7 +202,7 @@ public void testClear() { } @Test(timeout = 60000) - public void testIntegratedOperations() { + public void testIntegratedOperations() throws Exception { queue.add(1); queue.add(2); Assert.assertEquals(1, queue.tryRemoveBefore(1)); @@ -192,10 +214,18 @@ public void testIntegratedOperations() { final ConcurrentIterableLinkedQueue.DynamicIterator it2 = queue.iterateFromLatest(); Assert.assertEquals(2, it2.getNextIndex()); - final AtomicInteger value = new AtomicInteger(-1); - new Thread(() -> value.set(it2.next())).start(); - queue.add(3); - Awaitility.await().untilAsserted(() -> Assert.assertEquals(3, value.get())); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + final Future valueFuture = + executor.submit( + () -> { + return it2.next(); + }); + queue.add(3); + Assert.assertEquals(Integer.valueOf(3), valueFuture.get(10, TimeUnit.SECONDS)); + } finally { + shutdownExecutor(executor); + } Assert.assertEquals(1, it.seek(Integer.MIN_VALUE)); Assert.assertEquals(2, (int) it.next()); @@ -288,37 +318,37 @@ public void testBoundaryConditions() { } @Test(timeout = 60000) - public void testConcurrentExceptionHandling() throws InterruptedException { + public void testConcurrentExceptionHandling() throws Exception { final ExecutorService executor = Executors.newFixedThreadPool(2); + final CountDownLatch iteratorCreated = new CountDownLatch(1); + try { + final Future clearFuture = + executor.submit( + () -> { + queue.add(1); + if (!iteratorCreated.await(10, TimeUnit.SECONDS)) { + return false; + } + queue.clear(); + return true; + }); + + final Future iteratorFuture = + executor.submit( + () -> { + final ConcurrentIterableLinkedQueue.DynamicIterator itr = + queue.iterateFromEarliest(); + iteratorCreated.countDown(); + while (itr.hasNext()) { + itr.next(); + } + }); - executor.submit( - () -> { - queue.add(1); - try { - Thread.sleep(500); // Wait for the iterator to start - } catch (InterruptedException e) { - throw new RuntimeException(e); - } - queue.clear(); - }); - - final AtomicBoolean caughtException = new AtomicBoolean(false); - executor.submit( - () -> { - try { - final ConcurrentIterableLinkedQueue.DynamicIterator itr = - queue.iterateFromEarliest(); - while (itr.hasNext()) { - itr.next(); - } - } catch (Exception e) { - caughtException.set(true); - } - }); - - executor.shutdown(); - Assert.assertTrue(executor.awaitTermination(1, TimeUnit.MINUTES)); - Assert.assertFalse(caughtException.get()); + Assert.assertTrue(clearFuture.get(10, TimeUnit.SECONDS)); + iteratorFuture.get(10, TimeUnit.SECONDS); + } finally { + shutdownExecutor(executor); + } } @Test(timeout = 60000) @@ -337,18 +367,22 @@ public void testHighLoadPerformance() { } @Test(timeout = 60000) - public void testMultiThreadedConsistency() throws InterruptedException { + public void testMultiThreadedConsistency() throws Exception { final int numberOfElements = 1000; final ExecutorService executor = Executors.newFixedThreadPool(10); - - for (int i = 0; i < numberOfElements; i++) { - int finalI = i; - executor.submit(() -> queue.add(finalI)); + final List> addFutures = new ArrayList<>(numberOfElements); + try { + for (int i = 0; i < numberOfElements; i++) { + int finalI = i; + addFutures.add(executor.submit(() -> queue.add(finalI))); + } + for (final Future addFuture : addFutures) { + addFuture.get(10, TimeUnit.SECONDS); + } + } finally { + shutdownExecutor(executor); } - executor.shutdown(); - Assert.assertTrue(executor.awaitTermination(1, TimeUnit.MINUTES)); - final ConcurrentIterableLinkedQueue.DynamicIterator itr = queue.iterateFromEarliest(); final HashSet elements = new HashSet<>(); while (itr.hasNext()) { @@ -388,27 +422,32 @@ public void testIteratorAfterRemoval() { } @Test(timeout = 60000) - public void testIteratorConcurrentAccess() throws InterruptedException { + public void testIteratorConcurrentAccess() throws Exception { for (int i = 0; i < 100; i++) { queue.add(i); } final ExecutorService executor = Executors.newFixedThreadPool(10); final AtomicInteger count = new AtomicInteger(0); - - for (int i = 0; i < 10; i++) { - executor.submit( - () -> { - ConcurrentIterableLinkedQueue.DynamicIterator itr = queue.iterateFrom(0); - while (itr.hasNext()) { - itr.next(); - count.incrementAndGet(); - } - }); + final List> iteratorFutures = new ArrayList<>(10); + try { + for (int i = 0; i < 10; i++) { + iteratorFutures.add( + executor.submit( + () -> { + ConcurrentIterableLinkedQueue.DynamicIterator itr = queue.iterateFrom(0); + while (itr.hasNext()) { + itr.next(); + count.incrementAndGet(); + } + })); + } + for (final Future iteratorFuture : iteratorFutures) { + iteratorFuture.get(10, TimeUnit.SECONDS); + } + } finally { + shutdownExecutor(executor); } - - executor.shutdown(); - Assert.assertTrue(executor.awaitTermination(1, TimeUnit.MINUTES)); Assert.assertEquals(1000, count.get()); // 100 elements iterated by 10 threads } @@ -502,4 +541,9 @@ public void testSetFirstIndexReplacesEmptyQueueIndexes() { Assert.assertFalse(itr.hasNext()); } } + + private static void shutdownExecutor(final ExecutorService executor) throws InterruptedException { + executor.shutdownNow(); + Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); + } }