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);

Reply via email to