mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
Add notifyBrokerRoleChanged configuration
This commit is contained in:
@@ -30,7 +30,7 @@ import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.future.FutureTaskExt;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.common.protocol.header.NotifyBrokerRoleChangedRequestHeader;
|
||||
@@ -109,7 +109,10 @@ public class ControllerManager {
|
||||
if (StringUtils.isNotEmpty(responseHeader.getNewMasterAddress())) {
|
||||
heartbeatManager.changeBrokerMetadata(clusterName, responseHeader.getNewMasterAddress(), MixAll.MASTER_ID);
|
||||
}
|
||||
notifyBrokerMasterChanged(responseHeader, clusterName);
|
||||
|
||||
if (controllerConfig.isNotifyBrokerRoleChanged()) {
|
||||
notifyBrokerRoleChanged(responseHeader, clusterName);
|
||||
}
|
||||
}
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
@@ -126,7 +129,7 @@ public class ControllerManager {
|
||||
/**
|
||||
* Notify master and all slaves for a broker that the master role changed.
|
||||
*/
|
||||
public void notifyBrokerMasterChanged(final ElectMasterResponseHeader electMasterResult, final String clusterName) {
|
||||
public void notifyBrokerRoleChanged(final ElectMasterResponseHeader electMasterResult, final String clusterName) {
|
||||
final BrokerMemberGroup memberGroup = electMasterResult.getBrokerMemberGroup();
|
||||
if (memberGroup != null) {
|
||||
// First, inform the master
|
||||
|
||||
@@ -31,7 +31,7 @@ import org.apache.commons.cli.Options;
|
||||
import org.apache.commons.cli.PosixParser;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
|
||||
@@ -38,7 +38,7 @@ import java.util.function.Supplier;
|
||||
import org.apache.rocketmq.common.ServiceThread;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
@@ -140,6 +140,10 @@ public class DLedgerController implements Controller {
|
||||
return this.roleHandler.isLeaderState();
|
||||
}
|
||||
|
||||
public ControllerConfig getControllerConfig() {
|
||||
return controllerConfig;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<RemotingCommand> alterSyncStateSet(AlterSyncStateSetRequestHeader request,
|
||||
final SyncStateSet syncStateSet) {
|
||||
|
||||
+1
-1
@@ -29,7 +29,7 @@ import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.common.BrokerAddrInfo;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.controller.BrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
|
||||
+1
-1
@@ -27,7 +27,7 @@ import java.util.function.Predicate;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.common.protocol.body.InSyncStateData;
|
||||
|
||||
+1
-1
@@ -25,7 +25,7 @@ import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.client.exception.MQBrokerException;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.BrokerHeartbeatRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.BrokerRegisterRequestHeader;
|
||||
|
||||
+1
-1
@@ -25,7 +25,7 @@ import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.BrokerRegisterRequestHeader;
|
||||
|
||||
+1
-1
@@ -18,7 +18,7 @@ package org.apache.rocketmq.controller.impl.controller.impl;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.controller.BrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.controller.impl.DefaultBrokerHeartbeatManager;
|
||||
import org.junit.Before;
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ package org.apache.rocketmq.controller.impl.controller.impl.manager;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import org.apache.rocketmq.common.namesrv.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
|
||||
|
||||
Reference in New Issue
Block a user