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 b1b4ebd4be..aa96b673a5 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 @@ -30,6 +30,7 @@ import java.util.concurrent.TimeUnit; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.out.BrokerOuterAPI; +import org.apache.rocketmq.client.exception.MQBrokerException; import org.apache.rocketmq.common.BrokerConfig; import org.apache.rocketmq.common.EpochEntry; import org.apache.rocketmq.common.MixAll; @@ -45,6 +46,8 @@ import org.apache.rocketmq.logging.InternalLoggerFactory; import org.apache.rocketmq.store.config.BrokerRole; import org.apache.rocketmq.store.ha.autoswitch.AutoSwitchHAService; +import static org.apache.rocketmq.common.protocol.ResponseCode.CONTROLLER_BROKER_METADATA_NOT_EXIST; + /** * The manager of broker replicas, including: 0.regularly syncing controller metadata, change controller leader address, * both master and slave will start this timed task. 1.regularly syncing metadata from controllers, and changing broker @@ -95,7 +98,6 @@ public class ReplicasManager { return this.haService.getConfirmOffset(); } - enum State { INITIAL, FIRST_TIME_SYNC_CONTROLLER_METADATA_DONE, @@ -116,7 +118,8 @@ public class ReplicasManager { } retryTimes++; LOGGER.warn("Failed to start replicasManager, retry times:{}, current state:{}, try it again", retryTimes, this.state); - } while (!startBasicService()); + } + while (!startBasicService()); LOGGER.info("Start replicasManager success, retry times:{}", retryTimes); }); @@ -286,8 +289,8 @@ public class ReplicasManager { // Register this broker to controller, get brokerId and masterAddress. try { final RegisterBrokerToControllerResponseHeader registerResponse = this.brokerOuterAPI.registerBrokerToController(this.controllerLeaderAddress, - this.brokerConfig.getBrokerClusterName(), this.brokerConfig.getBrokerName(), this.localAddress, - this.haService.getLastEpoch(), this.brokerController.getMessageStore().getMaxPhyOffset()); + this.brokerConfig.getBrokerClusterName(), this.brokerConfig.getBrokerName(), this.localAddress, + this.haService.getLastEpoch(), this.brokerController.getMessageStore().getMaxPhyOffset()); final String newMasterAddress = registerResponse.getMasterAddress(); if (StringUtils.isNoneEmpty(newMasterAddress)) { if (StringUtils.equals(newMasterAddress, this.localAddress)) { @@ -345,6 +348,16 @@ public class ReplicasManager { } } } + } catch (final MQBrokerException exception) { + LOGGER.warn("Error happen when get broker {}'s metadata", this.brokerConfig.getBrokerName(), exception); + if (exception.getResponseCode() == CONTROLLER_BROKER_METADATA_NOT_EXIST) { + try { + registerBrokerToController(); + TimeUnit.SECONDS.sleep(2); + } catch (InterruptedException ignore) { + + } + } } catch (final Exception e) { LOGGER.warn("Error happen when get broker {}'s metadata", this.brokerConfig.getBrokerName(), e); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java index 20605a7dae..188440e049 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java @@ -128,6 +128,7 @@ import org.apache.rocketmq.store.timer.TimerCheckpoint; import org.apache.rocketmq.store.timer.TimerMetrics; import static org.apache.rocketmq.common.protocol.ResponseCode.CONTROLLER_NOT_LEADER; +import static org.apache.rocketmq.common.protocol.ResponseCode.CONTROLLER_BROKER_METADATA_NOT_EXIST; import static org.apache.rocketmq.remoting.protocol.RemotingSysResponseCode.SUCCESS; public class BrokerOuterAPI { @@ -1173,6 +1174,9 @@ public class BrokerOuterAPI { case CONTROLLER_NOT_LEADER: { throw new MQBrokerException(response.getCode(), "Controller leader was changed"); } + case CONTROLLER_BROKER_METADATA_NOT_EXIST: { + throw new MQBrokerException(response.getCode(), response.getRemark()); + } } throw new MQBrokerException(response.getCode(), response.getRemark()); } diff --git a/common/src/main/java/org/apache/rocketmq/common/protocol/ResponseCode.java b/common/src/main/java/org/apache/rocketmq/common/protocol/ResponseCode.java index a947854c4a..f68bac2790 100644 --- a/common/src/main/java/org/apache/rocketmq/common/protocol/ResponseCode.java +++ b/common/src/main/java/org/apache/rocketmq/common/protocol/ResponseCode.java @@ -109,4 +109,6 @@ public class ResponseCode extends RemotingSysResponseCode { public static final int CONTROLLER_BROKER_NOT_ALIVE = 2006; public static final int CONTROLLER_NOT_LEADER = 2007; + public static final int CONTROLLER_BROKER_METADATA_NOT_EXIST = 2008; + } 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 ed123aeb2b..ca36ae2e3a 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 @@ -288,7 +288,7 @@ public class ReplicasInfoManager { result.setBody(new SyncStateSet(syncStateInfo.getSyncStateSet(), syncStateInfo.getSyncStateSetEpoch()).encode()); return result; } - result.setCodeAndRemark(ResponseCode.CONTROLLER_INVALID_REQUEST, "Broker metadata is not existed"); + result.setCodeAndRemark(ResponseCode.CONTROLLER_BROKER_METADATA_NOT_EXIST, "Broker metadata is not existed"); return result; }