diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java index c9c5e04264..dc0339d0cd 100644 --- a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java +++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java @@ -321,9 +321,9 @@ public class ReplicasInfoManager { final ArrayList inSyncReplicas = new ArrayList<>(); final ArrayList notInSyncReplicas = new ArrayList<>(); - brokerInfo.getBrokerIdTable().forEach((brokerAddress, brokerId) -> { + brokerReplicaInfo.getBrokerIdTable().forEach((brokerAddress, brokerId) -> { if (syncStateSet.contains(brokerAddress)) { - long id = StringUtils.equals(master, brokerAddress) ? MixAll.MASTER_ID : brokerInfo.getBrokerId(brokerAddress); + long id = StringUtils.equals(master, brokerAddress) ? MixAll.MASTER_ID : brokerReplicaInfo.getBrokerId(brokerAddress); inSyncReplicas.add(new BrokerReplicasInfo.ReplicaIdentity(brokerAddress, id)); } else { notInSyncReplicas.add(new BrokerReplicasInfo.ReplicaIdentity(brokerAddress, brokerId)); diff --git a/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/impl/manager/ReplicasInfoManagerTest.java b/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/impl/manager/ReplicasInfoManagerTest.java index 57d372349c..8dc637842a 100644 --- a/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/impl/manager/ReplicasInfoManagerTest.java +++ b/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/impl/manager/ReplicasInfoManagerTest.java @@ -32,7 +32,7 @@ import org.apache.rocketmq.controller.impl.event.EventMessage; import org.apache.rocketmq.controller.impl.manager.ReplicasInfoManager; import org.apache.rocketmq.remoting.protocol.RemotingSerializable; import org.apache.rocketmq.remoting.protocol.ResponseCode; -import org.apache.rocketmq.remoting.protocol.body.InSyncStateData; +import org.apache.rocketmq.remoting.protocol.body.BrokerReplicasInfo; import org.apache.rocketmq.remoting.protocol.body.SyncStateSet; import org.apache.rocketmq.remoting.protocol.header.controller.AlterSyncStateSetRequestHeader; import org.apache.rocketmq.remoting.protocol.header.controller.AlterSyncStateSetResponseHeader; @@ -101,7 +101,7 @@ public class ReplicasInfoManagerTest { final GetReplicaInfoResponseHeader replicaInfoBefore = this.replicasInfoManager.getReplicaInfo(new GetReplicaInfoRequestHeader(brokerName, brokerAddress)).getResponse(); byte[] body = this.replicasInfoManager.getSyncStateData(Arrays.asList(brokerName)).getBody(); - InSyncStateData syncStateDataBefore = RemotingSerializable.decode(body, InSyncStateData.class); + BrokerReplicasInfo syncStateDataBefore = RemotingSerializable.decode(body, BrokerReplicasInfo.class); // Try elect itself as a master ElectMasterRequestHeader requestHeader = ElectMasterRequestHeader.ofBrokerTrigger(clusterName, brokerName, brokerAddress); final ControllerResult result = this.replicasInfoManager.electMaster(requestHeader, this.electPolicy); @@ -131,7 +131,7 @@ public class ReplicasInfoManagerTest { assertEquals(brokerId, replicaInfoAfter.getBrokerId()); return; } - if (syncStateDataBefore.getInSyncStateTable().containsKey(brokerAddress) || this.config.isEnableElectUncleanMaster()) { + if (syncStateDataBefore.getReplicasInfoTable().containsKey(brokerAddress) || this.config.isEnableElectUncleanMaster()) { // can be elected successfully assertEquals(ResponseCode.SUCCESS, result.getResponseCode()); assertEquals(MixAll.MASTER_ID, replicaInfoAfter.getBrokerId());