From 7b481ce9e4c06fb57394856d420d093e6ba0cc13 Mon Sep 17 00:00:00 2001 From: ccwss <1782935682@qq.com> Date: Wed, 16 Jul 2025 11:51:41 +0800 Subject: [PATCH] fix sync topic and group if attributes isn't null --- .../subscription/SubscriptionGroupManager.java | 14 +++++++++----- .../rocketmq/broker/topic/TopicConfigManager.java | 14 +++++++++----- .../config/v2/SubscriptionGroupManagerV2Test.java | 2 ++ 3 files changed, 20 insertions(+), 10 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/subscription/SubscriptionGroupManager.java b/broker/src/main/java/org/apache/rocketmq/broker/subscription/SubscriptionGroupManager.java index c7083365be4..0ee34e67fad 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/subscription/SubscriptionGroupManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/subscription/SubscriptionGroupManager.java @@ -47,6 +47,7 @@ import org.apache.rocketmq.remoting.protocol.DataVersion; import org.apache.rocketmq.remoting.protocol.RemotingSerializable; import org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig; +import org.apache.rocketmq.store.config.BrokerRole; public class SubscriptionGroupManager extends ConfigManager { protected static final Logger log = LoggerFactory.getLogger(LoggerName.BROKER_LOGGER_NAME); @@ -155,11 +156,14 @@ public void updateSubscriptionGroupConfigWithoutPersist(SubscriptionGroupConfig Map newAttributes = request(config); Map currentAttributes = current(config.getGroupName()); - Map finalAttributes = AttributeUtil.alterCurrentAttributes( - this.subscriptionGroupTable.get(config.getGroupName()) == null, - SubscriptionGroupAttributes.ALL, - ImmutableMap.copyOf(currentAttributes), - ImmutableMap.copyOf(newAttributes)); + Map finalAttributes = newAttributes; + if (this.brokerController.getMessageStoreConfig().getBrokerRole() != BrokerRole.SLAVE) { + finalAttributes = AttributeUtil.alterCurrentAttributes( + this.subscriptionGroupTable.get(config.getGroupName()) == null, + SubscriptionGroupAttributes.ALL, + ImmutableMap.copyOf(currentAttributes), + ImmutableMap.copyOf(newAttributes)); + } config.setAttributes(finalAttributes); diff --git a/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicConfigManager.java b/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicConfigManager.java index ed46dfdc49c..73b3958606a 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicConfigManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/topic/TopicConfigManager.java @@ -53,6 +53,7 @@ import org.apache.rocketmq.remoting.protocol.body.TopicConfigAndMappingSerializeWrapper; import org.apache.rocketmq.remoting.protocol.body.TopicConfigSerializeWrapper; import org.apache.rocketmq.remoting.protocol.statictopic.TopicQueueMappingInfo; +import org.apache.rocketmq.store.config.BrokerRole; import org.apache.rocketmq.store.timer.TimerMessageStore; import org.apache.rocketmq.tieredstore.TieredMessageStore; import org.apache.rocketmq.tieredstore.metadata.MetadataStore; @@ -504,11 +505,14 @@ public void updateSingleTopicConfigWithoutPersist(final TopicConfig topicConfig) Map newAttributes = request(topicConfig); Map currentAttributes = current(topicConfig.getTopicName()); - Map finalAttributes = AttributeUtil.alterCurrentAttributes( - this.topicConfigTable.get(topicConfig.getTopicName()) == null, - TopicAttributes.ALL, - ImmutableMap.copyOf(currentAttributes), - ImmutableMap.copyOf(newAttributes)); + Map finalAttributes = newAttributes; + if (this.brokerController.getMessageStoreConfig().getBrokerRole() != BrokerRole.SLAVE) { + finalAttributes = AttributeUtil.alterCurrentAttributes( + this.topicConfigTable.get(topicConfig.getTopicName()) == null, + TopicAttributes.ALL, + ImmutableMap.copyOf(currentAttributes), + ImmutableMap.copyOf(newAttributes)); + } topicConfig.setAttributes(finalAttributes); updateTieredStoreTopicMetadata(topicConfig, newAttributes); diff --git a/broker/src/test/java/org/apache/rocketmq/broker/config/v2/SubscriptionGroupManagerV2Test.java b/broker/src/test/java/org/apache/rocketmq/broker/config/v2/SubscriptionGroupManagerV2Test.java index 4ff8a81e60a..df46dd6aaec 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/config/v2/SubscriptionGroupManagerV2Test.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/config/v2/SubscriptionGroupManagerV2Test.java @@ -74,6 +74,8 @@ public void setUp() throws IOException { File configStoreDir = tf.newFolder(); messageStoreConfig = new MessageStoreConfig(); messageStoreConfig.setStorePathRootDir(configStoreDir.getAbsolutePath()); + Mockito.doReturn(messageStoreConfig).when(controller).getMessageStoreConfig(); + configStorage = new ConfigStorage(messageStoreConfig); configStorage.start(); subscriptionGroupManagerV2 = new SubscriptionGroupManagerV2(controller, configStorage);