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