mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
Resolve the conflict
This commit is contained in:
+2
-2
@@ -321,9 +321,9 @@ public class ReplicasInfoManager {
|
||||
final ArrayList<BrokerReplicasInfo.ReplicaIdentity> inSyncReplicas = new ArrayList<>();
|
||||
final ArrayList<BrokerReplicasInfo.ReplicaIdentity> 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));
|
||||
|
||||
+3
-3
@@ -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<ElectMasterResponseHeader> 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());
|
||||
|
||||
Reference in New Issue
Block a user