mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
[Summer of code] Let broker register to controller again when master not existed. (#4450)
* register to controller if master is null * stop electmaster when the old master is active * review * fix test bugs for brokerAlivePredicate
This commit is contained in:
@@ -69,13 +69,13 @@ public class DLedgerController implements Controller {
|
||||
private final DLedgerServer dLedgerServer;
|
||||
private final ControllerConfig controllerConfig;
|
||||
private final DLedgerConfig dLedgerConfig;
|
||||
// Usr for checking whether the broker is alive
|
||||
private final BiPredicate<String, String> brokerAlivePredicate;
|
||||
private final ReplicasInfoManager replicasInfoManager;
|
||||
private final EventScheduler scheduler;
|
||||
private final EventSerializer eventSerializer;
|
||||
private final RoleChangeHandler roleHandler;
|
||||
private final DLedgerControllerStateMachine statemachine;
|
||||
// Usr for checking whether the broker is alive
|
||||
private BiPredicate<String, String> brokerAlivePredicate;
|
||||
private volatile boolean isScheduling = false;
|
||||
|
||||
public DLedgerController(final ControllerConfig config, final BiPredicate<String, String> brokerAlivePredicate) {
|
||||
@@ -217,6 +217,10 @@ public class DLedgerController implements Controller {
|
||||
return this.dLedgerServer.getMemberState();
|
||||
}
|
||||
|
||||
public void setBrokerAlivePredicate(BiPredicate<String, String> brokerAlivePredicate) {
|
||||
this.brokerAlivePredicate = brokerAlivePredicate;
|
||||
}
|
||||
|
||||
/**
|
||||
* Event handler that handle event
|
||||
*/
|
||||
|
||||
+9
@@ -157,6 +157,15 @@ public class ReplicasInfoManager {
|
||||
final SyncStateInfo syncStateInfo = this.syncStateSetInfoTable.get(brokerName);
|
||||
final BrokerInfo brokerInfo = this.replicaInfoTable.get(brokerName);
|
||||
final Set<String> syncStateSet = syncStateInfo.getSyncStateSet();
|
||||
// First, check whether the master is still active
|
||||
final String oldMaster = syncStateInfo.getMasterAddress();
|
||||
if (StringUtils.isNoneEmpty(oldMaster) && brokerAlivePredicate.test(brokerInfo.getClusterName(), oldMaster)) {
|
||||
String err = String.format("The old master %s is still alive, not need to elect new master for broker %s", oldMaster, brokerInfo.getBrokerName());
|
||||
log.warn("{}", err);
|
||||
result.setCodeAndRemark(ResponseCode.CONTROLLER_INVALID_REQUEST, err);
|
||||
return result;
|
||||
}
|
||||
|
||||
// Try elect a master in syncStateSet
|
||||
if (syncStateSet.size() > 1) {
|
||||
boolean electSuccess = tryElectMaster(result, brokerName, syncStateSet, (candidate) ->
|
||||
|
||||
+14
@@ -164,10 +164,22 @@ public class DLedgerControllerTest {
|
||||
return leader;
|
||||
}
|
||||
|
||||
public void setBrokerAlivePredicate(DLedgerController controller, String... deathBroker) {
|
||||
controller.setBrokerAlivePredicate((clusterName, brokerAddress) -> {
|
||||
for (String broker : deathBroker) {
|
||||
if (broker.equals(brokerAddress)) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testElectMaster() throws Exception {
|
||||
final DLedgerController leader = mockMetaData(false);
|
||||
final ElectMasterRequestHeader request = new ElectMasterRequestHeader("broker1");
|
||||
setBrokerAlivePredicate(leader, "127.0.0.1:9000");
|
||||
final RemotingCommand resp = leader.electMaster(request).get(10, TimeUnit.SECONDS);
|
||||
final ElectMasterResponseHeader response = (ElectMasterResponseHeader) resp.readCustomHeader();
|
||||
assertEquals(response.getMasterEpoch(), 2);
|
||||
@@ -186,6 +198,7 @@ public class DLedgerControllerTest {
|
||||
// Now we trigger electMaster api, which means the old master is shutdown and want to elect a new master.
|
||||
// However, the syncStateSet in statemachine is {"127.0.0.1:9000"}, not more replicas can be elected as master, it will be failed.
|
||||
final ElectMasterRequestHeader electRequest = new ElectMasterRequestHeader("broker1");
|
||||
setBrokerAlivePredicate(leader, "127.0.0.1:9000");
|
||||
leader.electMaster(electRequest).get(10, TimeUnit.SECONDS);
|
||||
|
||||
final RemotingCommand resp = leader.getReplicaInfo(new GetReplicaInfoRequestHeader("broker1")).
|
||||
@@ -225,6 +238,7 @@ public class DLedgerControllerTest {
|
||||
// However, event if the syncStateSet in statemachine is {"127.0.0.1:9000"}
|
||||
// the option {enableElectUncleanMaster = true}, so the controller sill can elect a new master
|
||||
final ElectMasterRequestHeader electRequest = new ElectMasterRequestHeader("broker1");
|
||||
setBrokerAlivePredicate(leader, "127.0.0.1:9000");
|
||||
final CompletableFuture<RemotingCommand> future = leader.electMaster(electRequest);
|
||||
future.get(10, TimeUnit.SECONDS);
|
||||
|
||||
|
||||
+2
-2
@@ -108,7 +108,7 @@ public class ReplicasInfoManagerTest {
|
||||
public void testElectMaster() {
|
||||
mockMetaData();
|
||||
final ElectMasterRequestHeader request = new ElectMasterRequestHeader("broker1");
|
||||
final ControllerResult<ElectMasterResponseHeader> cResult = this.replicasInfoManager.electMaster(request, (va1, va2) -> true);
|
||||
final ControllerResult<ElectMasterResponseHeader> cResult = this.replicasInfoManager.electMaster(request, (clusterName, brokerAddress) -> !brokerAddress.equals("127.0.0.1:9000"));
|
||||
final ElectMasterResponseHeader response = cResult.getResponse();
|
||||
assertEquals(response.getMasterEpoch(), 2);
|
||||
assertFalse(response.getNewMasterAddress().isEmpty());
|
||||
@@ -125,7 +125,7 @@ public class ReplicasInfoManagerTest {
|
||||
// Now we trigger electMaster api, which means the old master is shutdown and want to elect a new master.
|
||||
// However, the syncStateSet in statemachine is {"127.0.0.1:9000"}, not more replicas can be elected as master, it will be failed.
|
||||
final ElectMasterRequestHeader electRequest = new ElectMasterRequestHeader("broker1");
|
||||
final ControllerResult<ElectMasterResponseHeader> cResult = this.replicasInfoManager.electMaster(electRequest, (va1, va2) -> true);
|
||||
final ControllerResult<ElectMasterResponseHeader> cResult = this.replicasInfoManager.electMaster(electRequest, (clusterName, brokerAddress) -> !brokerAddress.equals("127.0.0.1:9000"));
|
||||
final List<EventMessage> events = cResult.getEvents();
|
||||
assertEquals(events.size(), 1);
|
||||
final ElectMasterEvent event = (ElectMasterEvent) events.get(0);
|
||||
|
||||
Reference in New Issue
Block a user