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 593c3bc5c33 [Subscription] Skip runtime when subscription is disabled
(#18465)
593c3bc5c33 is described below
commit 593c3bc5c338abe27776a64ece7ba25dbab1e230
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 25 10:10:09 2026 +0800
[Subscription] Skip runtime when subscription is disabled (#18465)
* Subscription: skip runtime when disabled
* Use unsupported status when subscription is disabled
---
.../protocol/thrift/impl/ClientRPCServiceImpl.java | 4 +-
.../impl/DataNodeInternalRPCServiceImpl.java | 25 +++++++--
.../java/org/apache/iotdb/db/service/DataNode.java | 5 +-
.../agent/SubscriptionReceiverAgent.java | 2 +-
.../ConsensusSubscriptionSetupHandler.java | 9 +++
...rnalRPCServiceImplSubscriptionDisabledTest.java | 64 ++++++++++++++++++++++
.../agent/SubscriptionReceiverAgentTest.java | 47 ++++++++++++++++
.../ConsensusSubscriptionSetupHandlerTest.java | 15 +++++
8 files changed, 163 insertions(+), 8 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
index 790e022ded3..be0a00bddf8 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java
@@ -3635,7 +3635,9 @@ public class ClientRPCServiceImpl implements
IClientRPCServiceWithHandler {
}
PipeDataNodeAgent.receiver().thrift().handleClientExit();
PipeDataNodeAgent.receiver().legacy().handleClientExit();
- SubscriptionAgent.receiver().handleClientExit();
+ if (COMMON_CONFIG.getSubscriptionEnabled()) {
+ SubscriptionAgent.receiver().handleClientExit();
+ }
}
/**
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
index cde5b09b457..9e847795d64 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -1627,6 +1627,11 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
@Override
public TPullCommitProgressResp pullCommitProgress(TPullCommitProgressReq
req) {
+ if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ return new
TPullCommitProgressResp(RpcUtils.getStatus(TSStatusCode.UNSUPPORTED_OPERATION))
+ .setCommitRegionProgress(Collections.emptyMap());
+ }
+
try {
final int dataNodeId =
IoTDBDescriptor.getInstance().getConfig().getDataNodeId();
final Map<String, ByteBuffer> regionProgress =
@@ -1643,6 +1648,10 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
@Override
public TSStatus syncSubscriptionProgress(TSyncSubscriptionProgressReq req) {
+ if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ return RpcUtils.getStatus(TSStatusCode.UNSUPPORTED_OPERATION);
+ }
+
try {
SubscriptionAgent.broker()
.receiveSubscriptionProgress(
@@ -1664,6 +1673,10 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
@Override
public TSStatus pushSubscriptionRuntime(TPushSubscriptionRuntimeReq req) {
+ if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ return RpcUtils.getStatus(TSStatusCode.UNSUPPORTED_OPERATION);
+ }
+
try {
for (final TSubscriptionRuntimeStateEntry runtimeStateEntry :
req.getRuntimeStates()) {
ConsensusSubscriptionSetupHandler.applyRuntimeState(
@@ -2437,8 +2450,12 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
@Override
public TSStatus updateRegionCache(TRegionRouteReq req) {
- boolean result =
ClusterPartitionFetcher.getInstance().updateRegionCache(req);
- if (result) {
+ final boolean result =
ClusterPartitionFetcher.getInstance().updateRegionCache(req);
+ if (!result) {
+ return RpcUtils.getStatus(TSStatusCode.PARTITION_CACHE_UPDATE_ERROR);
+ }
+
+ if (SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
// Notify consensus subscription queues of any preferred-writer changes
try {
ConsensusSubscriptionSetupHandler.onRegionRouteChanged(
@@ -2449,10 +2466,8 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
.MISC_LOG_FAILED_TO_PROCESS_CONSENSUS_SUBSCRIPTION_ROUTE_UPDATE_80D73E2B,
e);
}
- return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
- } else {
- return RpcUtils.getStatus(TSStatusCode.PARTITION_CACHE_UPDATE_ERROR);
}
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
}
private Map<TConsensusGroupId, Boolean> getJudgedLeaders() {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNode.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNode.java
index 48484268623..760c5747c1a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNode.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNode.java
@@ -51,6 +51,7 @@ import org.apache.iotdb.commons.service.JMXService;
import org.apache.iotdb.commons.service.RegisterManager;
import org.apache.iotdb.commons.service.ServiceType;
import org.apache.iotdb.commons.service.metric.MetricService;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
import org.apache.iotdb.commons.trigger.TriggerInformation;
import org.apache.iotdb.commons.trigger.exception.TriggerManagementException;
import org.apache.iotdb.commons.trigger.service.TriggerExecutableManager;
@@ -941,7 +942,9 @@ public class DataNode extends ServerCommandLine implements
DataNodeMBean {
registerInternalRPCService();
// Register subscription agent before pipe agent
- registerManager.register(SubscriptionAgent.runtime());
+ if (SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ registerManager.register(SubscriptionAgent.runtime());
+ }
registerManager.register(PipeDataNodeAgent.runtime());
// Start GRASS Service
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
index d918b4aed47..1f26397cbcc 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
@@ -75,7 +75,7 @@ public class SubscriptionReceiverAgent {
private final ScheduledExecutorService receiverTimeoutChecker;
SubscriptionReceiverAgent() {
- this(SubscriptionReceiverV1::new, true);
+ this(SubscriptionReceiverV1::new,
SubscriptionConfig.getInstance().getSubscriptionEnabled());
}
SubscriptionReceiverAgent(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java
index 01e9b3d8232..87b3c66f698 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java
@@ -28,6 +28,7 @@ import
org.apache.iotdb.commons.pipe.datastructure.pattern.IoTDBTreePattern;
import org.apache.iotdb.commons.pipe.datastructure.pattern.PrefixTreePattern;
import org.apache.iotdb.commons.pipe.datastructure.pattern.TablePattern;
import org.apache.iotdb.commons.pipe.datastructure.pattern.TreePattern;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.consensus.IConsensus;
import org.apache.iotdb.consensus.iot.IoTConsensus;
@@ -758,6 +759,10 @@ public class ConsensusSubscriptionSetupHandler {
public static void applyRuntimeState(
final TConsensusGroupId groupId, final ConsensusRegionRuntimeState
runtimeState) {
+ if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ return;
+ }
+
final int newPreferredNodeId = runtimeState.getPreferredWriterNodeId();
final Integer oldPreferredBoxed = lastKnownPreferredWriter.put(groupId,
newPreferredNodeId);
final int oldPreferredNodeId = (oldPreferredBoxed != null) ?
oldPreferredBoxed : -1;
@@ -789,6 +794,10 @@ public class ConsensusSubscriptionSetupHandler {
public static void onRegionRouteChanged(
final Map<TConsensusGroupId, TRegionReplicaSet> newMap, final long
routingTimestamp) {
+ if (!SubscriptionConfig.getInstance().getSubscriptionEnabled()) {
+ return;
+ }
+
final int myNodeId = IOTDB_CONFIG.getDataNodeId();
for (final Map.Entry<TConsensusGroupId, TRegionReplicaSet> newEntry :
newMap.entrySet()) {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplSubscriptionDisabledTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplSubscriptionDisabledTest.java
new file mode 100644
index 00000000000..4b61cd72bfd
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplSubscriptionDisabledTest.java
@@ -0,0 +1,64 @@
+/*
+ * 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.db.protocol.thrift.impl;
+
+import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.service.DataNode.DataNodeContext;
+import org.apache.iotdb.mpp.rpc.thrift.TPullCommitProgressResp;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.BeforeClass;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+public class DataNodeInternalRPCServiceImplSubscriptionDisabledTest {
+
+ @BeforeClass
+ public static void setUp() {
+ IoTDBDescriptor.getInstance().getConfig().setDataNodeId(0);
+ }
+
+ @Test
+ public void testSubscriptionRuntimeRPCsAreNoOpWhenSubscriptionIsDisabled() {
+ final boolean subscriptionEnabled =
+ CommonDescriptor.getInstance().getConfig().getSubscriptionEnabled();
+ try {
+ CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(false);
+ final DataNodeInternalRPCServiceImpl service =
+ new
DataNodeInternalRPCServiceImpl(Mockito.mock(DataNodeContext.class));
+
+ final TPullCommitProgressResp pullResp =
service.pullCommitProgress(null);
+ Assert.assertEquals(
+ TSStatusCode.UNSUPPORTED_OPERATION.getStatusCode(),
pullResp.getStatus().getCode());
+ Assert.assertTrue(pullResp.isSetCommitRegionProgress());
+ Assert.assertTrue(pullResp.getCommitRegionProgress().isEmpty());
+ Assert.assertEquals(
+ TSStatusCode.UNSUPPORTED_OPERATION.getStatusCode(),
+ service.syncSubscriptionProgress(null).getCode());
+ Assert.assertEquals(
+ TSStatusCode.UNSUPPORTED_OPERATION.getStatusCode(),
+ service.pushSubscriptionRuntime(null).getCode());
+ } finally {
+
CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(subscriptionEnabled);
+ }
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgentTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgentTest.java
index a151e83e21d..29d4f4f47f8 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgentTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgentTest.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.subscription.agent;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.db.subscription.receiver.SubscriptionReceiver;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -38,10 +39,12 @@ import org.junit.Assert;
import org.junit.Test;
import java.io.IOException;
+import java.lang.reflect.Field;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
@@ -49,6 +52,43 @@ import java.util.function.Supplier;
public class SubscriptionReceiverAgentTest {
+ @Test
+ public void testTimeoutCheckerIsNotScheduledWhenSubscriptionIsDisabled()
throws Exception {
+ final boolean subscriptionEnabled =
+ CommonDescriptor.getInstance().getConfig().getSubscriptionEnabled();
+ try {
+ CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(false);
+
+ final SubscriptionReceiverAgent agent = new SubscriptionReceiverAgent();
+
+ Assert.assertNull(getReceiverTimeoutChecker(agent));
+ } finally {
+
CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(subscriptionEnabled);
+ }
+ }
+
+ @Test
+ public void testTimeoutCheckerIsScheduledWhenSubscriptionIsEnabled() throws
Exception {
+ final boolean subscriptionEnabled =
+ CommonDescriptor.getInstance().getConfig().getSubscriptionEnabled();
+ SubscriptionReceiverAgent agent = null;
+ try {
+ CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(true);
+
+ agent = new SubscriptionReceiverAgent();
+
+ Assert.assertNotNull(getReceiverTimeoutChecker(agent));
+ } finally {
+ if (agent != null) {
+ final ScheduledExecutorService receiverTimeoutChecker =
getReceiverTimeoutChecker(agent);
+ if (receiverTimeoutChecker != null) {
+ receiverTimeoutChecker.shutdownNow();
+ }
+ }
+
CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(subscriptionEnabled);
+ }
+ }
+
@Test
public void testDisconnectedReceiverIsRetainedUntilTimeout() throws
IOException {
final CopyOnWriteArrayList<FakeSubscriptionReceiver> receivers = new
CopyOnWriteArrayList<>();
@@ -176,6 +216,13 @@ public class SubscriptionReceiverAgentTest {
return new SubscriptionReceiverAgent(constructor, false);
}
+ private ScheduledExecutorService getReceiverTimeoutChecker(final
SubscriptionReceiverAgent agent)
+ throws Exception {
+ final Field field =
SubscriptionReceiverAgent.class.getDeclaredField("receiverTimeoutChecker");
+ field.setAccessible(true);
+ return (ScheduledExecutorService) field.get(agent);
+ }
+
private TPipeSubscribeReq createHandshakeRequest(
final String consumerGroupId, final String consumerId) throws
IOException {
final Map<String, String> attributes = new HashMap<>();
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandlerTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandlerTest.java
index 2bec37f65d2..3e0a0351ff8 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandlerTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandlerTest.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.subscription.broker.consensus;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.consensus.DataRegionId;
import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -49,6 +50,20 @@ public class ConsensusSubscriptionSetupHandlerTest {
@Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+ @Test
+ public void testRuntimeUpdatesAreIgnoredWhenSubscriptionIsDisabled() {
+ final boolean subscriptionEnabled =
+ CommonDescriptor.getInstance().getConfig().getSubscriptionEnabled();
+ try {
+ CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(false);
+
+ ConsensusSubscriptionSetupHandler.applyRuntimeState(null, null);
+ ConsensusSubscriptionSetupHandler.onRegionRouteChanged(null, 0);
+ } finally {
+
CommonDescriptor.getInstance().getConfig().setSubscriptionEnabled(subscriptionEnabled);
+ }
+ }
+
@Test
public void testSingleTopicSetupFailurePropagates() {
SubscriptionException failure = null;