This is an automated email from the ASF dual-hosted git repository.

rong 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 630597e502a Subscription: improve subscription meta management with 
sub-procedures strong exception handling semantics (#13698)
630597e502a is described below

commit 630597e502a4d053ee42385bbcab79a0036e5f58
Author: V_Galaxy <[email protected]>
AuthorDate: Mon Oct 28 10:36:12 2024 +0800

    Subscription: improve subscription meta management with sub-procedures 
strong exception handling semantics (#13698)
---
 .../subscription/SubscriptionMetaSyncer.java       |  38 +--
 .../persistence/subscription/SubscriptionInfo.java |  15 +-
 .../impl/pipe/runtime/PipeMetaSyncProcedure.java   |   3 +-
 .../AbstractOperateSubscriptionProcedure.java      |  16 +-
 .../consumer/AlterConsumerGroupProcedure.java      |  60 ++---
 .../runtime/ConsumerGroupMetaSyncProcedure.java    |   6 +-
 .../subscription/CreateSubscriptionProcedure.java  | 260 ++++++++-------------
 .../subscription/DropSubscriptionProcedure.java    | 147 ++----------
 .../subscription/topic/AlterTopicProcedure.java    |  44 ++--
 .../subscription/topic/CreateTopicProcedure.java   |  33 ++-
 .../subscription/topic/DropTopicProcedure.java     |   5 +-
 .../topic/runtime/TopicMetaSyncProcedure.java      |   6 +-
 .../CreateSubscriptionProcedureTest.java           |   9 -
 .../DropSubscriptionProcedureTest.java             |   8 -
 .../agent/SubscriptionBrokerAgent.java             |  78 +++----
 .../agent/SubscriptionConsumerAgent.java           |  27 ++-
 .../db/subscription/broker/SubscriptionBroker.java |  12 +-
 .../apache/iotdb/commons/conf/CommonConfig.java    |  22 ++
 .../iotdb/commons/conf/CommonDescriptor.java       |  12 +
 .../subscription/config/SubscriptionConfig.java    |  18 +-
 .../meta/consumer/ConsumerGroupMeta.java           |   4 +
 .../meta/consumer/ConsumerGroupMetaKeeper.java     |  25 ++
 .../commons/subscription/meta/topic/TopicMeta.java |  19 +-
 23 files changed, 384 insertions(+), 483 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/SubscriptionMetaSyncer.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/SubscriptionMetaSyncer.java
index 50a4ebb22b2..de49987e13f 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/SubscriptionMetaSyncer.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/SubscriptionMetaSyncer.java
@@ -23,7 +23,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.commons.concurrent.ThreadName;
 import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
-import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.confignode.manager.ConfigManager;
 import org.apache.iotdb.confignode.manager.ProcedureManager;
 import org.apache.iotdb.rpc.TSStatusCode;
@@ -43,9 +43,9 @@ public class SubscriptionMetaSyncer {
       IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
           ThreadName.SUBSCRIPTION_RUNTIME_META_SYNCER.getName());
   private static final long INITIAL_SYNC_DELAY_MINUTES =
-      PipeConfig.getInstance().getPipeMetaSyncerInitialSyncDelayMinutes();
+      
SubscriptionConfig.getInstance().getSubscriptionMetaSyncerInitialSyncDelayMinutes();
   private static final long SYNC_INTERVAL_MINUTES =
-      PipeConfig.getInstance().getPipeMetaSyncerSyncIntervalMinutes();
+      
SubscriptionConfig.getInstance().getSubscriptionMetaSyncerSyncIntervalMinutes();
 
   private final ConfigManager configManager;
 
@@ -89,22 +89,26 @@ public class SubscriptionMetaSyncer {
     }
 
     final ProcedureManager procedureManager = 
configManager.getProcedureManager();
-    final TSStatus consumerGroupMetaSyncStatus = 
procedureManager.consumerGroupMetaSync();
+
+    // sync topic meta firstly
+    // TODO: consider drop the topic which is subscribed by consumers
     final TSStatus topicMetaSyncStatus = procedureManager.topicMetaSync();
-    if (consumerGroupMetaSyncStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()
-        && topicMetaSyncStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      LOGGER.info(
-          "After this successful sync, if SubscriptionInfo is empty during 
this sync and has not been modified afterwards, all subsequent syncs will be 
skipped");
-      isLastSubscriptionSyncSuccessful = true;
-    } else {
-      if (consumerGroupMetaSyncStatus.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        LOGGER.warn(
-            "Failed to sync consumer group meta. Result status: {}.", 
consumerGroupMetaSyncStatus);
-      }
-      if (topicMetaSyncStatus.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        LOGGER.warn("Failed to sync topic meta. Result status: {}.", 
topicMetaSyncStatus);
-      }
+    if (topicMetaSyncStatus.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      LOGGER.warn("Failed to sync topic meta. Result status: {}.", 
topicMetaSyncStatus);
+      return;
     }
+
+    // sync consumer meta if syncing topic meta successfully
+    final TSStatus consumerGroupMetaSyncStatus = 
procedureManager.consumerGroupMetaSync();
+    if (consumerGroupMetaSyncStatus.getCode() != 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      LOGGER.warn(
+          "Failed to sync consumer group meta. Result status: {}.", 
consumerGroupMetaSyncStatus);
+      return;
+    }
+
+    LOGGER.info(
+        "After this successful sync, if SubscriptionInfo is empty during this 
sync and has not been modified afterwards, all subsequent syncs will be 
skipped");
+    isLastSubscriptionSyncSuccessful = true;
   }
 
   public synchronized void stop() {
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java
index d1071c511af..2ccde563866 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java
@@ -194,7 +194,7 @@ public class SubscriptionInfo implements SnapshotProcessor {
       // executed on all nodes to ensure the consistency.
       return;
     } else {
-      if (!topicMeta.hasSubscribedConsumerGroup()) {
+      if 
(!consumerGroupMetaKeeper.isTopicSubscribedByConsumerGroup(topicName)) {
         return;
       }
     }
@@ -506,6 +506,16 @@ public class SubscriptionInfo implements SnapshotProcessor 
{
     }
   }
 
+  public boolean isTopicSubscribedByConsumerGroup(
+      final String topicName, final String consumerGroupId) {
+    acquireReadLock();
+    try {
+      return 
consumerGroupMetaKeeper.isTopicSubscribedByConsumerGroup(topicName, 
consumerGroupId);
+    } finally {
+      releaseReadLock();
+    }
+  }
+
   public TSStatus alterConsumerGroup(AlterConsumerGroupPlan plan) {
     acquireWriteLock();
     try {
@@ -630,7 +640,8 @@ public class SubscriptionInfo implements SnapshotProcessor {
   private List<SubscriptionMeta> getAllSubscriptionMeta() {
     List<SubscriptionMeta> allSubscriptions = new ArrayList<>();
     for (TopicMeta topicMeta : topicMetaKeeper.getAllTopicMeta()) {
-      for (String consumerGroupId : topicMeta.getSubscribedConsumerGroupIds()) 
{
+      for (String consumerGroupId :
+          
consumerGroupMetaKeeper.getSubscribedConsumerGroupIds(topicMeta.getTopicName()))
 {
         Set<String> subscribedConsumerIDs =
             consumerGroupMetaKeeper.getConsumersSubscribingTopic(
                 consumerGroupId, topicMeta.getTopicName());
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java
index 431858f0679..c8a734a9e35 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/runtime/PipeMetaSyncProcedure.java
@@ -130,7 +130,8 @@ public class PipeMetaSyncProcedure extends 
AbstractOperatePipeProcedureV2 {
   }
 
   @Override
-  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws 
IOException {
+  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws PipeException, IOException {
     LOGGER.info("PipeMetaSyncProcedure: executeFromOperateOnDataNodes");
 
     Map<Integer, TPushPipeMetaResp> respMap = pushPipeMetaToDataNodes(env);
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/AbstractOperateSubscriptionProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/AbstractOperateSubscriptionProcedure.java
index e661f3b3d3f..dcefa8d5f45 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/AbstractOperateSubscriptionProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/AbstractOperateSubscriptionProcedure.java
@@ -59,6 +59,11 @@ public abstract class AbstractOperateSubscriptionProcedure
 
   private static final int RETRY_THRESHOLD = 1;
 
+  // Only used in rollback to reduce the number of network calls
+  // Pure in-memory object, not involved in snapshot serialization and 
deserialization.
+  // TODO: consider serializing this variable later
+  protected boolean isRollbackFromOperateOnDataNodesSuccessful = false;
+
   // Only used in rollback to avoid executing rollbackFromValidate multiple 
times
   // Pure in-memory object, not involved in snapshot serialization and 
deserialization.
   // TODO: consider serializing this variable later
@@ -271,7 +276,9 @@ public abstract class AbstractOperateSubscriptionProcedure
         break;
       case OPERATE_ON_CONFIG_NODES:
         try {
-          rollbackFromOperateOnConfigNodes(env);
+          if (!isRollbackFromOperateOnDataNodesSuccessful) {
+            rollbackFromOperateOnConfigNodes(env);
+          }
         } catch (Exception e) {
           LOGGER.warn(
               "ProcedureId {}: Failed to rollback from state [{}], because {}",
@@ -283,7 +290,9 @@ public abstract class AbstractOperateSubscriptionProcedure
         break;
       case OPERATE_ON_DATA_NODES:
         try {
+          rollbackFromOperateOnConfigNodes(env);
           rollbackFromOperateOnDataNodes(env);
+          isRollbackFromOperateOnDataNodesSuccessful = true;
         } catch (Exception e) {
           LOGGER.warn(
               "ProcedureId {}: Failed to rollback from state [{}], because {}",
@@ -301,10 +310,11 @@ public abstract class AbstractOperateSubscriptionProcedure
 
   protected abstract void rollbackFromValidate(ConfigNodeProcedureEnv env);
 
-  protected abstract void 
rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env);
+  protected abstract void 
rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException;
 
   protected abstract void 
rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
-      throws IOException;
+      throws SubscriptionException, IOException;
 
   /**
    * Pushing all the topicMeta's to all the dataNodes.
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/AlterConsumerGroupProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/AlterConsumerGroupProcedure.java
index 6e6031b8102..69017422505 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/AlterConsumerGroupProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/AlterConsumerGroupProcedure.java
@@ -108,7 +108,6 @@ public class AlterConsumerGroupProcedure extends 
AbstractOperateSubscriptionProc
           new TSStatus(TSStatusCode.ALTER_CONSUMER_ERROR.getStatusCode())
               .setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
@@ -119,27 +118,20 @@ public class AlterConsumerGroupProcedure extends 
AbstractOperateSubscriptionProc
 
   @Override
   public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
-      throws SubscriptionException {
+      throws SubscriptionException, IOException {
     LOGGER.info(
         "AlterConsumerGroupProcedure: executeFromOperateOnDataNodes({})",
         updatedConsumerGroupMeta.getConsumerGroupId());
 
-    try {
-      final List<TSStatus> statuses =
-          
env.pushSingleConsumerGroupOnDataNode(updatedConsumerGroupMeta.serialize());
-      if (RpcUtils.squashResponseStatusList(statuses).getCode()
-          != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        throw new SubscriptionException(
-            String.format(
-                "Failed to alter consumer group %s on data nodes, because %s",
-                updatedConsumerGroupMeta.getConsumerGroupId(), statuses));
-      }
-    } catch (IOException e) {
-      LOGGER.warn("Failed to serialize the consumer group meta due to: ", e);
+    final List<TSStatus> statuses =
+        
env.pushSingleConsumerGroupOnDataNode(updatedConsumerGroupMeta.serialize());
+    if (RpcUtils.squashResponseStatusList(statuses).getCode()
+        != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
-              "Failed to alter consumer group %s on data nodes, because %s",
-              updatedConsumerGroupMeta.getConsumerGroupId(), e.getMessage()));
+              "Failed to alter consumer group (%s -> %s) on data nodes, 
because %s",
+              existingConsumerGroupMeta, updatedConsumerGroupMeta, statuses));
     }
   }
 
@@ -149,7 +141,8 @@ public class AlterConsumerGroupProcedure extends 
AbstractOperateSubscriptionProc
   }
 
   @Override
-  public void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  public void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info(
         "AlterConsumerGroupProcedure: rollbackFromOperateOnConfigNodes({})",
         updatedConsumerGroupMeta.getConsumerGroupId());
@@ -166,35 +159,28 @@ public class AlterConsumerGroupProcedure extends 
AbstractOperateSubscriptionProc
           new TSStatus(TSStatusCode.ALTER_CONSUMER_ERROR.getStatusCode())
               .setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      LOGGER.warn(
-          "Failed to rollback from altering consumer group {} on config nodes, 
because {}",
-          updatedConsumerGroupMeta.getConsumerGroupId(),
-          response);
+      throw new SubscriptionException(
+          String.format(
+              "Failed to rollback from altering consumer group (%s -> %s) on 
config nodes, because %s",
+              existingConsumerGroupMeta, updatedConsumerGroupMeta, response));
     }
   }
 
   @Override
-  public void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
+  public void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException, IOException {
     LOGGER.info("AlterConsumerGroupProcedure: rollbackFromOperateOnDataNodes");
 
-    try {
-      final List<TSStatus> statuses =
-          
env.pushSingleConsumerGroupOnDataNode(existingConsumerGroupMeta.serialize());
-      if (RpcUtils.squashResponseStatusList(statuses).getCode()
-          != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        throw new SubscriptionException(
-            String.format(
-                "Failed to rollback from altering consumer group %s on data 
nodes, because %s",
-                updatedConsumerGroupMeta.getConsumerGroupId(), statuses));
-      }
-    } catch (IOException e) {
-      LOGGER.warn("Failed to serialize the consumer group meta due to: ", e);
+    final List<TSStatus> statuses =
+        
env.pushSingleConsumerGroupOnDataNode(existingConsumerGroupMeta.serialize());
+    if (RpcUtils.squashResponseStatusList(statuses).getCode()
+        != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
-              "Failed to rollback from altering consumer group %s on data 
nodes, because %s",
-              updatedConsumerGroupMeta.getConsumerGroupId(), e.getMessage()));
+              "Failed to rollback from altering consumer group (%s -> %s) on 
data nodes, because %s",
+              existingConsumerGroupMeta, updatedConsumerGroupMeta, statuses));
     }
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/ConsumerGroupMetaSyncProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/ConsumerGroupMetaSyncProcedure.java
index bb010eaca40..93eb6c5a5fc 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/ConsumerGroupMetaSyncProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/ConsumerGroupMetaSyncProcedure.java
@@ -99,7 +99,8 @@ public class ConsumerGroupMetaSyncProcedure extends 
AbstractOperateSubscriptionP
   }
 
   @Override
-  public void executeFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  public void executeFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info("ConsumerGroupMetaSyncProcedure: 
executeFromOperateOnConfigNodes");
 
     final List<ConsumerGroupMeta> consumerGroupMetaList =
@@ -122,7 +123,8 @@ public class ConsumerGroupMetaSyncProcedure extends 
AbstractOperateSubscriptionP
   }
 
   @Override
-  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws 
IOException {
+  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException, IOException {
     LOGGER.info("ConsumerGroupMetaSyncProcedure: 
executeFromOperateOnDataNodes");
 
     Map<Integer, TPushConsumerGroupMetaResp> respMap = 
pushConsumerGroupMetaToDataNodes(env);
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java
index 2f977b9a47b..4a48ebdd35d 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java
@@ -27,8 +27,6 @@ import org.apache.iotdb.commons.utils.TestOnly;
 import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
 import 
org.apache.iotdb.confignode.consensus.request.write.pipe.task.DropPipePlanV2;
 import 
org.apache.iotdb.confignode.consensus.request.write.pipe.task.OperateMultiplePipesPlanV2;
-import 
org.apache.iotdb.confignode.consensus.request.write.subscription.topic.AlterMultipleTopicsPlan;
-import 
org.apache.iotdb.confignode.consensus.request.write.subscription.topic.AlterTopicPlan;
 import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import 
org.apache.iotdb.confignode.procedure.impl.pipe.AbstractOperatePipeProcedureV2;
 import 
org.apache.iotdb.confignode.procedure.impl.pipe.task.CreatePipeProcedureV2;
@@ -55,28 +53,26 @@ import java.util.List;
 import java.util.Objects;
 import java.util.stream.Collectors;
 
-// TODO: check if it also needs meta sync to keep CN and DN in sync
 public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPipeProcedure {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(CreateSubscriptionProcedure.class);
 
   private TSubscribeReq subscribeReq;
 
+  // execution order: alter consumer group -> create pipe
+  // rollback order: create pipe -> alter consumer group
+  // NOTE: The 'alter consumer group' operation must be performed before 
'create pipe'.
   private AlterConsumerGroupProcedure alterConsumerGroupProcedure;
-  private List<AlterTopicProcedure> alterTopicProcedures = new ArrayList<>();
   private List<CreatePipeProcedureV2> createPipeProcedures = new ArrayList<>();
 
-  // Record failed index of procedures to rollback properly.
-  // We only record fail index when executing on config nodes, because when 
executing on data nodes
-  // fails, we just push all meta to data nodes.
-  private int alterTopicProcedureFailIndexOnCN = -1;
-  private int createPipeProcedureFailIndexOnCN = -1;
+  // TODO: remove this variable later
+  private final List<AlterTopicProcedure> alterTopicProcedures = new 
ArrayList<>(); // unused now
 
   public CreateSubscriptionProcedure() {
     super();
   }
 
-  public CreateSubscriptionProcedure(TSubscribeReq subscribeReq) {
+  public CreateSubscriptionProcedure(final TSubscribeReq subscribeReq) {
     this.subscribeReq = subscribeReq;
   }
 
@@ -86,229 +82,171 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
   }
 
   @Override
-  protected boolean executeFromValidate(ConfigNodeProcedureEnv env) throws 
SubscriptionException {
+  protected boolean executeFromValidate(final ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info("CreateSubscriptionProcedure: executeFromValidate");
 
     subscriptionInfo.get().validateBeforeSubscribe(subscribeReq);
 
-    // alterConsumerGroupProcedure
+    // Construct AlterConsumerGroupProcedure
+    final String consumerGroupId = subscribeReq.getConsumerGroupId();
     final ConsumerGroupMeta updatedConsumerGroupMeta =
-        
subscriptionInfo.get().deepCopyConsumerGroupMeta(subscribeReq.getConsumerGroupId());
+        subscriptionInfo.get().deepCopyConsumerGroupMeta(consumerGroupId);
     updatedConsumerGroupMeta.addSubscription(
         subscribeReq.getConsumerId(), subscribeReq.getTopicNames());
     alterConsumerGroupProcedure =
         new AlterConsumerGroupProcedure(updatedConsumerGroupMeta, 
subscriptionInfo);
 
-    // alterTopicProcedures & createPipeProcedures
-    for (String topic : subscribeReq.getTopicNames()) {
-      TopicMeta updatedTopicMeta = 
subscriptionInfo.get().deepCopyTopicMeta(topic);
-
-      if 
(updatedTopicMeta.addSubscribedConsumerGroup(subscribeReq.getConsumerGroupId()))
 {
+    // Construct CreatePipeProcedureV2s
+    for (final String topicName : subscribeReq.getTopicNames()) {
+      final String pipeName =
+          PipeStaticMeta.generateSubscriptionPipeName(topicName, 
consumerGroupId);
+      if (!subscriptionInfo.get().isTopicSubscribedByConsumerGroup(topicName, 
consumerGroupId)
+          // even if there existed subscription meta, if there is no 
corresponding pipe meta, it
+          // will try to create the pipe
+          || !pipeTaskInfo.get().isPipeExisted(pipeName)) {
+        final TopicMeta topicMeta = 
subscriptionInfo.get().deepCopyTopicMeta(topicName);
         createPipeProcedures.add(
             new CreatePipeProcedureV2(
                 new TCreatePipeReq()
-                    .setPipeName(
-                        PipeStaticMeta.generateSubscriptionPipeName(
-                            topic, subscribeReq.getConsumerGroupId()))
-                    
.setExtractorAttributes(updatedTopicMeta.generateExtractorAttributes())
-                    
.setProcessorAttributes(updatedTopicMeta.generateProcessorAttributes())
-                    .setConnectorAttributes(
-                        updatedTopicMeta.generateConnectorAttributes(
-                            subscribeReq.getConsumerGroupId())),
+                    .setPipeName(pipeName)
+                    
.setExtractorAttributes(topicMeta.generateExtractorAttributes())
+                    
.setProcessorAttributes(topicMeta.generateProcessorAttributes())
+                    
.setConnectorAttributes(topicMeta.generateConnectorAttributes(consumerGroupId)),
                 pipeTaskInfo));
-
-        alterTopicProcedures.add(new AlterTopicProcedure(updatedTopicMeta, 
subscriptionInfo));
       }
     }
 
+    // Validate AlterConsumerGroupProcedure
     alterConsumerGroupProcedure.executeFromValidate(env);
 
-    for (AlterTopicProcedure alterTopicProcedure : alterTopicProcedures) {
-      alterTopicProcedure.executeFromValidate(env);
-    }
-
-    for (CreatePipeProcedureV2 createPipeProcedure : createPipeProcedures) {
+    // Validate CreatePipeProcedureV2s
+    for (final CreatePipeProcedureV2 createPipeProcedure : 
createPipeProcedures) {
       createPipeProcedure.executeFromValidateTask(env);
       createPipeProcedure.executeFromCalculateInfoForTask(env);
     }
+
     return true;
   }
 
-  // TODO: check periodically if the subscription is still valid but no 
working pipe?
   @Override
-  protected void executeFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+  protected void executeFromOperateOnConfigNodes(final ConfigNodeProcedureEnv 
env)
       throws SubscriptionException {
     LOGGER.info("CreateSubscriptionProcedure: 
executeFromOperateOnConfigNodes");
 
+    // Execute AlterConsumerGroupProcedure
     alterConsumerGroupProcedure.executeFromOperateOnConfigNodes(env);
 
-    TSStatus response;
-
-    List<AlterTopicPlan> alterTopicPlans =
-        alterTopicProcedures.stream()
-            .map(AlterTopicProcedure::getUpdatedTopicMeta)
-            .map(AlterTopicPlan::new)
-            .collect(Collectors.toList());
-    try {
-      response =
-          env.getConfigManager()
-              .getConsensusManager()
-              .write(new AlterMultipleTopicsPlan(alterTopicPlans));
-    } catch (ConsensusException e) {
-      LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
-      response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
-      response.setMessage(e.getMessage());
-    }
-    if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()
-        && response.getSubStatusSize() > 0) {
-      // Record the failed index for rollback
-      alterTopicProcedureFailIndexOnCN = response.getSubStatusSize() - 1;
-    }
-
-    List<ConfigPhysicalPlan> createPipePlans =
+    // Execute CreatePipeProcedureV2s
+    final List<ConfigPhysicalPlan> createPipePlans =
         createPipeProcedures.stream()
             .map(CreatePipeProcedureV2::constructPlan)
             .collect(Collectors.toList());
+    TSStatus response;
     try {
       response =
           env.getConfigManager()
               .getConsensusManager()
               .write(new OperateMultiplePipesPlanV2(createPipePlans));
-    } catch (ConsensusException e) {
+    } catch (final ConsensusException e) {
       LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
       response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
       response.setMessage(e.getMessage());
     }
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()
         && response.getSubStatusSize() > 0) {
-      // Record the failed index for rollback
-      createPipeProcedureFailIndexOnCN = response.getSubStatusSize() - 1;
+      throw new SubscriptionException(
+          String.format(
+              "Failed to create subscription with request %s on config nodes, 
because %s",
+              subscribeReq, response));
     }
   }
 
   @Override
-  protected void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+  protected void executeFromOperateOnDataNodes(final ConfigNodeProcedureEnv 
env)
       throws SubscriptionException, IOException {
     LOGGER.info("CreateSubscriptionProcedure: executeFromOperateOnDataNodes");
 
+    // Push consumer group meta to data nodes
     alterConsumerGroupProcedure.executeFromOperateOnDataNodes(env);
 
-    // push topic meta to data nodes
-    List<ByteBuffer> topicMetaBinaryList = new ArrayList<>();
-    for (AlterTopicProcedure alterTopicProcedure : alterTopicProcedures) {
-      
topicMetaBinaryList.add(alterTopicProcedure.getUpdatedTopicMeta().serialize());
-    }
-    if 
(pushTopicMetaHasException(env.pushMultiTopicMetaToDataNodes(topicMetaBinaryList)))
 {
-      // If not all topic meta are pushed successfully, the meta can be pushed 
during meta sync.
-      LOGGER.warn(
-          "Failed to alter topics when creating subscription, metadata will be 
synchronized later.");
-    }
-
-    // push pipe meta to data nodes
-    List<String> pipeNames =
+    // Push pipe meta to data nodes
+    final List<String> pipeNames =
         createPipeProcedures.stream()
             .map(CreatePipeProcedureV2::getPipeName)
             .collect(Collectors.toList());
-    String exceptionMessage =
+    final String exceptionMessage =
         AbstractOperatePipeProcedureV2.parsePushPipeMetaExceptionForPipe(
             null, pushMultiPipeMetaToDataNodes(pipeNames, env));
     if (!exceptionMessage.isEmpty()) {
-      // If not all pipe meta are pushed successfully, the meta can be pushed 
during meta sync.
-      LOGGER.warn(
-          "Failed to create pipes {} when creating subscription, details: {}, 
metadata will be synchronized later.",
-          pipeNames,
-          exceptionMessage);
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
+      throw new SubscriptionException(
+          String.format(
+              "Failed to create pipes %s when creating subscription with 
request %s, details: %s, metadata will be synchronized later.",
+              pipeNames, subscribeReq, exceptionMessage));
     }
   }
 
   @Override
-  protected void rollbackFromValidate(ConfigNodeProcedureEnv env) {
+  protected void rollbackFromValidate(final ConfigNodeProcedureEnv env) {
     LOGGER.info("CreateSubscriptionProcedure: rollbackFromValidate");
   }
 
   @Override
-  protected void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  protected void rollbackFromOperateOnConfigNodes(final ConfigNodeProcedureEnv 
env)
+      throws SubscriptionException {
     LOGGER.info("CreateSubscriptionProcedure: 
rollbackFromOperateOnConfigNodes");
 
-    // TODO: roll back from the last executed procedure to the first executed
-    alterConsumerGroupProcedure.rollbackFromOperateOnConfigNodes(env);
-
+    // Rollback CreatePipeProcedureV2s
+    final List<ConfigPhysicalPlan> dropPipePlans =
+        createPipeProcedures.stream()
+            .map(procedure -> new DropPipePlanV2(procedure.getPipeName()))
+            .collect(Collectors.toList());
     TSStatus response;
-
-    // rollback alterTopicProcedures
-    List<AlterTopicPlan> alterTopicRollbackPlans = new ArrayList<>();
-    for (int i = 0;
-        i <= Math.min(alterTopicProcedureFailIndexOnCN, 
alterTopicProcedures.size());
-        i++) {
-      alterTopicRollbackPlans.add(
-          new 
AlterTopicPlan(alterTopicProcedures.get(i).getExistedTopicMeta()));
-    }
-    try {
-      response =
-          env.getConfigManager()
-              .getConsensusManager()
-              .write(new AlterMultipleTopicsPlan(alterTopicRollbackPlans));
-    } catch (ConsensusException e) {
-      LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
-      response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
-      response.setMessage(e.getMessage());
-    }
-    // if failed to rollback, throw exception
-    if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      throw new SubscriptionException(response.getMessage());
-    }
-
-    // rollback createPipeProcedures
-    List<ConfigPhysicalPlan> dropPipePlans = new ArrayList<>();
-    for (int i = 0;
-        i <= Math.min(createPipeProcedureFailIndexOnCN, 
createPipeProcedures.size());
-        i++) {
-      dropPipePlans.add(new 
DropPipePlanV2(createPipeProcedures.get(i).getPipeName()));
-    }
     try {
       response =
           env.getConfigManager()
               .getConsensusManager()
               .write(new OperateMultiplePipesPlanV2(dropPipePlans));
-    } catch (ConsensusException e) {
+    } catch (final ConsensusException e) {
       LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
       response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
       response.setMessage(e.getMessage());
     }
-    // if failed to rollback, throw exception
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      throw new SubscriptionException(response.getMessage());
+      throw new SubscriptionException(
+          String.format(
+              "Failed to rollback creating subscription with request %s on 
config nodes, because %s",
+              subscribeReq, response));
     }
+
+    // Rollback AlterConsumerGroupProcedure
+    alterConsumerGroupProcedure.rollbackFromOperateOnConfigNodes(env);
   }
 
   @Override
-  protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) 
throws IOException {
+  protected void rollbackFromOperateOnDataNodes(final ConfigNodeProcedureEnv 
env)
+      throws SubscriptionException, IOException {
     LOGGER.info("CreateSubscriptionProcedure: rollbackFromOperateOnDataNodes");
 
-    // TODO: roll back from the last executed procedure to the first executed
-    alterConsumerGroupProcedure.rollbackFromOperateOnDataNodes(env);
-
-    // Push all topic metas to datanode, may be time-consuming
-    if (pushTopicMetaHasException(pushTopicMetaToDataNodes(env))) {
-      LOGGER.warn(
-          "Failed to rollback alter topics when creating subscription, 
metadata will be synchronized later.");
-    }
-
     // Push all pipe metas to datanode, may be time-consuming
-    String exceptionMessage =
+    final String exceptionMessage =
         AbstractOperatePipeProcedureV2.parsePushPipeMetaExceptionForPipe(
             null, AbstractOperatePipeProcedureV2.pushPipeMetaToDataNodes(env, 
pipeTaskInfo));
     if (!exceptionMessage.isEmpty()) {
-      LOGGER.warn(
-          "Failed to rollback create pipes when creating subscription, 
details: {}, metadata will be synchronized later.",
-          exceptionMessage);
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
+      throw new SubscriptionException(
+          String.format(
+              "Failed to rollback create pipes when creating subscription with 
request %s, because %s",
+              subscribeReq, exceptionMessage));
     }
-  }
 
-  // TODO: we still need some strategies to clean the subscription if it's not 
valid anymore
+    // Rollback AlterConsumerGroupProcedure
+    alterConsumerGroupProcedure.rollbackFromOperateOnDataNodes(env);
+  }
 
   @Override
-  public void serialize(DataOutputStream stream) throws IOException {
+  public void serialize(final DataOutputStream stream) throws IOException {
     
stream.writeShort(ProcedureType.CREATE_SUBSCRIPTION_PROCEDURE.getTypeCode());
 
     super.serialize(stream);
@@ -318,12 +256,12 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
     final int size = subscribeReq.getTopicNamesSize();
     ReadWriteIOUtils.write(size, stream);
     if (size != 0) {
-      for (String topicName : subscribeReq.getTopicNames()) {
+      for (final String topicName : subscribeReq.getTopicNames()) {
         ReadWriteIOUtils.write(topicName, stream);
       }
     }
 
-    // serialize consumerGroupProcedure
+    // Serialize AlterConsumerGroupProcedure
     if (alterConsumerGroupProcedure != null) {
       ReadWriteIOUtils.write(true, stream);
       alterConsumerGroupProcedure.serialize(stream);
@@ -331,22 +269,22 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
       ReadWriteIOUtils.write(false, stream);
     }
 
-    // serialize topic procedures
+    // Serialize AlterTopicProcedures
     if (alterTopicProcedures != null) {
       ReadWriteIOUtils.write(true, stream);
       ReadWriteIOUtils.write(alterTopicProcedures.size(), stream);
-      for (AlterTopicProcedure topicProcedure : alterTopicProcedures) {
+      for (final AlterTopicProcedure topicProcedure : alterTopicProcedures) {
         topicProcedure.serialize(stream);
       }
     } else {
       ReadWriteIOUtils.write(false, stream);
     }
 
-    // serialize pipe procedures
+    // Serialize CreatePipeProcedureV2s
     if (createPipeProcedures != null) {
       ReadWriteIOUtils.write(true, stream);
       ReadWriteIOUtils.write(createPipeProcedures.size(), stream);
-      for (CreatePipeProcedureV2 pipeProcedure : createPipeProcedures) {
+      for (final CreatePipeProcedureV2 pipeProcedure : createPipeProcedures) {
         pipeProcedure.serialize(stream);
       }
     } else {
@@ -355,7 +293,7 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
   }
 
   @Override
-  public void deserialize(ByteBuffer byteBuffer) {
+  public void deserialize(final ByteBuffer byteBuffer) {
     super.deserialize(byteBuffer);
 
     subscribeReq =
@@ -368,7 +306,7 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
       
subscribeReq.getTopicNames().add(ReadWriteIOUtils.readString(byteBuffer));
     }
 
-    // deserialize consumerGroupProcedure
+    // Deserialize AlterConsumerGroupProcedure
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       // This readShort should return ALTER_CONSUMER_GROUP_PROCEDURE, and we 
ignore it.
       ReadWriteIOUtils.readShort(byteBuffer);
@@ -377,27 +315,27 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
       alterConsumerGroupProcedure.deserialize(byteBuffer);
     }
 
-    // deserialize topic procedures
+    // Deserialize AlterTopicProcedures
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       size = ReadWriteIOUtils.readInt(byteBuffer);
       for (int i = 0; i < size; ++i) {
         // This readShort should return ALTER_TOPIC_PROCEDURE, and we ignore 
it.
         ReadWriteIOUtils.readShort(byteBuffer);
 
-        AlterTopicProcedure topicProcedure = new AlterTopicProcedure();
+        final AlterTopicProcedure topicProcedure = new AlterTopicProcedure();
         topicProcedure.deserialize(byteBuffer);
         alterTopicProcedures.add(topicProcedure);
       }
     }
 
-    // deserialize pipe procedures
+    // Deserialize CreatePipeProcedureV2s
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       size = ReadWriteIOUtils.readInt(byteBuffer);
       for (int i = 0; i < size; ++i) {
         // This readShort should return CREATE_PIPE_PROCEDURE or 
START_PIPE_PROCEDURE.
-        short typeCode = ReadWriteIOUtils.readShort(byteBuffer);
+        final short typeCode = ReadWriteIOUtils.readShort(byteBuffer);
         if (typeCode == ProcedureType.CREATE_PIPE_PROCEDURE_V2.getTypeCode()) {
-          CreatePipeProcedureV2 createPipeProcedureV2 = new 
CreatePipeProcedureV2();
+          final CreatePipeProcedureV2 createPipeProcedureV2 = new 
CreatePipeProcedureV2();
           createPipeProcedureV2.deserialize(byteBuffer);
           createPipeProcedures.add(createPipeProcedureV2);
         }
@@ -406,20 +344,19 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
   }
 
   @Override
-  public boolean equals(Object o) {
+  public boolean equals(final Object o) {
     if (this == o) {
       return true;
     }
     if (o == null || getClass() != o.getClass()) {
       return false;
     }
-    CreateSubscriptionProcedure that = (CreateSubscriptionProcedure) o;
+    final CreateSubscriptionProcedure that = (CreateSubscriptionProcedure) o;
     return Objects.equals(getProcId(), that.getProcId())
         && Objects.equals(getCurrentState(), that.getCurrentState())
         && getCycles() == that.getCycles()
         && Objects.equals(subscribeReq, that.subscribeReq)
         && Objects.equals(alterConsumerGroupProcedure, 
that.alterConsumerGroupProcedure)
-        && Objects.equals(alterTopicProcedures, that.alterTopicProcedures)
         && Objects.equals(createPipeProcedures, that.createPipeProcedures);
   }
 
@@ -431,13 +368,12 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
         getCycles(),
         subscribeReq,
         alterConsumerGroupProcedure,
-        alterTopicProcedures,
         createPipeProcedures);
   }
 
   @TestOnly
   public void setAlterConsumerGroupProcedure(
-      AlterConsumerGroupProcedure alterConsumerGroupProcedure) {
+      final AlterConsumerGroupProcedure alterConsumerGroupProcedure) {
     this.alterConsumerGroupProcedure = alterConsumerGroupProcedure;
   }
 
@@ -447,17 +383,7 @@ public class CreateSubscriptionProcedure extends 
AbstractOperateSubscriptionAndP
   }
 
   @TestOnly
-  public void setAlterTopicProcedures(List<AlterTopicProcedure> 
alterTopicProcedures) {
-    this.alterTopicProcedures = alterTopicProcedures;
-  }
-
-  @TestOnly
-  public List<AlterTopicProcedure> getAlterTopicProcedures() {
-    return this.alterTopicProcedures;
-  }
-
-  @TestOnly
-  public void setCreatePipeProcedures(List<CreatePipeProcedureV2> 
createPipeProcedures) {
+  public void setCreatePipeProcedures(final List<CreatePipeProcedureV2> 
createPipeProcedures) {
     this.createPipeProcedures = createPipeProcedures;
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java
index 1fa775844d2..6741a6c1e2a 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java
@@ -22,13 +22,10 @@ package 
org.apache.iotdb.confignode.procedure.impl.subscription.subscription;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
 import org.apache.iotdb.commons.subscription.meta.consumer.ConsumerGroupMeta;
-import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta;
 import org.apache.iotdb.commons.utils.TestOnly;
 import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
 import 
org.apache.iotdb.confignode.consensus.request.write.pipe.task.DropPipePlanV2;
 import 
org.apache.iotdb.confignode.consensus.request.write.pipe.task.OperateMultiplePipesPlanV2;
-import 
org.apache.iotdb.confignode.consensus.request.write.subscription.topic.AlterMultipleTopicsPlan;
-import 
org.apache.iotdb.confignode.consensus.request.write.subscription.topic.AlterTopicPlan;
 import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
 import 
org.apache.iotdb.confignode.procedure.impl.pipe.AbstractOperatePipeProcedureV2;
 import 
org.apache.iotdb.confignode.procedure.impl.pipe.task.DropPipeProcedureV2;
@@ -55,23 +52,20 @@ import java.util.Objects;
 import java.util.Set;
 import java.util.stream.Collectors;
 
-// TODO: check if it also needs meta sync to keep CN and DN in sync
 public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPipeProcedure {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(DropSubscriptionProcedure.class);
 
   private TUnsubscribeReq unsubscribeReq;
 
-  // NOTE: The 'drop pipe' operation should be performed before 'alter 
consumer group'.
+  // execution order: drop pipe -> alter consumer group
+  // rollback order: alter consumer group -> drop pipe (no-op)
+  // NOTE: The 'drop pipe' operation must be performed before 'alter consumer 
group'.
   private List<DropPipeProcedureV2> dropPipeProcedures = new ArrayList<>();
-  private List<AlterTopicProcedure> alterTopicProcedures = new ArrayList<>();
   private AlterConsumerGroupProcedure alterConsumerGroupProcedure;
 
-  // Record failed index of procedures to rollback properly.
-  // We only record fail index when executing on config nodes, because when 
executing on data nodes
-  // fails, we just push all meta to data nodes.
-  private int alterTopicProcedureFailIndexOnCN = -1;
-  private int dropPipeProcedureFailIndexOnCN = -1;
+  // TODO: remove this variable later
+  private final List<AlterTopicProcedure> alterTopicProcedures = new 
ArrayList<>(); // unused now
 
   public DropSubscriptionProcedure() {
     super();
@@ -107,11 +101,6 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
     for (final String topic : unsubscribeReq.getTopicNames()) {
       if (topicsUnsubByGroup.contains(topic)) {
         // Topic will be subscribed by no consumers in this group
-
-        final TopicMeta updatedTopicMeta = 
subscriptionInfo.get().deepCopyTopicMeta(topic);
-        
updatedTopicMeta.removeSubscribedConsumerGroup(unsubscribeReq.getConsumerGroupId());
-
-        alterTopicProcedures.add(new AlterTopicProcedure(updatedTopicMeta, 
subscriptionInfo));
         dropPipeProcedures.add(
             new DropPipeProcedureV2(
                 PipeStaticMeta.generateSubscriptionPipeName(
@@ -126,11 +115,6 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       dropPipeProcedure.executeFromCalculateInfoForTask(env);
     }
 
-    // Validate AlterTopicProcedures
-    for (final AlterTopicProcedure alterTopicProcedure : alterTopicProcedures) 
{
-      alterTopicProcedure.executeFromValidate(env);
-    }
-
     // Validate AlterConsumerGroupProcedure
     alterConsumerGroupProcedure.executeFromValidate(env);
     return true;
@@ -141,13 +125,12 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       throws SubscriptionException {
     LOGGER.info("DropSubscriptionProcedure: executeFromOperateOnConfigNodes");
 
-    TSStatus response;
-
     // Execute DropPipeProcedureV2s
     final List<ConfigPhysicalPlan> dropPipePlans =
         dropPipeProcedures.stream()
             .map(proc -> new DropPipePlanV2(proc.getPipeName()))
             .collect(Collectors.toList());
+    TSStatus response;
     try {
       response =
           env.getConfigManager()
@@ -160,30 +143,10 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
     }
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()
         && response.getSubStatusSize() > 0) {
-      // Record the failed index for rollback
-      dropPipeProcedureFailIndexOnCN = response.getSubStatusSize() - 1;
-    }
-
-    // Execute AlterTopicProcedures
-    final List<AlterTopicPlan> alterTopicPlans =
-        alterTopicProcedures.stream()
-            .map(AlterTopicProcedure::getUpdatedTopicMeta)
-            .map(AlterTopicPlan::new)
-            .collect(Collectors.toList());
-    try {
-      response =
-          env.getConfigManager()
-              .getConsensusManager()
-              .write(new AlterMultipleTopicsPlan(alterTopicPlans));
-    } catch (final ConsensusException e) {
-      LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
-      response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
-      response.setMessage(e.getMessage());
-    }
-    if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()
-        && response.getSubStatusSize() > 0) {
-      // Record the failed index for rollback
-      alterTopicProcedureFailIndexOnCN = response.getSubStatusSize() - 1;
+      throw new SubscriptionException(
+          String.format(
+              "Failed to drop subscription with request %s on config nodes, 
because %s",
+              unsubscribeReq, response));
     }
 
     // Execute AlterConsumerGroupProcedure
@@ -204,22 +167,11 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
         AbstractOperatePipeProcedureV2.parsePushPipeMetaExceptionForPipe(
             null, dropMultiPipeOnDataNodes(pipeNames, env));
     if (!exceptionMessage.isEmpty()) {
-      // If not all pipe meta are pushed successfully, the meta can be pushed 
during meta sync.
-      LOGGER.warn(
-          "Failed to drop pipes {} when dropping subscription, details: {}, 
metadata will be synchronized later.",
-          pipeNames,
-          exceptionMessage);
-    }
-
-    // Push topic meta to data nodes
-    final List<ByteBuffer> topicMetaBinaryList = new ArrayList<>();
-    for (final AlterTopicProcedure alterTopicProcedure : alterTopicProcedures) 
{
-      
topicMetaBinaryList.add(alterTopicProcedure.getUpdatedTopicMeta().serialize());
-    }
-    if 
(pushTopicMetaHasException(env.pushMultiTopicMetaToDataNodes(topicMetaBinaryList)))
 {
-      // If not all topic meta are pushed successfully, the meta can be pushed 
during meta sync.
-      LOGGER.warn(
-          "Failed to alter topics when creating subscription, metadata will be 
synchronized later.");
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
+      throw new SubscriptionException(
+          String.format(
+              "Failed to drop pipes %s when dropping subscription with request 
%s, because %s",
+              pipeNames, unsubscribeReq, exceptionMessage));
     }
 
     // Push consumer group meta to data nodes
@@ -232,62 +184,25 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
   }
 
   @Override
-  protected void rollbackFromOperateOnConfigNodes(final ConfigNodeProcedureEnv 
env) {
+  protected void rollbackFromOperateOnConfigNodes(final ConfigNodeProcedureEnv 
env)
+      throws SubscriptionException {
     LOGGER.info("DropSubscriptionProcedure: rollbackFromOperateOnConfigNodes");
 
     // Rollback AlterConsumerGroupProcedure
     alterConsumerGroupProcedure.rollbackFromOperateOnConfigNodes(env);
 
-    // Rollback AlterTopicProcedures
-    TSStatus response;
-    final List<AlterTopicPlan> alterTopicRollbackPlans = new ArrayList<>();
-    for (int i = 0;
-        i <= Math.min(alterTopicProcedureFailIndexOnCN, 
alterTopicProcedures.size());
-        i++) {
-      alterTopicRollbackPlans.add(
-          new 
AlterTopicPlan(alterTopicProcedures.get(i).getExistedTopicMeta()));
-    }
-    try {
-      response =
-          env.getConfigManager()
-              .getConsensusManager()
-              .write(new AlterMultipleTopicsPlan(alterTopicRollbackPlans));
-    } catch (final ConsensusException e) {
-      LOGGER.warn("Failed in the write API executing the consensus layer due 
to: ", e);
-      response = new 
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
-      response.setMessage(e.getMessage());
-    }
-    // If failed to rollback, throw exception
-    if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      throw new SubscriptionException(response.getMessage());
-    }
-
     // Do nothing to rollback DropPipeProcedureV2s
   }
 
   @Override
   protected void rollbackFromOperateOnDataNodes(final ConfigNodeProcedureEnv 
env)
-      throws IOException {
+      throws SubscriptionException, IOException {
     LOGGER.info("DropSubscriptionProcedure: rollbackFromOperateOnDataNodes");
 
     // Rollback AlterConsumerGroupProcedure
     alterConsumerGroupProcedure.rollbackFromOperateOnDataNodes(env);
 
-    // Push all topic metas to datanode, may be time-consuming
-    if (pushTopicMetaHasException(pushTopicMetaToDataNodes(env))) {
-      LOGGER.warn(
-          "Failed to rollback alter topics when dropping subscription, 
metadata will be synchronized later.");
-    }
-
-    // Push all pipe metas to datanode, may be time-consuming
-    final String exceptionMessage =
-        AbstractOperatePipeProcedureV2.parsePushPipeMetaExceptionForPipe(
-            null, AbstractOperatePipeProcedureV2.pushPipeMetaToDataNodes(env, 
pipeTaskInfo));
-    if (!exceptionMessage.isEmpty()) {
-      LOGGER.warn(
-          "Failed to rollback create pipes when dropping subscription, 
details: {}, metadata will be synchronized later.",
-          exceptionMessage);
-    }
+    // Do nothing to rollback DropPipeProcedureV2s
   }
 
   @Override
@@ -306,7 +221,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       }
     }
 
-    // serialize consumerGroupProcedure
+    // Serialize AlterConsumerGroupProcedure
     if (alterConsumerGroupProcedure != null) {
       ReadWriteIOUtils.write(true, stream);
       alterConsumerGroupProcedure.serialize(stream);
@@ -314,7 +229,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       ReadWriteIOUtils.write(false, stream);
     }
 
-    // serialize topic procedures
+    // Serialize AlterTopicProcedures
     if (alterTopicProcedures != null) {
       ReadWriteIOUtils.write(true, stream);
       ReadWriteIOUtils.write(alterTopicProcedures.size(), stream);
@@ -325,7 +240,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       ReadWriteIOUtils.write(false, stream);
     }
 
-    // serialize pipe procedures
+    // Serialize DropPipeProcedureV2s
     if (dropPipeProcedures != null) {
       ReadWriteIOUtils.write(true, stream);
       ReadWriteIOUtils.write(dropPipeProcedures.size(), stream);
@@ -351,7 +266,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       
unsubscribeReq.getTopicNames().add(ReadWriteIOUtils.readString(byteBuffer));
     }
 
-    // deserialize consumerGroupProcedure
+    // Deserialize AlterConsumerGroupProcedure
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       // This readShort should return ALTER_CONSUMER_GROUP_PROCEDURE, and we 
ignore it.
       ReadWriteIOUtils.readShort(byteBuffer);
@@ -360,7 +275,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       alterConsumerGroupProcedure.deserialize(byteBuffer);
     }
 
-    // deserialize topic procedures
+    // Deserialize AlterTopicProcedures
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       size = ReadWriteIOUtils.readInt(byteBuffer);
       for (int i = 0; i < size; ++i) {
@@ -373,7 +288,7 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
       }
     }
 
-    // deserialize pipe procedures
+    // Deserialize DropPipeProcedureV2s
     if (ReadWriteIOUtils.readBool(byteBuffer)) {
       size = ReadWriteIOUtils.readInt(byteBuffer);
       for (int i = 0; i < size; ++i) {
@@ -402,7 +317,6 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
         && getCycles() == that.getCycles()
         && Objects.equals(unsubscribeReq, that.unsubscribeReq)
         && Objects.equals(alterConsumerGroupProcedure, 
that.alterConsumerGroupProcedure)
-        && Objects.equals(alterTopicProcedures, that.alterTopicProcedures)
         && Objects.equals(dropPipeProcedures, that.dropPipeProcedures);
   }
 
@@ -414,7 +328,6 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
         getCycles(),
         unsubscribeReq,
         alterConsumerGroupProcedure,
-        alterTopicProcedures,
         dropPipeProcedures);
   }
 
@@ -429,16 +342,6 @@ public class DropSubscriptionProcedure extends 
AbstractOperateSubscriptionAndPip
     return this.alterConsumerGroupProcedure;
   }
 
-  @TestOnly
-  public void setAlterTopicProcedures(final List<AlterTopicProcedure> 
alterTopicProcedures) {
-    this.alterTopicProcedures = alterTopicProcedures;
-  }
-
-  @TestOnly
-  public List<AlterTopicProcedure> getAlterTopicProcedures() {
-    return this.alterTopicProcedures;
-  }
-
   @TestOnly
   public void setDropPipeProcedures(final List<DropPipeProcedureV2> 
dropPipeProcedures) {
     this.dropPipeProcedures = dropPipeProcedures;
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/AlterTopicProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/AlterTopicProcedure.java
index f8fbdab72e0..4faa2cfa0c1 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/AlterTopicProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/AlterTopicProcedure.java
@@ -108,7 +108,6 @@ public class AlterTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
       response =
           new 
TSStatus(TSStatusCode.ALTER_TOPIC_ERROR.getStatusCode()).setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
@@ -119,25 +118,18 @@ public class AlterTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
 
   @Override
   public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
-      throws SubscriptionException {
+      throws SubscriptionException, IOException {
     LOGGER.info(
         "AlterTopicProcedure: executeFromOperateOnDataNodes({})", 
updatedTopicMeta.getTopicName());
 
-    try {
-      final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(updatedTopicMeta.serialize());
-      if (RpcUtils.squashResponseStatusList(statuses).getCode()
-          != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        throw new SubscriptionException(
-            String.format(
-                "Failed to alter topic (%s -> %s) on data nodes, because %s",
-                existedTopicMeta, updatedTopicMeta, statuses));
-      }
-    } catch (IOException e) {
-      LOGGER.warn("Failed to serialize the topic meta due to: ", e);
+    final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(updatedTopicMeta.serialize());
+    if (RpcUtils.squashResponseStatusList(statuses).getCode()
+        != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
               "Failed to alter topic (%s -> %s) on data nodes, because %s",
-              existedTopicMeta, updatedTopicMeta, e.getMessage()));
+              existedTopicMeta, updatedTopicMeta, statuses));
     }
   }
 
@@ -147,7 +139,8 @@ public class AlterTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
   }
 
   @Override
-  public void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  public void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info(
         "AlterTopicProcedure: rollbackFromOperateOnConfigNodes({})",
         updatedTopicMeta.getTopicName());
@@ -161,7 +154,6 @@ public class AlterTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
       response =
           new 
TSStatus(TSStatusCode.ALTER_TOPIC_ERROR.getStatusCode()).setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
@@ -171,25 +163,19 @@ public class AlterTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
   }
 
   @Override
-  public void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
+  public void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException, IOException {
     LOGGER.info(
         "AlterTopicProcedure: rollbackFromOperateOnDataNodes({})", 
updatedTopicMeta.getTopicName());
 
-    try {
-      final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(existedTopicMeta.serialize());
-      if (RpcUtils.squashResponseStatusList(statuses).getCode()
-          != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        throw new SubscriptionException(
-            String.format(
-                "Failed to rollback from altering topic (%s -> %s) on data 
nodes, because %s",
-                updatedTopicMeta, existedTopicMeta, statuses));
-      }
-    } catch (IOException e) {
-      LOGGER.warn("Failed to serialize the topic meta due to: ", e);
+    final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(existedTopicMeta.serialize());
+    if (RpcUtils.squashResponseStatusList(statuses).getCode()
+        != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
               "Failed to rollback from altering topic (%s -> %s) on data 
nodes, because %s",
-              updatedTopicMeta, existedTopicMeta, e.getMessage()));
+              updatedTopicMeta, existedTopicMeta, statuses));
     }
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/CreateTopicProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/CreateTopicProcedure.java
index 65035f996bf..c27205290a8 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/CreateTopicProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/CreateTopicProcedure.java
@@ -96,7 +96,6 @@ public class CreateTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
       response =
           new 
TSStatus(TSStatusCode.CREATE_TOPIC_ERROR.getStatusCode()).setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
@@ -106,22 +105,16 @@ public class CreateTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
 
   @Override
   protected void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
-      throws SubscriptionException {
+      throws SubscriptionException, IOException {
     LOGGER.info("CreateTopicProcedure: executeFromOperateOnDataNodes({})", 
topicMeta);
 
-    try {
-      final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(topicMeta.serialize());
-      if (RpcUtils.squashResponseStatusList(statuses).getCode()
-          != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        throw new SubscriptionException(
-            String.format(
-                "Failed to create topic %s on data nodes, because %s", 
topicMeta, statuses));
-      }
-    } catch (IOException e) {
-      LOGGER.warn("Failed to serialize the topic meta due to: ", e);
+    final List<TSStatus> statuses = 
env.pushSingleTopicOnDataNode(topicMeta.serialize());
+    if (RpcUtils.squashResponseStatusList(statuses).getCode()
+        != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
-              "Failed to create topic %s on data nodes, because %s", 
topicMeta, e.getMessage()));
+              "Failed to create topic %s on data nodes, because %s", 
topicMeta, statuses));
     }
   }
 
@@ -131,7 +124,8 @@ public class CreateTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
   }
 
   @Override
-  protected void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  protected void rollbackFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info("CreateTopicProcedure: rollbackFromCreateOnConfigNodes({})", 
topicMeta);
 
     TSStatus response;
@@ -145,24 +139,27 @@ public class CreateTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
       response =
           new 
TSStatus(TSStatusCode.DROP_TOPIC_ERROR.getStatusCode()).setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
-              "Failed to rollback topic %s on config nodes, because %s", 
topicMeta, response));
+              "Failed to rollback creating topic %s on config nodes, because 
%s",
+              topicMeta, response));
     }
   }
 
   @Override
-  protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
+  protected void rollbackFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info("CreateTopicProcedure: rollbackFromCreateOnDataNodes({})", 
topicMeta);
 
     final List<TSStatus> statuses = 
env.dropSingleTopicOnDataNode(topicMeta.getTopicName());
     if (RpcUtils.squashResponseStatusList(statuses).getCode()
         != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format(
-              "Failed to rollback topic %s on data nodes, because %s", 
topicMeta, statuses));
+              "Failed to rollback creating topic %s on data nodes, because %s",
+              topicMeta, statuses));
     }
   }
 
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/DropTopicProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/DropTopicProcedure.java
index f1cfbb59d10..363e5716ee8 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/DropTopicProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/DropTopicProcedure.java
@@ -81,7 +81,6 @@ public class DropTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
       response =
           new 
TSStatus(TSStatusCode.DROP_TOPIC_ERROR.getStatusCode()).setMessage(e.getMessage());
     }
-
     if (response.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
       throw new SubscriptionException(
           String.format(
@@ -90,13 +89,13 @@ public class DropTopicProcedure extends 
AbstractOperateSubscriptionProcedure {
   }
 
   @Override
-  protected void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
-      throws SubscriptionException {
+  protected void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) {
     LOGGER.info("DropTopicProcedure: executeFromOperateOnDataNodes({})", 
topicName);
 
     final List<TSStatus> statuses = env.dropSingleTopicOnDataNode(topicName);
     if (RpcUtils.squashResponseStatusList(statuses).getCode()
         != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+      // throw exception instead of logging warn, do not rely on metadata 
synchronization
       throw new SubscriptionException(
           String.format("Failed to drop topic %s on data nodes, because %s", 
topicName, statuses));
     }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/runtime/TopicMetaSyncProcedure.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/runtime/TopicMetaSyncProcedure.java
index 40920e43936..27919bbcb9e 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/runtime/TopicMetaSyncProcedure.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/topic/runtime/TopicMetaSyncProcedure.java
@@ -98,7 +98,8 @@ public class TopicMetaSyncProcedure extends 
AbstractOperateSubscriptionProcedure
   }
 
   @Override
-  public void executeFromOperateOnConfigNodes(ConfigNodeProcedureEnv env) {
+  public void executeFromOperateOnConfigNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException {
     LOGGER.info("TopicMetaSyncProcedure: executeFromOperateOnConfigNodes");
 
     final List<TopicMeta> topicMetaList = new ArrayList<>();
@@ -121,7 +122,8 @@ public class TopicMetaSyncProcedure extends 
AbstractOperateSubscriptionProcedure
   }
 
   @Override
-  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) throws 
IOException {
+  public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env)
+      throws SubscriptionException, IOException {
     LOGGER.info("TopicMetaSyncProcedure: executeFromOperateOnDataNodes");
 
     Map<Integer, TPushTopicMetaResp> respMap = pushTopicMetaToDataNodes(env);
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedureTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedureTest.java
index 3305e58d831..93d9941fbf3 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedureTest.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedureTest.java
@@ -21,10 +21,8 @@ package 
org.apache.iotdb.confignode.procedure.impl.subscription.subscription;
 
 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.confignode.procedure.impl.pipe.task.CreatePipeProcedureV2;
 import 
org.apache.iotdb.confignode.procedure.impl.subscription.consumer.AlterConsumerGroupProcedure;
-import 
org.apache.iotdb.confignode.procedure.impl.subscription.topic.AlterTopicProcedure;
 import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
 import org.apache.iotdb.confignode.rpc.thrift.TCreatePipeReq;
 import org.apache.iotdb.confignode.rpc.thrift.TSubscribeReq;
@@ -74,11 +72,6 @@ public class CreateSubscriptionProcedureTest {
     AlterConsumerGroupProcedure alterConsumerGroupProcedure =
         new AlterConsumerGroupProcedure(newConsumerGroupMeta);
 
-    List<AlterTopicProcedure> topicProcedures = new ArrayList<>();
-    TopicMeta newTopicMeta = new TopicMeta("t1", 1, topicAttributes);
-    newTopicMeta.addSubscribedConsumerGroup("cg1");
-    topicProcedures.add(new AlterTopicProcedure(newTopicMeta));
-
     List<CreatePipeProcedureV2> pipeProcedures = new ArrayList<>();
     pipeProcedures.add(
         new CreatePipeProcedureV2(
@@ -92,7 +85,6 @@ public class CreateSubscriptionProcedureTest {
                 .setProcessorAttributes(Collections.singletonMap("processor", 
"pro"))));
 
     proc.setAlterConsumerGroupProcedure(alterConsumerGroupProcedure);
-    proc.setAlterTopicProcedures(topicProcedures);
     proc.setCreatePipeProcedures(pipeProcedures);
 
     try {
@@ -104,7 +96,6 @@ public class CreateSubscriptionProcedureTest {
 
       assertEquals(proc, proc2);
       assertEquals(alterConsumerGroupProcedure, 
proc2.getAlterConsumerGroupProcedure());
-      assertEquals(topicProcedures, proc2.getAlterTopicProcedures());
       assertEquals(pipeProcedures, proc2.getCreatePipeProcedures());
     } catch (Exception e) {
       e.printStackTrace();
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedureTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedureTest.java
index 9519bf2ed15..9ecce2a522c 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedureTest.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedureTest.java
@@ -21,10 +21,8 @@ package 
org.apache.iotdb.confignode.procedure.impl.subscription.subscription;
 
 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.confignode.procedure.impl.pipe.task.DropPipeProcedureV2;
 import 
org.apache.iotdb.confignode.procedure.impl.subscription.consumer.AlterConsumerGroupProcedure;
-import 
org.apache.iotdb.confignode.procedure.impl.subscription.topic.AlterTopicProcedure;
 import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
 import org.apache.iotdb.confignode.rpc.thrift.TUnsubscribeReq;
 
@@ -72,16 +70,11 @@ public class DropSubscriptionProcedureTest {
     AlterConsumerGroupProcedure alterConsumerGroupProcedure =
         new AlterConsumerGroupProcedure(newConsumerGroupMeta);
 
-    List<AlterTopicProcedure> topicProcedures = new ArrayList<>();
-    topicProcedures.add(new AlterTopicProcedure(new TopicMeta("t1", 1, 
topicAttributes)));
-    topicProcedures.add(new AlterTopicProcedure(new TopicMeta("t2", 2, 
topicAttributes)));
-
     List<DropPipeProcedureV2> pipeProcedures = new ArrayList<>();
     pipeProcedures.add(new DropPipeProcedureV2("pipe_topic1"));
     pipeProcedures.add(new DropPipeProcedureV2("pipe_topic2"));
 
     proc.setAlterConsumerGroupProcedure(alterConsumerGroupProcedure);
-    proc.setAlterTopicProcedures(topicProcedures);
     proc.setDropPipeProcedures(pipeProcedures);
 
     try {
@@ -93,7 +86,6 @@ public class DropSubscriptionProcedureTest {
 
       assertEquals(proc, proc2);
       assertEquals(alterConsumerGroupProcedure, 
proc2.getAlterConsumerGroupProcedure());
-      assertEquals(topicProcedures, proc2.getAlterTopicProcedures());
       assertEquals(pipeProcedures, proc2.getDropPipeProcedures());
     } catch (Exception e) {
       fail();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java
index 13770fdbf7c..ee784fa76c4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java
@@ -19,7 +19,6 @@
 
 package org.apache.iotdb.db.subscription.agent;
 
-import 
org.apache.iotdb.commons.subscription.meta.consumer.ConsumerGroupMetaKeeper;
 import org.apache.iotdb.db.subscription.broker.SubscriptionBroker;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import 
org.apache.iotdb.db.subscription.task.subtask.SubscriptionConnectorSubtask;
@@ -35,6 +34,7 @@ import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 public class SubscriptionBrokerAgent {
 
@@ -117,61 +117,61 @@ public class SubscriptionBrokerAgent {
 
   /////////////////////////////// broker ///////////////////////////////
 
-  /**
-   * Caller should ensure that the method is called in the lock {@link
-   * ConsumerGroupMetaKeeper#acquireWriteLock}.
-   */
   public boolean isBrokerExist(final String consumerGroupId) {
     return consumerGroupIdToSubscriptionBroker.containsKey(consumerGroupId);
   }
 
-  /**
-   * Caller should ensure that the method is called in the lock {@link
-   * ConsumerGroupMetaKeeper#acquireWriteLock}.
-   */
-  public void createBroker(final String consumerGroupId) {
-    final SubscriptionBroker broker = new SubscriptionBroker(consumerGroupId);
-    consumerGroupIdToSubscriptionBroker.put(consumerGroupId, broker);
+  public void createBrokerIfNotExist(final String consumerGroupId) {
+    consumerGroupIdToSubscriptionBroker.computeIfAbsent(consumerGroupId, 
SubscriptionBroker::new);
     LOGGER.info("Subscription: create broker bound to consumer group [{}]", 
consumerGroupId);
   }
 
   /**
-   * Caller should ensure that the method is called in the lock {@link
-   * ConsumerGroupMetaKeeper#acquireWriteLock}.
-   *
    * @return {@code true} if drop broker success, {@code false} otherwise
    */
   public boolean dropBroker(final String consumerGroupId) {
-    final SubscriptionBroker broker = 
consumerGroupIdToSubscriptionBroker.get(consumerGroupId);
-    if (Objects.isNull(broker)) {
-      LOGGER.warn(
-          "Subscription: broker bound to consumer group [{}] does not exist", 
consumerGroupId);
-      // do nothing
-      return true;
-    }
-    if (!broker.isEmpty()) {
-      LOGGER.warn(
-          "Subscription: broker bound to consumer group [{}] is not empty when 
dropping",
-          consumerGroupId);
-      // do nothing
-      return false;
-    }
-    consumerGroupIdToSubscriptionBroker.remove(consumerGroupId);
-    LOGGER.info("Subscription: drop broker bound to consumer group [{}]", 
consumerGroupId);
-    return true;
+    final AtomicBoolean dropped = new AtomicBoolean(false);
+    consumerGroupIdToSubscriptionBroker.compute(
+        consumerGroupId,
+        (id, broker) -> {
+          if (Objects.isNull(broker)) {
+            LOGGER.warn(
+                "Subscription: broker bound to consumer group [{}] does not 
exist",
+                consumerGroupId);
+            dropped.set(true);
+            return null;
+          }
+          if (!broker.isEmpty()) {
+            LOGGER.warn(
+                "Subscription: broker bound to consumer group [{}] is not 
empty when dropping",
+                consumerGroupId);
+            return broker;
+          }
+          dropped.set(true);
+          LOGGER.info("Subscription: drop broker bound to consumer group 
[{}]", consumerGroupId);
+          return null; // remove this entry
+        });
+    return dropped.get();
   }
 
   /////////////////////////////// prefetching queue 
///////////////////////////////
 
   public void bindPrefetchingQueue(final SubscriptionConnectorSubtask subtask) 
{
     final String consumerGroupId = subtask.getConsumerGroupId();
-    final SubscriptionBroker broker = 
consumerGroupIdToSubscriptionBroker.get(consumerGroupId);
-    if (Objects.isNull(broker)) {
-      LOGGER.warn(
-          "Subscription: broker bound to consumer group [{}] does not exist", 
consumerGroupId);
-      return;
-    }
-    broker.bindPrefetchingQueue(subtask.getTopicName(), 
subtask.getInputPendingQueue());
+    consumerGroupIdToSubscriptionBroker
+        .compute(
+            consumerGroupId,
+            (id, broker) -> {
+              if (Objects.isNull(broker)) {
+                LOGGER.info(
+                    "Subscription: broker bound to consumer group [{}] does 
not exist, create new for binding prefetching queue",
+                    consumerGroupId);
+                // TODO: consider more robust metadata semantics
+                return new SubscriptionBroker(consumerGroupId);
+              }
+              return broker;
+            })
+        .bindPrefetchingQueue(subtask.getTopicName(), 
subtask.getInputPendingQueue());
   }
 
   public void unbindPrefetchingQueue(final String consumerGroupId, final 
String topicName) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionConsumerAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionConsumerAgent.java
index 0105700d3c8..fee23cf6af4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionConsumerAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionConsumerAgent.java
@@ -94,14 +94,23 @@ public class SubscriptionConsumerAgent {
     final ConsumerGroupMeta metaInAgent =
         consumerGroupMetaKeeper.getConsumerGroupMeta(consumerGroupId);
 
-    // if consumer group meta does not exist on local agent or creation time 
is inconsistent with
-    // meta from coordinator
-    if (Objects.isNull(metaInAgent)
-        || metaInAgent.getCreationTime() != 
metaFromCoordinator.getCreationTime()) {
+    // if consumer group meta does not exist on local agent
+    if (Objects.isNull(metaInAgent)) {
+      consumerGroupMetaKeeper.removeConsumerGroupMeta(consumerGroupId);
+      consumerGroupMetaKeeper.addConsumerGroupMeta(consumerGroupId, 
metaFromCoordinator);
+      SubscriptionAgent.broker().createBrokerIfNotExist(consumerGroupId);
+      return;
+    }
+
+    // if the creation time of consumer group meta on local agent is 
inconsistent with meta from
+    // coordinator
+    if (metaInAgent.getCreationTime() != 
metaFromCoordinator.getCreationTime()) {
       if (SubscriptionAgent.broker().isBrokerExist(consumerGroupId)) {
         LOGGER.warn(
-            "Subscription: broker bound to consumer group [{}] has already 
existed when the corresponding consumer group meta does not exist on local 
agent, drop it",
-            consumerGroupId);
+            "Subscription: broker bound to consumer group [{}] has already 
existed when the creation time of consumer group meta on local agent {} is 
inconsistent with meta from coordinator {}, drop it",
+            consumerGroupId,
+            metaInAgent,
+            metaFromCoordinator);
         if (!SubscriptionAgent.broker().dropBroker(consumerGroupId)) {
           final String exceptionMessage =
               String.format(
@@ -113,11 +122,11 @@ public class SubscriptionConsumerAgent {
 
       consumerGroupMetaKeeper.removeConsumerGroupMeta(consumerGroupId);
       consumerGroupMetaKeeper.addConsumerGroupMeta(consumerGroupId, 
metaFromCoordinator);
-      SubscriptionAgent.broker().createBroker(consumerGroupId);
+      // no need to create broker manually
       return;
     }
 
-    // remove prefetching queue
+    // remove prefetching queues for topics unsubscribed by the consumer group
     final Set<String> topicsUnsubByGroup =
         ConsumerGroupMeta.getTopicsUnsubByGroup(metaInAgent, 
metaFromCoordinator);
     for (final String topicName : topicsUnsubByGroup) {
@@ -125,7 +134,7 @@ public class SubscriptionConsumerAgent {
     }
 
     // TODO: Currently we fully replace the entire ConsumerGroupMeta without 
carefully checking the
-    // changes in its fields.
+    //       changes in its fields.
     consumerGroupMetaKeeper.removeConsumerGroupMeta(consumerGroupId);
     consumerGroupMetaKeeper.addConsumerGroupMeta(consumerGroupId, 
metaFromCoordinator);
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionBroker.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionBroker.java
index 420b571103a..afc8f2f2290 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionBroker.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionBroker.java
@@ -117,8 +117,9 @@ public class SubscriptionBroker {
         }
         // There are two reasons for not printing logs here:
         // 1. There will be a delay in the creation of the prefetching queue 
after subscription.
-        // 2. There is no corresponding prefetching queue on this DN 
(currently the consumer is
-        // fully connected to all DNs).
+        // 2. There is no corresponding prefetching queue on this DN:
+        //   2.1. the consumer is fully connected to all DNs currently...
+        //   2.2. potential disorder of unbind and close prefetching queue...
         continue;
       }
       if (prefetchingQueue.isClosed()) {
@@ -338,11 +339,12 @@ public class SubscriptionBroker {
     final SubscriptionPrefetchingQueue prefetchingQueue =
         topicNameToPrefetchingQueue.get(topicName);
     if (Objects.nonNull(prefetchingQueue)) {
-      LOGGER.warn(
-          "Subscription: prefetching queue bound to topic [{}] for consumer 
group [{}] still exists",
+      LOGGER.info(
+          "Subscription: prefetching queue bound to topic [{}] for consumer 
group [{}] still exists, unbind it before closing",
           topicName,
           brokerId);
-      return;
+      // TODO: consider more robust metadata semantics
+      unbindPrefetchingQueue(topicName);
     }
 
     completedTopicNames.remove(topicName);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
index b0f67d6493e..936e8777a22 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
@@ -291,6 +291,9 @@ public class CommonConfig {
   private long subscriptionReadTabletBufferSize = 8 * MB;
   private long subscriptionTsFileDeduplicationWindowSeconds = 120; // 120s
 
+  private long subscriptionMetaSyncerInitialSyncDelayMinutes = 3;
+  private long subscriptionMetaSyncerSyncIntervalMinutes = 3;
+
   /** Whether to use persistent schema mode. */
   private String schemaEngineMode = "Memory";
 
@@ -1322,6 +1325,25 @@ public class CommonConfig {
         subscriptionTsFileDeduplicationWindowSeconds;
   }
 
+  public long getSubscriptionMetaSyncerInitialSyncDelayMinutes() {
+    return subscriptionMetaSyncerInitialSyncDelayMinutes;
+  }
+
+  public void setSubscriptionMetaSyncerInitialSyncDelayMinutes(
+      long subscriptionMetaSyncerInitialSyncDelayMinutes) {
+    this.subscriptionMetaSyncerInitialSyncDelayMinutes =
+        subscriptionMetaSyncerInitialSyncDelayMinutes;
+  }
+
+  public long getSubscriptionMetaSyncerSyncIntervalMinutes() {
+    return subscriptionMetaSyncerSyncIntervalMinutes;
+  }
+
+  public void setSubscriptionMetaSyncerSyncIntervalMinutes(
+      long subscriptionMetaSyncerSyncIntervalMinutes) {
+    this.subscriptionMetaSyncerSyncIntervalMinutes = 
subscriptionMetaSyncerSyncIntervalMinutes;
+  }
+
   public String getSchemaEngineMode() {
     return schemaEngineMode;
   }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
index cccee5e7a23..4ebfb6ef307 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
@@ -623,6 +623,7 @@ public class CommonDescriptor {
             properties.getProperty(
                 "subscription_cache_memory_usage_percentage",
                 
String.valueOf(config.getSubscriptionCacheMemoryUsagePercentage()))));
+
     config.setSubscriptionSubtaskExecutorMaxThreadNum(
         Integer.parseInt(
             properties.getProperty(
@@ -686,6 +687,17 @@ public class CommonDescriptor {
             properties.getProperty(
                 "subscription_ts_file_deduplication_window_seconds",
                 
String.valueOf(config.getSubscriptionTsFileDeduplicationWindowSeconds()))));
+
+    config.setSubscriptionMetaSyncerInitialSyncDelayMinutes(
+        Long.parseLong(
+            properties.getProperty(
+                "subscription_meta_syncer_initial_sync_delay_minutes",
+                
String.valueOf(config.getSubscriptionMetaSyncerInitialSyncDelayMinutes()))));
+    config.setSubscriptionMetaSyncerSyncIntervalMinutes(
+        Long.parseLong(
+            properties.getProperty(
+                "subscription_meta_syncer_sync_interval_minutes",
+                
String.valueOf(config.getSubscriptionMetaSyncerSyncIntervalMinutes()))));
   }
 
   public void loadRetryProperties(Properties properties) throws IOException {
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
index 8a894bc5df5..6a930e29ff6 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
@@ -29,8 +29,6 @@ public class SubscriptionConfig {
 
   private static final CommonConfig COMMON_CONFIG = 
CommonDescriptor.getInstance().getConfig();
 
-  /////////////////////////////// Subtask Executor 
///////////////////////////////
-
   public float getSubscriptionCacheMemoryUsagePercentage() {
     return COMMON_CONFIG.getSubscriptionCacheMemoryUsagePercentage();
   }
@@ -83,6 +81,14 @@ public class SubscriptionConfig {
     return COMMON_CONFIG.getSubscriptionTsFileDeduplicationWindowSeconds();
   }
 
+  public long getSubscriptionMetaSyncerInitialSyncDelayMinutes() {
+    return COMMON_CONFIG.getSubscriptionMetaSyncerInitialSyncDelayMinutes();
+  }
+
+  public long getSubscriptionMetaSyncerSyncIntervalMinutes() {
+    return COMMON_CONFIG.getSubscriptionMetaSyncerSyncIntervalMinutes();
+  }
+
   /////////////////////////////// Utils ///////////////////////////////
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(SubscriptionConfig.class);
@@ -90,6 +96,7 @@ public class SubscriptionConfig {
   public void printAllConfigs() {
     LOGGER.info(
         "SubscriptionCacheMemoryUsagePercentage: {}", 
getSubscriptionCacheMemoryUsagePercentage());
+
     LOGGER.info(
         "SubscriptionSubtaskExecutorMaxThreadNum: {}",
         getSubscriptionSubtaskExecutorMaxThreadNum());
@@ -117,6 +124,13 @@ public class SubscriptionConfig {
     LOGGER.info(
         "SubscriptionTsFileDeduplicationWindowSeconds: {}",
         getSubscriptionTsFileDeduplicationWindowSeconds());
+
+    LOGGER.info(
+        "SubscriptionMetaSyncerInitialSyncDelayMinutes: {}",
+        getSubscriptionMetaSyncerInitialSyncDelayMinutes());
+    LOGGER.info(
+        "SubscriptionMetaSyncerSyncIntervalMinutes: {}",
+        getSubscriptionMetaSyncerSyncIntervalMinutes());
   }
 
   /////////////////////////////// Singleton ///////////////////////////////
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMeta.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMeta.java
index 764fbfca501..e902368cc2d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMeta.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMeta.java
@@ -149,6 +149,10 @@ public class ConsumerGroupMeta {
     return topics;
   }
 
+  public Set<String> getTopicsSubscribedByConsumerGroup() {
+    return topicNameToSubscribedConsumerIdSet.keySet();
+  }
+
   public void addSubscription(final String consumerId, final Set<String> 
topics) {
     if (!consumerIdToConsumerMeta.containsKey(consumerId)) {
       throw new SubscriptionException(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMetaKeeper.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMetaKeeper.java
index c6556076698..84166834570 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMetaKeeper.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/consumer/ConsumerGroupMetaKeeper.java
@@ -26,10 +26,12 @@ import java.io.FileOutputStream;
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Map;
+import java.util.Map.Entry;
 import java.util.Objects;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
+import java.util.stream.Collectors;
 
 public class ConsumerGroupMetaKeeper {
 
@@ -113,6 +115,29 @@ public class ConsumerGroupMetaKeeper {
     return consumerGroupIdToConsumerGroupMetaMap.isEmpty();
   }
 
+  /////////////////////////////////  TopicMeta  
/////////////////////////////////
+
+  public Set<String> getSubscribedConsumerGroupIds(final String topicName) {
+    return consumerGroupIdToConsumerGroupMetaMap.entrySet().stream()
+        .filter(entry -> 
entry.getValue().getTopicsSubscribedByConsumerGroup().contains(topicName))
+        .map(Entry::getKey)
+        .collect(Collectors.toSet());
+  }
+
+  public boolean isTopicSubscribedByConsumerGroup(
+      final String topicName, final String consumerGroupId) {
+    return consumerGroupIdToConsumerGroupMetaMap.containsKey(consumerGroupId)
+        && consumerGroupIdToConsumerGroupMetaMap
+            .get(consumerGroupId)
+            .getTopicsSubscribedByConsumerGroup()
+            .contains(topicName);
+  }
+
+  public boolean isTopicSubscribedByConsumerGroup(final String topicName) {
+    return consumerGroupIdToConsumerGroupMetaMap.values().stream()
+        .anyMatch(meta -> 
meta.getTopicsSubscribedByConsumerGroup().contains(topicName));
+  }
+
   /////////////////////////////////  Snapshot  
/////////////////////////////////
 
   public void processTakeSnapshot(FileOutputStream fileOutputStream) throws 
IOException {
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/topic/TopicMeta.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/topic/TopicMeta.java
index 7d328fbc135..3e409a9160e 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/topic/TopicMeta.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/meta/topic/TopicMeta.java
@@ -20,6 +20,7 @@
 package org.apache.iotdb.commons.subscription.meta.topic;
 
 import org.apache.iotdb.commons.pipe.config.constant.PipeConnectorConstant;
+import org.apache.iotdb.commons.utils.TestOnly;
 import org.apache.iotdb.rpc.subscription.config.TopicConfig;
 
 import org.apache.tsfile.utils.PublicBAOS;
@@ -42,7 +43,8 @@ public class TopicMeta {
   private long creationTime; // unit in ms
   private TopicConfig config;
 
-  private Set<String> subscribedConsumerGroupIds;
+  // TODO: remove this variable later
+  private Set<String> subscribedConsumerGroupIds; // unused now
 
   private TopicMeta() {
     this.config = new TopicConfig(new HashMap<>());
@@ -84,22 +86,27 @@ public class TopicMeta {
   /**
    * @return true if the consumer group did not already subscribe this topic
    */
+  @TestOnly
   public boolean addSubscribedConsumerGroup(final String consumerGroupId) {
     return subscribedConsumerGroupIds.add(consumerGroupId);
   }
 
+  @TestOnly
   public void removeSubscribedConsumerGroup(final String consumerGroupId) {
     subscribedConsumerGroupIds.remove(consumerGroupId);
   }
 
+  @TestOnly
   public Set<String> getSubscribedConsumerGroupIds() {
     return subscribedConsumerGroupIds;
   }
 
+  @TestOnly
   public boolean isSubscribedByConsumerGroup(final String consumerGroupId) {
     return subscribedConsumerGroupIds.contains(consumerGroupId);
   }
 
+  @TestOnly
   public boolean hasSubscribedConsumerGroup() {
     return !subscribedConsumerGroupIds.isEmpty();
   }
@@ -218,13 +225,12 @@ public class TopicMeta {
     final TopicMeta that = (TopicMeta) obj;
     return creationTime == that.creationTime
         && Objects.equals(topicName, that.topicName)
-        && Objects.equals(config, that.config)
-        && Objects.equals(subscribedConsumerGroupIds, 
that.subscribedConsumerGroupIds);
+        && Objects.equals(config, that.config);
   }
 
   @Override
   public int hashCode() {
-    return Objects.hash(topicName, creationTime, subscribedConsumerGroupIds, 
config);
+    return Objects.hash(topicName, creationTime, config);
   }
 
   @Override
@@ -232,13 +238,10 @@ public class TopicMeta {
     return "TopicMeta{"
         + "topicName='"
         + topicName
-        + '\''
-        + ", creationTime="
+        + "', creationTime="
         + creationTime
         + ", config="
         + config
-        + ", subscribedConsumerGroupIds="
-        + subscribedConsumerGroupIds
         + '}';
   }
 }


Reply via email to