This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 d6b602e4463 [Subscription] Bound heartbeat provider failures (#18712)
d6b602e4463 is described below
commit d6b602e4463a9c01196e84368b6a647b00bf8856
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 17:10:15 2026 +0800
[Subscription] Bound heartbeat provider failures (#18712)
* [Subscription] Bound heartbeat provider failures
* [Subscription] Bound provider heartbeat executor
---
.../subscription/i18n/SubscriptionMessages.java | 2 +
.../subscription/i18n/SubscriptionMessages.java | 2 +
.../rpc/subscription/config/ConsumerConstant.java | 2 +-
.../base/AbstractSubscriptionConsumer.java | 2 +-
.../base/AbstractSubscriptionConsumerBuilder.java | 2 +-
.../base/AbstractSubscriptionProviders.java | 72 +++--
.../base/SubscriptionExecutorServiceManager.java | 68 ++++
.../tree/SubscriptionTreePullConsumer.java | 2 +-
.../tree/SubscriptionTreePushConsumer.java | 2 +-
...SubscriptionConsumerHeartbeatIsolationTest.java | 350 +++++++++++++++++++++
.../base/SubscriptionConsumerLifecycleTest.java | 10 +-
11 files changed, 478 insertions(+), 36 deletions(-)
diff --git
a/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
b/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
index 3693c7da2ca..41bcc525d18 100644
---
a/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
+++
b/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
@@ -53,6 +53,8 @@ public final class SubscriptionMessages {
// --- SubscriptionExecutorServiceManager ---
public static final String EXECUTOR_LAUNCHING = "Launching {} with core pool
size {}...";
+ public static final String
LOG_SUBSCRIPTION_HEARTBEAT_EXECUTOR_REACHED_ITS_THREAD_OR_QUEUE_LIMIT_SKIP_13157249
=
+ "Subscription heartbeat executor reached its thread or queue limit;
skipping heartbeat tasks until the next interval (thread limit {}, queue
capacity {}).";
public static final String EXECUTOR_SHUTTING_DOWN = "Shutting down {}...";
public static final String EXECUTOR_NOT_LAUNCHED_SUBMIT =
"{} has not been launched, ignore submit task";
diff --git
a/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
b/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
index e3eea001964..89e4d59750a 100644
---
a/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
+++
b/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java
@@ -52,6 +52,8 @@ public final class SubscriptionMessages {
// --- SubscriptionExecutorServiceManager ---
public static final String EXECUTOR_LAUNCHING = "正在启动 {},核心线程池大小:{}...";
+ public static final String
LOG_SUBSCRIPTION_HEARTBEAT_EXECUTOR_REACHED_ITS_THREAD_OR_QUEUE_LIMIT_SKIP_13157249
=
+ "Subscription heartbeat executor 达到线程数或队列容量上限;本轮跳过 heartbeat
任务,下个周期重试(线程上限 {},队列容量 {})。";
public static final String EXECUTOR_SHUTTING_DOWN = "正在关闭 {}...";
public static final String EXECUTOR_NOT_LAUNCHED_SUBMIT =
"{} 尚未启动,忽略提交任务";
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/ConsumerConstant.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/ConsumerConstant.java
index 89878efb49a..6ecb56ba63c 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/ConsumerConstant.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/ConsumerConstant.java
@@ -62,7 +62,7 @@ public class ConsumerConstant {
public static final String THRIFT_MAX_FRAME_SIZE_KEY =
"thrift-max-frame-size";
public static final String CONNECTION_TIMEOUT_MS_KEY =
"connection-timeout-ms";
- public static final int CONNECTION_TIMEOUT_MS_DEFAULT_VALUE = 0;
+ public static final int CONNECTION_TIMEOUT_MS_DEFAULT_VALUE = 10_000;
public static final String MAX_POLL_PARALLELISM_KEY = "max-poll-parallelism";
public static final int MAX_POLL_PARALLELISM_DEFAULT_VALUE = 1;
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
index e2ab3476606..a57c7fcd1d3 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
@@ -296,7 +296,7 @@ abstract class AbstractSubscriptionConsumer implements
AutoCloseable {
(Integer)
properties.getOrDefault(
ConsumerConstant.CONNECTION_TIMEOUT_MS_KEY,
- SessionConfig.DEFAULT_CONNECTION_TIMEOUT_MS))
+ ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE))
.maxPollParallelism(
(Integer)
properties.getOrDefault(
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumerBuilder.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumerBuilder.java
index 07fde6a88a8..a9e9d38f334 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumerBuilder.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumerBuilder.java
@@ -52,7 +52,7 @@ public class AbstractSubscriptionConsumerBuilder {
protected boolean fileSaveFsync =
ConsumerConstant.FILE_SAVE_FSYNC_DEFAULT_VALUE;
protected int thriftMaxFrameSize = SessionConfig.DEFAULT_MAX_FRAME_SIZE;
- protected int connectionTimeoutInMs =
SessionConfig.DEFAULT_CONNECTION_TIMEOUT_MS;
+ protected int connectionTimeoutInMs =
ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE;
protected int maxPollParallelism =
ConsumerConstant.MAX_POLL_PARALLELISM_DEFAULT_VALUE;
public AbstractSubscriptionConsumerBuilder host(final String host) {
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
index 3eaf6e6906e..cae7ded2cf9 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
@@ -38,6 +38,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.SortedMap;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@@ -50,6 +51,9 @@ final class AbstractSubscriptionProviders {
private final SortedMap<Integer, AbstractSubscriptionProvider>
subscriptionProviders =
new ConcurrentSkipListMap<>();
private final AtomicBoolean isClosing = new AtomicBoolean(false);
+ private final Set<AbstractSubscriptionProvider>
providersWithHeartbeatInFlight =
+ ConcurrentHashMap.newKeySet();
+ private final Object heartbeatResponseLock = new Object();
private int nextDataNodeId = -1;
private final ReentrantReadWriteLock subscriptionProvidersLock = new
ReentrantReadWriteLock(true);
@@ -316,29 +320,33 @@ final class AbstractSubscriptionProviders {
return;
}
- acquireWriteLock();
- try {
- if (consumer.isClosed() || consumer.isFenced()) {
- return;
+ for (final AbstractSubscriptionProvider provider : getAllProviders()) {
+ if (!providersWithHeartbeatInFlight.add(provider)) {
+ continue;
+ }
+ if (Objects.isNull(
+ SubscriptionExecutorServiceManager.submitProviderHeartbeat(
+ () -> heartbeatProvider(consumer, provider)))) {
+ providersWithHeartbeatInFlight.remove(provider);
}
- heartbeatInternal(consumer);
- } finally {
- releaseWriteLock();
}
}
- private void heartbeatInternal(final AbstractSubscriptionConsumer consumer) {
- for (final AbstractSubscriptionProvider provider : getAllProviders()) {
- if (consumer.isFenced()) {
+ private void heartbeatProvider(
+ final AbstractSubscriptionConsumer consumer, final
AbstractSubscriptionProvider provider) {
+ try {
+ if (consumer.isClosed() || consumer.isFenced() ||
!containsProvider(provider)) {
return;
}
- try {
- final List<SubscriptionCommitContext> processorBufferedCommitContexts =
-
consumer.getProcessorBufferedCommitContexts(provider.getDataNodeId());
- final PipeSubscribeHeartbeatResp resp =
provider.heartbeat(processorBufferedCommitContexts);
- // update subscribed topics
+ final List<SubscriptionCommitContext> processorBufferedCommitContexts =
+
consumer.getProcessorBufferedCommitContexts(provider.getDataNodeId());
+ final PipeSubscribeHeartbeatResp resp =
provider.heartbeat(processorBufferedCommitContexts);
+ if (consumer.isClosed() || consumer.isFenced() ||
!containsProvider(provider)) {
+ return;
+ }
+ provider.setAvailable();
+ synchronized (heartbeatResponseLock) {
consumer.subscribedTopics = resp.getTopics();
- // unsubscribe completed topics
for (final String topicName : resp.getTopicNamesToUnsubscribe()) {
LOGGER.info(
SubscriptionMessages
@@ -347,24 +355,28 @@ final class AbstractSubscriptionProviders {
topicName);
consumer.unsubscribe(topicName);
}
- provider.setAvailable();
- } catch (final SubscriptionConsumerFencedException e) {
- consumer.fence(e);
- provider.setUnavailable();
- return;
- } catch (final Exception e) {
- LOGGER.warn(
- SubscriptionMessages
-
.LOG_ARG_FAILED_SENDING_HEARTBEAT_SUBSCRIPTION_PROVIDER_ARG_BECAUSE_ARG_SET_0B38FB1F,
- consumer,
- provider,
- e,
- e);
- provider.setUnavailable();
}
+ } catch (final SubscriptionConsumerFencedException e) {
+ consumer.fence(e);
+ provider.setUnavailable();
+ } catch (final Exception e) {
+ LOGGER.warn(
+ SubscriptionMessages
+
.LOG_ARG_FAILED_SENDING_HEARTBEAT_SUBSCRIPTION_PROVIDER_ARG_BECAUSE_ARG_SET_0B38FB1F,
+ consumer,
+ provider,
+ e,
+ e);
+ provider.setUnavailable();
+ } finally {
+ providersWithHeartbeatInFlight.remove(provider);
}
}
+ private boolean containsProvider(final AbstractSubscriptionProvider
provider) {
+ return subscriptionProviders.get(provider.getDataNodeId()) == provider;
+ }
+
/////////////////////////////// sync endpoints
///////////////////////////////
void sync(final AbstractSubscriptionConsumer consumer) {
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionExecutorServiceManager.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionExecutorServiceManager.java
index 4f48f9d6f4d..6987418826c 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionExecutorServiceManager.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionExecutorServiceManager.java
@@ -27,13 +27,18 @@ import org.slf4j.LoggerFactory;
import java.util.Collection;
import java.util.List;
import java.util.Objects;
+import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
public final class SubscriptionExecutorServiceManager {
@@ -41,12 +46,30 @@ public final class SubscriptionExecutorServiceManager {
LoggerFactory.getLogger(SubscriptionExecutorServiceManager.class);
private static final long AWAIT_TERMINATION_TIMEOUT_MS = 15_000L;
+ private static final long HEARTBEAT_EXECUTOR_REJECTION_LOG_INTERVAL_NANOS =
+ TimeUnit.MINUTES.toNanos(1);
private static final String CONTROL_FLOW_EXECUTOR_NAME =
"SubscriptionControlFlowExecutor";
private static final String UPSTREAM_DATA_FLOW_EXECUTOR_NAME =
"SubscriptionUpstreamDataFlowExecutor";
private static final String DOWNSTREAM_DATA_FLOW_EXECUTOR_NAME =
"SubscriptionDownstreamDataFlowExecutor";
+ private static final String HEARTBEAT_EXECUTOR_NAME =
"SubscriptionHeartbeatExecutor";
+ private static final int HEARTBEAT_EXECUTOR_MIN_THREAD_COUNT = 4;
+ private static final int HEARTBEAT_EXECUTOR_MAX_THREAD_COUNT = 16;
+ private static final int HEARTBEAT_EXECUTOR_THREAD_COUNT =
+ Math.min(
+ Math.max(Runtime.getRuntime().availableProcessors(),
HEARTBEAT_EXECUTOR_MIN_THREAD_COUNT),
+ HEARTBEAT_EXECUTOR_MAX_THREAD_COUNT);
+ private static final int HEARTBEAT_EXECUTOR_QUEUE_CAPACITY =
HEARTBEAT_EXECUTOR_THREAD_COUNT;
+ private static final AtomicLong LAST_HEARTBEAT_EXECUTOR_REJECTION_LOG_TIME =
new AtomicLong();
+ private static final ThreadFactory HEARTBEAT_EXECUTOR_THREAD_FACTORY =
+ r -> {
+ final Thread t =
+ new Thread(Thread.currentThread().getThreadGroup(), r,
HEARTBEAT_EXECUTOR_NAME, 0);
+ t.setDaemon(true);
+ return t;
+ };
/** Control Flow Executor: execute heartbeat worker, endpoints syncer and
auto poll worker */
private static final SubscriptionScheduledExecutorService
CONTROL_FLOW_EXECUTOR =
@@ -65,6 +88,31 @@ public final class SubscriptionExecutorServiceManager {
DOWNSTREAM_DATA_FLOW_EXECUTOR_NAME,
Math.max(Runtime.getRuntime().availableProcessors(), 1));
+ /** Heartbeat Executor: isolate a slow provider from the heartbeat
control-flow scheduler. */
+ private static final SubscriptionExecutorService HEARTBEAT_EXECUTOR =
+ new SubscriptionExecutorService(HEARTBEAT_EXECUTOR_NAME,
HEARTBEAT_EXECUTOR_THREAD_COUNT) {
+ @Override
+ void launchIfNeeded() {
+ if (isShutdown()) {
+ synchronized (this) {
+ if (isShutdown()) {
+ LOGGER.info(SubscriptionMessages.EXECUTOR_LAUNCHING,
this.name, this.corePoolSize);
+ this.executor =
+ new ThreadPoolExecutor(
+ HEARTBEAT_EXECUTOR_THREAD_COUNT,
+ HEARTBEAT_EXECUTOR_THREAD_COUNT,
+ 60L,
+ TimeUnit.SECONDS,
+ new
ArrayBlockingQueue<>(HEARTBEAT_EXECUTOR_QUEUE_CAPACITY),
+ HEARTBEAT_EXECUTOR_THREAD_FACTORY,
+ new ThreadPoolExecutor.AbortPolicy());
+ ((ThreadPoolExecutor)
this.executor).allowCoreThreadTimeOut(true);
+ }
+ }
+ }
+ }
+ };
+
/////////////////////////////// set core pool size
///////////////////////////////
public static void setControlFlowExecutorCorePoolSize(final int
corePoolSize) {
@@ -98,6 +146,7 @@ public final class SubscriptionExecutorServiceManager {
CONTROL_FLOW_EXECUTOR.shutdown();
UPSTREAM_DATA_FLOW_EXECUTOR.shutdown();
DOWNSTREAM_DATA_FLOW_EXECUTOR.shutdown();
+ HEARTBEAT_EXECUTOR.shutdown();
}
}
@@ -152,6 +201,25 @@ public final class SubscriptionExecutorServiceManager {
UPSTREAM_DATA_FLOW_EXECUTOR.submit(task);
}
+ static Future<?> submitProviderHeartbeat(final Runnable task) {
+ HEARTBEAT_EXECUTOR.launchIfNeeded();
+ try {
+ return HEARTBEAT_EXECUTOR.submit(task);
+ } catch (final RejectedExecutionException e) {
+ final long now = System.nanoTime();
+ final long previous = LAST_HEARTBEAT_EXECUTOR_REJECTION_LOG_TIME.get();
+ if ((previous == 0 || now - previous >=
HEARTBEAT_EXECUTOR_REJECTION_LOG_INTERVAL_NANOS)
+ &&
LAST_HEARTBEAT_EXECUTOR_REJECTION_LOG_TIME.compareAndSet(previous, now)) {
+ LOGGER.warn(
+ SubscriptionMessages
+
.LOG_SUBSCRIPTION_HEARTBEAT_EXECUTOR_REACHED_ITS_THREAD_OR_QUEUE_LIMIT_SKIP_13157249,
+ HEARTBEAT_EXECUTOR_THREAD_COUNT,
+ HEARTBEAT_EXECUTOR_QUEUE_CAPACITY);
+ }
+ return null;
+ }
+ }
+
public static <T> List<Future<T>> submitMultiplePollTasks(
final Collection<? extends Callable<T>> tasks, final long timeoutMs)
throws InterruptedException {
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePullConsumer.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePullConsumer.java
index 551eb4c3298..3e529d16d89 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePullConsumer.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePullConsumer.java
@@ -282,7 +282,7 @@ public class SubscriptionTreePullConsumer extends
AbstractSubscriptionPullConsum
private boolean fileSaveFsync =
ConsumerConstant.FILE_SAVE_FSYNC_DEFAULT_VALUE;
private int thriftMaxFrameSize = SessionConfig.DEFAULT_MAX_FRAME_SIZE;
- private int connectionTimeoutInMs =
SessionConfig.DEFAULT_CONNECTION_TIMEOUT_MS;
+ private int connectionTimeoutInMs =
ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE;
private int maxPollParallelism =
ConsumerConstant.MAX_POLL_PARALLELISM_DEFAULT_VALUE;
private boolean autoCommit = ConsumerConstant.AUTO_COMMIT_DEFAULT_VALUE;
diff --git
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePushConsumer.java
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePushConsumer.java
index 8fbc33d14ce..895a88467a5 100644
---
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePushConsumer.java
+++
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/tree/SubscriptionTreePushConsumer.java
@@ -209,7 +209,7 @@ public class SubscriptionTreePushConsumer extends
AbstractSubscriptionPushConsum
private boolean fileSaveFsync =
ConsumerConstant.FILE_SAVE_FSYNC_DEFAULT_VALUE;
private int thriftMaxFrameSize = SessionConfig.DEFAULT_MAX_FRAME_SIZE;
- private int connectionTimeoutInMs =
SessionConfig.DEFAULT_CONNECTION_TIMEOUT_MS;
+ private int connectionTimeoutInMs =
ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE;
private int maxPollParallelism =
ConsumerConstant.MAX_POLL_PARALLELISM_DEFAULT_VALUE;
private AckStrategy ackStrategy = AckStrategy.defaultValue();
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
new file mode 100644
index 00000000000..1c7d0af6b66
--- /dev/null
+++
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerHeartbeatIsolationTest.java
@@ -0,0 +1,350 @@
+/*
+ * 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.session.subscription.consumer.base;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.rpc.subscription.config.ConsumerConstant;
+import org.apache.iotdb.rpc.subscription.config.TopicConfig;
+import
org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionCommitContext;
+import
org.apache.iotdb.rpc.subscription.payload.response.PipeSubscribeHeartbeatResp;
+import org.apache.iotdb.session.AbstractSessionBuilder;
+import org.apache.iotdb.session.subscription.SubscriptionTreeSessionBuilder;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.lang.reflect.Field;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class SubscriptionConsumerHeartbeatIsolationTest {
+
+ private static final String HOST = "127.0.0.1";
+ private static final int FIRST_PORT = 6667;
+ private static final long LONG_INTERVAL_MS = 86_400_000L;
+ private static final String TOPIC = "topic";
+
+ @Test
+ public void testSubscriptionConsumerUsesFiniteConnectionTimeoutByDefault() {
+ Assert.assertEquals(
+ ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE,
+ new AbstractSubscriptionPullConsumerBuilder().connectionTimeoutInMs);
+ Assert.assertTrue(ConsumerConstant.CONNECTION_TIMEOUT_MS_DEFAULT_VALUE >
0);
+ }
+
+ @Test
+ public void testBlockedProviderDoesNotDelayOtherHeartbeatsOrProviderReads()
throws Exception {
+ final CountDownLatch blockedHeartbeatStarted = new CountDownLatch(1);
+ final CountDownLatch releaseBlockedHeartbeat = new CountDownLatch(1);
+ final CountDownLatch healthyHeartbeatCompleted = new CountDownLatch(1);
+ final TestPullConsumer consumer =
+ new TestPullConsumer(
+ blockedHeartbeatStarted, releaseBlockedHeartbeat,
healthyHeartbeatCompleted);
+ final ExecutorService executor = Executors.newSingleThreadExecutor();
+
+ try {
+ consumer.open();
+ consumer.blockProviderOne.set(true);
+
+ final AbstractSubscriptionProviders providers = getProviders(consumer);
+ final int blockedProviderInitialHeartbeatCount =
consumer.getHeartbeatCount(1);
+ final int healthyProviderInitialHeartbeatCount =
consumer.getHeartbeatCount(2);
+ final long startNanos = System.nanoTime();
+ final Future<?> heartbeat = executor.submit(() ->
providers.heartbeat(consumer));
+ heartbeat.get(1, TimeUnit.SECONDS);
+ final long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() -
startNanos);
+
+ Assert.assertTrue(
+ "Heartbeat scheduling was blocked for " + elapsedMs + " ms",
elapsedMs < 500);
+ Assert.assertTrue(blockedHeartbeatStarted.await(5, TimeUnit.SECONDS));
+ Assert.assertTrue(healthyHeartbeatCompleted.await(5, TimeUnit.SECONDS));
+
+ providers.acquireReadLock();
+ try {
+ Assert.assertEquals(2, providers.getAllProviders().size());
+ } finally {
+ providers.releaseReadLock();
+ }
+
+ providers.heartbeat(consumer);
+ Thread.sleep(100L);
+ Assert.assertEquals(blockedProviderInitialHeartbeatCount + 1,
consumer.getHeartbeatCount(1));
+ Assert.assertTrue(consumer.getHeartbeatCount(2) >=
healthyProviderInitialHeartbeatCount + 2);
+ } finally {
+ releaseBlockedHeartbeat.countDown();
+ consumer.close();
+ executor.shutdownNow();
+ executor.awaitTermination(5, TimeUnit.SECONDS);
+ }
+ }
+
+ @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);
+
+ final CountDownLatch releaseTasks = new CountDownLatch(1);
+ final List<Future<?>> futures = new ArrayList<>();
+ try {
+ for (int i = 0; i < maximumPoolSize; i++) {
+ futures.add(
+ SubscriptionExecutorServiceManager.submitProviderHeartbeat(() ->
await(releaseTasks)));
+ }
+ for (int i = 0; i < queueCapacity; i++) {
+ futures.add(
+ SubscriptionExecutorServiceManager.submitProviderHeartbeat(() ->
await(releaseTasks)));
+ }
+
+
Assert.assertNull(SubscriptionExecutorServiceManager.submitProviderHeartbeat(()
-> {}));
+ } finally {
+ releaseTasks.countDown();
+ for (final Future<?> future : futures) {
+ future.get(5, TimeUnit.SECONDS);
+ }
+ }
+ }
+
+ private void await(final CountDownLatch latch) {
+ try {
+ latch.await();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ private ThreadPoolExecutor getHeartbeatExecutor() throws Exception {
+ SubscriptionExecutorServiceManager.submitProviderHeartbeat(() -> {});
+ final Field holderField =
+
SubscriptionExecutorServiceManager.class.getDeclaredField("HEARTBEAT_EXECUTOR");
+ holderField.setAccessible(true);
+ final Object holder = holderField.get(null);
+ final Field executorField =
holder.getClass().getSuperclass().getDeclaredField("executor");
+ executorField.setAccessible(true);
+ 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");
+ field.setAccessible(true);
+ return (AbstractSubscriptionProviders) field.get(consumer);
+ }
+
+ private static class TestPullConsumer extends
AbstractSubscriptionPullConsumer {
+
+ private final CountDownLatch blockedHeartbeatStarted;
+ private final CountDownLatch releaseBlockedHeartbeat;
+ private final CountDownLatch healthyHeartbeatCompleted;
+ private final AtomicBoolean blockProviderOne = new AtomicBoolean(false);
+ private final Map<Integer, AtomicInteger> heartbeatCounts = new
ConcurrentHashMap<>();
+
+ private TestPullConsumer(
+ final CountDownLatch blockedHeartbeatStarted,
+ final CountDownLatch releaseBlockedHeartbeat,
+ final CountDownLatch healthyHeartbeatCompleted) {
+ super(
+ new AbstractSubscriptionPullConsumerBuilder()
+ .host(HOST)
+ .port(FIRST_PORT)
+ .consumerId("consumer")
+ .consumerGroupId("group")
+ .heartbeatIntervalMs(LONG_INTERVAL_MS)
+ .endpointsSyncIntervalMs(LONG_INTERVAL_MS)
+ .autoCommit(false));
+ this.blockedHeartbeatStarted = blockedHeartbeatStarted;
+ this.releaseBlockedHeartbeat = releaseBlockedHeartbeat;
+ this.healthyHeartbeatCompleted = healthyHeartbeatCompleted;
+ }
+
+ @Override
+ protected AbstractSubscriptionProvider constructSubscriptionProvider(
+ final TEndPoint endPoint,
+ final String username,
+ final String password,
+ final String encryptedPassword,
+ final String consumerId,
+ final String consumerGroupId,
+ final String ownerId,
+ final Long ownerEpoch,
+ final int thriftMaxFrameSize,
+ final long heartbeatIntervalMs,
+ final int connectionTimeoutInMs) {
+ return new TestSubscriptionProvider(
+ endPoint,
+ username,
+ password,
+ encryptedPassword,
+ consumerId,
+ consumerGroupId,
+ ownerId,
+ ownerEpoch,
+ thriftMaxFrameSize,
+ heartbeatIntervalMs,
+ connectionTimeoutInMs,
+ blockedHeartbeatStarted,
+ releaseBlockedHeartbeat,
+ healthyHeartbeatCompleted,
+ blockProviderOne,
+ heartbeatCounts);
+ }
+
+ private int getHeartbeatCount(final int dataNodeId) {
+ return heartbeatCounts.getOrDefault(dataNodeId, new
AtomicInteger()).get();
+ }
+ }
+
+ private static class TestSubscriptionProvider extends
AbstractSubscriptionProvider {
+
+ private final int dataNodeId;
+ private final CountDownLatch blockedHeartbeatStarted;
+ private final CountDownLatch releaseBlockedHeartbeat;
+ private final CountDownLatch healthyHeartbeatCompleted;
+ private final AtomicBoolean blockProviderOne;
+ private final Map<Integer, AtomicInteger> heartbeatCounts;
+
+ private TestSubscriptionProvider(
+ final TEndPoint endPoint,
+ final String username,
+ final String password,
+ final String encryptedPassword,
+ final String consumerId,
+ final String consumerGroupId,
+ final String ownerId,
+ final Long ownerEpoch,
+ final int thriftMaxFrameSize,
+ final long heartbeatIntervalMs,
+ final int connectionTimeoutInMs,
+ final CountDownLatch blockedHeartbeatStarted,
+ final CountDownLatch releaseBlockedHeartbeat,
+ final CountDownLatch healthyHeartbeatCompleted,
+ final AtomicBoolean blockProviderOne,
+ final Map<Integer, AtomicInteger> heartbeatCounts) {
+ super(
+ endPoint,
+ username,
+ password,
+ encryptedPassword,
+ consumerId,
+ consumerGroupId,
+ ownerId,
+ ownerEpoch,
+ thriftMaxFrameSize,
+ heartbeatIntervalMs,
+ connectionTimeoutInMs);
+ this.dataNodeId = endPoint.port - FIRST_PORT + 1;
+ this.blockedHeartbeatStarted = blockedHeartbeatStarted;
+ this.releaseBlockedHeartbeat = releaseBlockedHeartbeat;
+ this.healthyHeartbeatCompleted = healthyHeartbeatCompleted;
+ this.blockProviderOne = blockProviderOne;
+ this.heartbeatCounts = heartbeatCounts;
+ }
+
+ @Override
+ protected AbstractSessionBuilder constructSubscriptionSessionBuilder(
+ final String host,
+ final int port,
+ final String username,
+ final String password,
+ final String encryptedPassword,
+ final int thriftMaxFrameSize,
+ final int connectionTimeoutInMs) {
+ final boolean useEncryptedPassword = Objects.nonNull(encryptedPassword);
+ return new SubscriptionTreeSessionBuilder()
+ .host(host)
+ .port(port)
+ .username(username)
+ .password(useEncryptedPassword ? encryptedPassword : password)
+ .useEncryptedPassword(useEncryptedPassword)
+ .thriftMaxFrameSize(thriftMaxFrameSize)
+ .connectionTimeoutInMs(connectionTimeoutInMs);
+ }
+
+ @Override
+ synchronized void handshake() {
+ setAvailable();
+ }
+
+ @Override
+ synchronized void close() {
+ setUnavailable();
+ }
+
+ @Override
+ synchronized void closeSession() {
+ setUnavailable();
+ }
+
+ @Override
+ int getDataNodeId() {
+ return dataNodeId;
+ }
+
+ @Override
+ PipeSubscribeHeartbeatResp heartbeat(
+ final List<SubscriptionCommitContext> processorBufferedCommitContexts)
{
+ heartbeatCounts.computeIfAbsent(dataNodeId, ignored -> new
AtomicInteger()).incrementAndGet();
+ if (dataNodeId == 1 && blockProviderOne.get()) {
+ blockedHeartbeatStarted.countDown();
+ try {
+ releaseBlockedHeartbeat.await();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ if (dataNodeId == 2) {
+ healthyHeartbeatCompleted.countDown();
+ }
+
+ final PipeSubscribeHeartbeatResp response = new
PipeSubscribeHeartbeatResp();
+ response.getTopics().put(TOPIC, new TopicConfig());
+ response.getEndPoints().put(1, new TEndPoint(HOST, FIRST_PORT));
+ response.getEndPoints().put(2, new TEndPoint(HOST, FIRST_PORT + 1));
+ return response;
+ }
+ }
+}
diff --git
a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
index f67d8829263..8f734b24629 100644
---
a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
+++
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
@@ -125,7 +125,7 @@ public class SubscriptionConsumerLifecycleTest {
consumer.fenceOnHeartbeat = true;
providers.heartbeat(consumer);
- Assert.assertTrue(consumer.isFenced());
+ waitUntil(consumer::isFenced);
providers.sync(consumer);
providers.heartbeat(consumer);
Assert.assertEquals(1, consumer.createdProviders.size());
@@ -375,6 +375,14 @@ public class SubscriptionConsumerLifecycleTest {
return (AbstractSubscriptionProviders) field.get(consumer);
}
+ private void waitUntil(final BooleanSupplier condition) throws
InterruptedException {
+ final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
+ while (!condition.getAsBoolean() && System.nanoTime() < deadline) {
+ Thread.sleep(10L);
+ }
+ Assert.assertTrue(condition.getAsBoolean());
+ }
+
@Test
public void testConcurrentPullConsumerCloseReturnsWithoutWaiting() throws
Exception {
final CountDownLatch providerCloseStarted = new CountDownLatch(1);