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 473672dfc32 Fix fenced owner heartbeat lease refresh (#18717)
473672dfc32 is described below

commit 473672dfc326e1325e64df25c3bacbfb03012328
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 16:58:53 2026 +0800

    Fix fenced owner heartbeat lease refresh (#18717)
---
 .../receiver/SubscriptionReceiverV1.java           | 32 +++++----
 .../receiver/SubscriptionReceiverV1Test.java       | 78 +++++++++++++++++++++-
 2 files changed, 95 insertions(+), 15 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
index 697645cde13..57c768e4c7c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
@@ -408,19 +408,6 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
 
     final List<SubscriptionCommitContext> processorBufferedCommitContexts =
         req.getProcessorBufferedCommitContexts();
-    final int refreshedCount =
-        SubscriptionAgent.broker()
-            .refreshInFlightEventLeases(consumerConfig, 
processorBufferedCommitContexts);
-    if (Objects.nonNull(processorBufferedCommitContexts)
-        && !processorBufferedCommitContexts.isEmpty()) {
-      LOGGER.debug(
-          DataNodePipeMessages
-              
.PIPE_LOG_SUBSCRIPTION_CONSUMER_REFRESHED_OF_PROCESSOR_BUFFERED_COMMIT_8C7A352A,
-          consumerConfig,
-          refreshedCount,
-          processorBufferedCommitContexts.size());
-    }
-
     final Set<String> subscribedTopicNames =
         SubscriptionAgent.consumer()
             .getTopicNamesSubscribedByConsumer(
@@ -439,6 +426,18 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
       return 
PipeSubscribeHeartbeatResp.toTPipeSubscribeResp(readPermissionStatus);
     }
 
+    final int refreshedCount =
+        refreshInFlightEventLeases(consumerConfig, 
processorBufferedCommitContexts);
+    if (Objects.nonNull(processorBufferedCommitContexts)
+        && !processorBufferedCommitContexts.isEmpty()) {
+      LOGGER.debug(
+          DataNodePipeMessages
+              
.PIPE_LOG_SUBSCRIPTION_CONSUMER_REFRESHED_OF_PROCESSOR_BUFFERED_COMMIT_8C7A352A,
+          consumerConfig,
+          refreshedCount,
+          processorBufferedCommitContexts.size());
+    }
+
     LOGGER.info(DataNodeMiscMessages.SUBSCRIPTION_CONSUMER_HEARTBEAT_SUCCESS, 
consumerConfig);
 
     // fetch subscribed topics
@@ -487,6 +486,13 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
         RpcUtils.SUCCESS_STATUS, topics, endPoints, topicNamesToUnsubscribe);
   }
 
+  protected int refreshInFlightEventLeases(
+      final ConsumerConfig consumerConfig,
+      final List<SubscriptionCommitContext> processorBufferedCommitContexts) {
+    return SubscriptionAgent.broker()
+        .refreshInFlightEventLeases(consumerConfig, 
processorBufferedCommitContexts);
+  }
+
   private TPipeSubscribeResp handlePipeSubscribeSubscribe(final 
PipeSubscribeSubscribeReq req) {
     try {
       return handlePipeSubscribeSubscribeInternal(req);
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1Test.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1Test.java
index 807d7231514..53f42666c65 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1Test.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1Test.java
@@ -20,12 +20,15 @@
 package org.apache.iotdb.db.subscription.receiver;
 
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
+import org.apache.iotdb.commons.subscription.meta.consumer.ConsumerGroupMeta;
+import org.apache.iotdb.commons.subscription.meta.consumer.ConsumerMeta;
 import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta;
 import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.rpc.subscription.config.ConsumerConfig;
 import org.apache.iotdb.rpc.subscription.config.ConsumerConstant;
 import org.apache.iotdb.rpc.subscription.config.TopicConstant;
+import 
org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionCommitContext;
 import 
org.apache.iotdb.rpc.subscription.payload.request.PipeSubscribeHandshakeReq;
 import 
org.apache.iotdb.rpc.subscription.payload.request.SubscriptionHeartbeatReq;
 
@@ -289,6 +292,63 @@ public class SubscriptionReceiverV1Test {
     }
   }
 
+  @Test
+  @SuppressWarnings("unchecked")
+  public void testFencedOwnerHeartbeatDoesNotRefreshBufferedEventLeases() 
throws Exception {
+    final String topicName = "topic-" + UUID.randomUUID();
+    final String consumerGroupId = "group-" + UUID.randomUUID();
+    final String consumerId = "consumer-" + UUID.randomUUID();
+    final ConsumerConfig oldOwnerConsumer =
+        createConsumerConfig(1_000L, "owner1", 5L, consumerId, 
consumerGroupId);
+    final SubscriptionCommitContext bufferedContext =
+        new SubscriptionCommitContext(1, 0, topicName, consumerGroupId, 1L);
+
+    final TopicMeta topicMeta = createTopicMeta(topicName, "owner1", 5L);
+    SubscriptionAgent.topic().handleSingleTopicMetaChanges(topicMeta);
+    final ConsumerGroupMeta consumerGroupMeta =
+        new ConsumerGroupMeta(
+            consumerGroupId,
+            System.currentTimeMillis(),
+            new ConsumerMeta(
+                consumerId, System.currentTimeMillis(), 
oldOwnerConsumer.getAttribute()));
+    consumerGroupMeta.addSubscription(consumerId, 
Collections.singleton(topicName));
+    
SubscriptionAgent.consumer().handleSingleConsumerGroupMetaChanges(consumerGroupMeta);
+
+    try {
+      final TopicMeta transferredTopicMeta = topicMeta.deepCopy();
+      transferredTopicMeta.transferOwner("owner2", 6L);
+      
SubscriptionAgent.topic().handleSingleTopicMetaChanges(transferredTopicMeta);
+
+      final AtomicLong refreshCount = new AtomicLong();
+      final SubscriptionReceiverV1 receiver =
+          new SubscriptionReceiverV1() {
+            @Override
+            protected int refreshInFlightEventLeases(
+                final ConsumerConfig consumerConfig,
+                final java.util.List<SubscriptionCommitContext> 
processorBufferedCommitContexts) {
+              refreshCount.incrementAndGet();
+              return 1;
+            }
+          };
+      setField(receiver, "sharedConsumerConfig", oldOwnerConsumer);
+      final ThreadLocal<ConsumerConfig> consumerConfigThreadLocal =
+          (ThreadLocal<ConsumerConfig>) getField(receiver, 
"consumerConfigThreadLocal");
+      consumerConfigThreadLocal.set(oldOwnerConsumer);
+
+      Assert.assertEquals(
+          TSStatusCode.SUBSCRIPTION_OWNER_FENCED.getStatusCode(),
+          receiver
+              .handle(
+                  
SubscriptionHeartbeatReq.toThriftReq(Collections.singletonList(bufferedContext)))
+              .getStatus()
+              .getCode());
+      Assert.assertEquals(0L, refreshCount.get());
+    } finally {
+      SubscriptionAgent.consumer().handleDropConsumerGroup(consumerGroupId);
+      SubscriptionAgent.topic().handleDropTopic(topicName);
+    }
+  }
+
   private long invokeCalculateConsumerInactivityTimeoutMs(
       final SubscriptionReceiverV1 receiver, final ConsumerConfig 
consumerConfig) throws Exception {
     final Method method =
@@ -312,9 +372,23 @@ public class SubscriptionReceiverV1Test {
 
   private ConsumerConfig createConsumerConfig(
       final long heartbeatIntervalMs, final String ownerId, final Long 
ownerEpoch) {
+    return createConsumerConfig(
+        heartbeatIntervalMs,
+        ownerId,
+        ownerEpoch,
+        "consumer-" + UUID.randomUUID(),
+        "group-" + UUID.randomUUID());
+  }
+
+  private ConsumerConfig createConsumerConfig(
+      final long heartbeatIntervalMs,
+      final String ownerId,
+      final Long ownerEpoch,
+      final String consumerId,
+      final String consumerGroupId) {
     final Map<String, String> attributes = new HashMap<>();
-    attributes.put(ConsumerConstant.CONSUMER_ID_KEY, "consumer-" + 
UUID.randomUUID());
-    attributes.put(ConsumerConstant.CONSUMER_GROUP_ID_KEY, "group-" + 
UUID.randomUUID());
+    attributes.put(ConsumerConstant.CONSUMER_ID_KEY, consumerId);
+    attributes.put(ConsumerConstant.CONSUMER_GROUP_ID_KEY, consumerGroupId);
     attributes.put(ConsumerConstant.HEARTBEAT_INTERVAL_MS_KEY, 
String.valueOf(heartbeatIntervalMs));
     if (ownerId != null) {
       attributes.put(ConsumerConstant.OWNER_ID_KEY, ownerId);

Reply via email to