diff --git a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java index 842f7e7daf25..662386f4dc0c 100644 --- a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java +++ b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java @@ -109,33 +109,34 @@ public void testBlockedProviderDoesNotDelayOtherHeartbeatsOrProviderReads() thro @Test public void testHeartbeatExecutorIsBounded() throws Exception { final ThreadPoolExecutor heartbeatExecutor = getHeartbeatExecutor(); - waitUntilHeartbeatExecutorIdle(heartbeatExecutor); final int maximumPoolSize = heartbeatExecutor.getMaximumPoolSize(); - final int queueCapacity = - heartbeatExecutor.getQueue().size() + heartbeatExecutor.getQueue().remainingCapacity(); Assert.assertTrue(maximumPoolSize >= 4); Assert.assertTrue(maximumPoolSize <= 16); - Assert.assertEquals(maximumPoolSize, queueCapacity); + Assert.assertTrue(heartbeatExecutor.getQueue().remainingCapacity() <= maximumPoolSize); final CountDownLatch releaseTasks = new CountDownLatch(1); final List> futures = new ArrayList<>(); try { // The executor is shared, so other heartbeat tasks may consume capacity concurrently. - for (int i = 0; i <= maximumPoolSize + queueCapacity; i++) { + boolean rejected = false; + for (int i = 0; i <= maximumPoolSize * 2; i++) { final Future future = SubscriptionExecutorServiceManager.submitProviderHeartbeat(() -> await(releaseTasks)); if (Objects.isNull(future)) { + rejected = true; break; } futures.add(future); } - Assert.assertTrue(futures.size() <= maximumPoolSize + queueCapacity); + Assert.assertTrue( + "Heartbeat executor accepted more tasks than its configured capacity", rejected); } finally { releaseTasks.countDown(); for (final Future future : futures) { - future.get(5, TimeUnit.SECONDS); + future.cancel(true); } + heartbeatExecutor.purge(); } } @@ -158,17 +159,6 @@ private ThreadPoolExecutor getHeartbeatExecutor() throws Exception { return (ThreadPoolExecutor) executorField.get(holder); } - private void waitUntilHeartbeatExecutorIdle(final ThreadPoolExecutor executor) - throws InterruptedException { - final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); - while ((executor.getActiveCount() != 0 || !executor.getQueue().isEmpty()) - && System.nanoTime() < deadline) { - Thread.sleep(10L); - } - Assert.assertEquals(0, executor.getActiveCount()); - Assert.assertTrue(executor.getQueue().isEmpty()); - } - private AbstractSubscriptionProviders getProviders(final AbstractSubscriptionConsumer consumer) throws Exception { final Field field = AbstractSubscriptionConsumer.class.getDeclaredField("providers");