diff --git a/broker/src/main/java/org/apache/rocketmq/broker/controller/ReplicasManager.java b/broker/src/main/java/org/apache/rocketmq/broker/controller/ReplicasManager.java
index f22f22a12bd..bc9130087bd 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/controller/ReplicasManager.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/controller/ReplicasManager.java
@@ -206,7 +206,7 @@ private boolean startBasicService() {
return false;
}
}
-
+ // 同步controller的元数据
schedulingSyncBrokerMetadata();
// Register syncStateSet changed listener.
diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/JRaftController.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/JRaftController.java
index e40a6349450..98620654773 100644
--- a/controller/src/main/java/org/apache/rocketmq/controller/impl/JRaftController.java
+++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/JRaftController.java
@@ -116,9 +116,15 @@ private void initPeerIdMap() {
String[] peers = this.controllerConfig.getJraftConfig().getjRaftInitConf().split(",");
String[] rpcAddrs = this.controllerConfig.getJraftConfig().getjRaftControllerRPCAddr().split(",");
for (int i = 0; i < peers.length; i++) {
+ String peerStr = peers[i];
+ // Remove /learner suffix if present, as PeerId.parse() cannot handle it
+ int learnerIndex;
+ if ((learnerIndex = peerStr.indexOf("/learner")) > 0) {
+ peerStr = peerStr.substring(0, learnerIndex);
+ }
PeerId peerId = new PeerId();
- if (!peerId.parse(peers[i])) {
- throw new IllegalArgumentException("Fail to parse peerId:" + peers[i]);
+ if (!peerId.parse(peerStr)) {
+ throw new IllegalArgumentException("Fail to parse peerId:" + peerStr);
}
this.peerIdToAddr.put(peerId, rpcAddrs[i]);
}
diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/BrokerReplicaInfo.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/BrokerReplicaInfo.java
index 1623a05908d..49d265dcbec 100644
--- a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/BrokerReplicaInfo.java
+++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/BrokerReplicaInfo.java
@@ -37,7 +37,7 @@ public class BrokerReplicaInfo implements Serializable {
// Start from 1
private final AtomicLong nextAssignBrokerId;
-
+ //registerCheckCode值:this.brokerAddress + ";" + System.currentTimeMillis(),作用后续ApplyBrokerId动作的身份校验
private final Map> brokerIdInfo;
public BrokerReplicaInfo(String clusterName, String brokerName) {
diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/SyncStateInfo.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/SyncStateInfo.java
index 0b2bc1385b0..70e1a1bca77 100644
--- a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/SyncStateInfo.java
+++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/SyncStateInfo.java
@@ -28,9 +28,11 @@
public class SyncStateInfo implements Serializable {
private final String clusterName;
private final String brokerName;
+ //master版本号(任期号)
private final AtomicInteger masterEpoch;
+ //同步副本数集合的版本(任期号)
private final AtomicInteger syncStateSetEpoch;
-
+ //同步副本数集合包含master
private Set syncStateSet;
private Long masterBrokerId;
diff --git a/pom.xml b/pom.xml
index 003e29e0ffd..b3e7ba8ccdc 100644
--- a/pom.xml
+++ b/pom.xml
@@ -138,7 +138,7 @@
1.29.0-alpha
2.0.6
2.20.29
- 1.0.2
+ 1.0.3
2.13.4.2
1.3.14