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