Skip to content

Commit e163f3b

Browse files
authored
Stabilize subscription heartbeat executor test (#18782)
1 parent b777d3c commit e163f3b

1 file changed

Lines changed: 8 additions & 18 deletions

File tree

‎iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java‎

Lines changed: 8 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -109,33 +109,34 @@ public void testBlockedProviderDoesNotDelayOtherHeartbeatsOrProviderReads() thro
109109
@Test
110110
public void testHeartbeatExecutorIsBounded() throws Exception {
111111
final ThreadPoolExecutor heartbeatExecutor = getHeartbeatExecutor();
112-
waitUntilHeartbeatExecutorIdle(heartbeatExecutor);
113112
final int maximumPoolSize = heartbeatExecutor.getMaximumPoolSize();
114-
final int queueCapacity =
115-
heartbeatExecutor.getQueue().size() + heartbeatExecutor.getQueue().remainingCapacity();
116113
Assert.assertTrue(maximumPoolSize >= 4);
117114
Assert.assertTrue(maximumPoolSize <= 16);
118-
Assert.assertEquals(maximumPoolSize, queueCapacity);
115+
Assert.assertTrue(heartbeatExecutor.getQueue().remainingCapacity() <= maximumPoolSize);
119116

120117
final CountDownLatch releaseTasks = new CountDownLatch(1);
121118
final List<Future<?>> futures = new ArrayList<>();
122119
try {
123120
// The executor is shared, so other heartbeat tasks may consume capacity concurrently.
124-
for (int i = 0; i <= maximumPoolSize + queueCapacity; i++) {
121+
boolean rejected = false;
122+
for (int i = 0; i <= maximumPoolSize * 2; i++) {
125123
final Future<?> future =
126124
SubscriptionExecutorServiceManager.submitProviderHeartbeat(() -> await(releaseTasks));
127125
if (Objects.isNull(future)) {
126+
rejected = true;
128127
break;
129128
}
130129
futures.add(future);
131130
}
132131

133-
Assert.assertTrue(futures.size() <= maximumPoolSize + queueCapacity);
132+
Assert.assertTrue(
133+
"Heartbeat executor accepted more tasks than its configured capacity", rejected);
134134
} finally {
135135
releaseTasks.countDown();
136136
for (final Future<?> future : futures) {
137-
future.get(5, TimeUnit.SECONDS);
137+
future.cancel(true);
138138
}
139+
heartbeatExecutor.purge();
139140
}
140141
}
141142

@@ -158,17 +159,6 @@ private ThreadPoolExecutor getHeartbeatExecutor() throws Exception {
158159
return (ThreadPoolExecutor) executorField.get(holder);
159160
}
160161

161-
private void waitUntilHeartbeatExecutorIdle(final ThreadPoolExecutor executor)
162-
throws InterruptedException {
163-
final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
164-
while ((executor.getActiveCount() != 0 || !executor.getQueue().isEmpty())
165-
&& System.nanoTime() < deadline) {
166-
Thread.sleep(10L);
167-
}
168-
Assert.assertEquals(0, executor.getActiveCount());
169-
Assert.assertTrue(executor.getQueue().isEmpty());
170-
}
171-
172162
private AbstractSubscriptionProviders getProviders(final AbstractSubscriptionConsumer consumer)
173163
throws Exception {
174164
final Field field = AbstractSubscriptionConsumer.class.getDeclaredField("providers");

0 commit comments

Comments
 (0)