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
+ '}';
}
}