[ISSUE #9206] Fix slave sync topic sub in rocksdb ha (#9207)

* typo int readme[ecosystem]

* fix slave sunc topic and sub in rocksdb ha mode
This commit is contained in:
fujian-zfj
2025-07-19 18:17:56 +08:00
committed by RongtongJin
parent 137e7bf17b
commit e0580bcd3f
4 changed files with 34 additions and 19 deletions
@@ -78,7 +78,6 @@ public class RocksDBSubscriptionGroupManager extends SubscriptionGroupManager {
return true;
}
private boolean merge() {
if (!UtilAll.isPathExists(this.configFilePath()) && !UtilAll.isPathExists(this.configFilePath() + ".bak")) {
log.info("subGroup json file does not exist, so skip merge");
@@ -17,6 +17,8 @@
package org.apache.rocketmq.broker.slave;
import java.io.IOException;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -24,6 +26,7 @@ import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.broker.BrokerController;
import org.apache.rocketmq.broker.loadbalance.MessageRequestModeManager;
import org.apache.rocketmq.broker.subscription.SubscriptionGroupManager;
import org.apache.rocketmq.broker.topic.TopicConfigManager;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.TopicConfig;
import org.apache.rocketmq.common.constant.LoggerName;
@@ -77,20 +80,28 @@ public class SlaveSynchronize {
try {
TopicConfigAndMappingSerializeWrapper topicWrapper =
this.brokerController.getBrokerOuterAPI().getAllTopicConfig(masterAddrBak);
if (!this.brokerController.getTopicConfigManager().getDataVersion()
.equals(topicWrapper.getDataVersion())) {
TopicConfigManager topicConfigManager = this.brokerController.getTopicConfigManager();
if (!topicConfigManager.getDataVersion().equals(topicWrapper.getDataVersion())) {
this.brokerController.getTopicConfigManager().getDataVersion()
.assignNewOne(topicWrapper.getDataVersion());
topicConfigManager.getDataVersion().assignNewOne(topicWrapper.getDataVersion());
ConcurrentMap<String, TopicConfig> newTopicConfigTable = topicWrapper.getTopicConfigTable();
//delete
ConcurrentMap<String, TopicConfig> topicConfigTable = this.brokerController.getTopicConfigManager().getTopicConfigTable();
topicConfigTable.entrySet().removeIf(item -> !newTopicConfigTable.containsKey(item.getKey()));
//update
topicConfigTable.putAll(newTopicConfigTable);
ConcurrentMap<String, TopicConfig> topicConfigTable = topicConfigManager.getTopicConfigTable();
this.brokerController.getTopicConfigManager().persist();
//delete
Iterator<Map.Entry<String, TopicConfig>> iterator = topicConfigTable.entrySet().iterator();
while (iterator.hasNext()) {
Map.Entry<String, TopicConfig> entry = iterator.next();
if (!newTopicConfigTable.containsKey(entry.getKey())) {
iterator.remove();
}
topicConfigManager.deleteTopicConfig(entry.getKey());
}
//update
newTopicConfigTable.values().forEach(topicConfigManager::updateSingleTopicConfigWithoutPersist);
topicConfigManager.persist();
}
if (topicWrapper.getTopicQueueMappingDetailMap() != null
&& !topicWrapper.getMappingDataVersion().equals(this.brokerController.getTopicQueueMappingManager().getDataVersion())) {
@@ -165,19 +176,24 @@ public class SlaveSynchronize {
if (!this.brokerController.getSubscriptionGroupManager().getDataVersion()
.equals(subscriptionWrapper.getDataVersion())) {
SubscriptionGroupManager subscriptionGroupManager =
this.brokerController.getSubscriptionGroupManager();
subscriptionGroupManager.getDataVersion().assignNewOne(
subscriptionWrapper.getDataVersion());
SubscriptionGroupManager subscriptionGroupManager = this.brokerController.getSubscriptionGroupManager();
subscriptionGroupManager.getDataVersion().assignNewOne(subscriptionWrapper.getDataVersion());
ConcurrentMap<String, SubscriptionGroupConfig> curSubscriptionGroupTable =
subscriptionGroupManager.getSubscriptionGroupTable();
ConcurrentMap<String, SubscriptionGroupConfig> newSubscriptionGroupTable =
subscriptionWrapper.getSubscriptionGroupTable();
// delete
curSubscriptionGroupTable.entrySet().removeIf(e -> !newSubscriptionGroupTable.containsKey(e.getKey()));
Iterator<Map.Entry<String, SubscriptionGroupConfig>> iterator = curSubscriptionGroupTable.entrySet().iterator();
while (iterator.hasNext()) {
Map.Entry<String, SubscriptionGroupConfig> configEntry = iterator.next();
if (!newSubscriptionGroupTable.containsKey(configEntry.getKey())) {
iterator.remove();
}
subscriptionGroupManager.deleteSubscriptionGroupConfig(configEntry.getKey());
}
// update
curSubscriptionGroupTable.putAll(newSubscriptionGroupTable);
newSubscriptionGroupTable.values().forEach(subscriptionGroupManager::updateSubscriptionGroupConfigWithoutPersist);
// persist
subscriptionGroupManager.persist();
LOGGER.info("Update slave Subscription Group from master, {}", masterAddrBak);
@@ -143,7 +143,7 @@ public class SubscriptionGroupManager extends ConfigManager {
this.persist();
}
protected void updateSubscriptionGroupConfigWithoutPersist(SubscriptionGroupConfig config) {
public void updateSubscriptionGroupConfigWithoutPersist(SubscriptionGroupConfig config) {
Map<String, String> newAttributes = request(config);
Map<String, String> currentAttributes = current(config.getGroupName());
@@ -497,7 +497,7 @@ public class TopicConfigManager extends ConfigManager {
}
}
protected void updateSingleTopicConfigWithoutPersist(final TopicConfig topicConfig) {
public void updateSingleTopicConfigWithoutPersist(final TopicConfig topicConfig) {
checkNotNull(topicConfig, "topicConfig shouldn't be null");
Map<String, String> newAttributes = request(topicConfig);