mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
* Fix the invalid of heartbeat detection after controller switch * Pass the checkstyle * Format ReplicasInfoManagerTest code style
This commit is contained in:
+5
-13
@@ -23,12 +23,9 @@ public interface BrokerHeartbeatManager {
|
||||
/**
|
||||
* Broker new heartbeat.
|
||||
*/
|
||||
void onBrokerHeartbeat(final String clusterName, final String brokerAddr, final Integer epoch, final Long maxOffset, final Long confirmOffset);
|
||||
|
||||
/**
|
||||
* Change the metadata(brokerId ..) for a broker.
|
||||
*/
|
||||
void changeBrokerMetadata(final String clusterName, final String brokerAddr, final Long brokerId);
|
||||
void onBrokerHeartbeat(final String clusterName, final String brokerName, final String brokerAddr,
|
||||
final Long brokerId, final Long timeoutMillis, final Channel channel, final Integer epoch,
|
||||
final Long maxOffset, final Long confirmOffset, final Integer electionPriority);
|
||||
|
||||
/**
|
||||
* Start heartbeat manager.
|
||||
@@ -45,12 +42,6 @@ public interface BrokerHeartbeatManager {
|
||||
*/
|
||||
void addBrokerLifecycleListener(final BrokerLifecycleListener listener);
|
||||
|
||||
/**
|
||||
* Register new broker to heartManager.
|
||||
*/
|
||||
void registerBroker(final String clusterName, final String brokerName, final String brokerAddr, final long brokerId,
|
||||
final Long timeoutMillis, final Channel channel, final Integer epoch, final Long maxOffset, final Integer electionPriority);
|
||||
|
||||
/**
|
||||
* Broker channel close
|
||||
*/
|
||||
@@ -70,6 +61,7 @@ public interface BrokerHeartbeatManager {
|
||||
/**
|
||||
* Trigger when broker inactive.
|
||||
*/
|
||||
void onBrokerInactive(final String clusterName, final String brokerName, final String brokerAddress, final long brokerId);
|
||||
void onBrokerInactive(final String clusterName, final String brokerName, final String brokerAddress,
|
||||
final long brokerId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,15 +16,13 @@
|
||||
*/
|
||||
package org.apache.rocketmq.controller;
|
||||
|
||||
|
||||
import io.netty.channel.Channel;
|
||||
|
||||
|
||||
public class BrokerLiveInfo {
|
||||
private final String brokerName;
|
||||
|
||||
private final String brokerAddr;
|
||||
private final long heartbeatTimeoutMillis;
|
||||
private long heartbeatTimeoutMillis;
|
||||
private final Channel channel;
|
||||
private long brokerId;
|
||||
private long lastUpdateTimestamp;
|
||||
@@ -33,8 +31,8 @@ public class BrokerLiveInfo {
|
||||
private long confirmOffset;
|
||||
private Integer electionPriority;
|
||||
|
||||
public BrokerLiveInfo(String brokerName, String brokerAddr,long brokerId, long lastUpdateTimestamp, long heartbeatTimeoutMillis,
|
||||
Channel channel, int epoch, long maxOffset, Integer electionPriority) {
|
||||
public BrokerLiveInfo(String brokerName, String brokerAddr, long brokerId, long lastUpdateTimestamp,
|
||||
long heartbeatTimeoutMillis, Channel channel, int epoch, long maxOffset, Integer electionPriority) {
|
||||
this.brokerName = brokerName;
|
||||
this.brokerAddr = brokerAddr;
|
||||
this.brokerId = brokerId;
|
||||
@@ -46,8 +44,8 @@ public class BrokerLiveInfo {
|
||||
this.maxOffset = maxOffset;
|
||||
}
|
||||
|
||||
public BrokerLiveInfo(String brokerName, String brokerAddr,long brokerId, long lastUpdateTimestamp, long heartbeatTimeoutMillis,
|
||||
Channel channel, int epoch, long maxOffset, Integer electionPriority, long confirmOffset) {
|
||||
public BrokerLiveInfo(String brokerName, String brokerAddr, long brokerId, long lastUpdateTimestamp,
|
||||
long heartbeatTimeoutMillis, Channel channel, int epoch, long maxOffset, Integer electionPriority, long confirmOffset) {
|
||||
this.brokerName = brokerName;
|
||||
this.brokerAddr = brokerAddr;
|
||||
this.brokerId = brokerId;
|
||||
@@ -63,16 +61,16 @@ public class BrokerLiveInfo {
|
||||
@Override
|
||||
public String toString() {
|
||||
return "BrokerLiveInfo{" +
|
||||
"brokerName='" + brokerName + '\'' +
|
||||
", brokerAddr='" + brokerAddr + '\'' +
|
||||
", heartbeatTimeoutMillis=" + heartbeatTimeoutMillis +
|
||||
", channel=" + channel +
|
||||
", brokerId=" + brokerId +
|
||||
", lastUpdateTimestamp=" + lastUpdateTimestamp +
|
||||
", epoch=" + epoch +
|
||||
", maxOffset=" + maxOffset +
|
||||
", confirmOffset=" + confirmOffset +
|
||||
'}';
|
||||
"brokerName='" + brokerName + '\'' +
|
||||
", brokerAddr='" + brokerAddr + '\'' +
|
||||
", heartbeatTimeoutMillis=" + heartbeatTimeoutMillis +
|
||||
", channel=" + channel +
|
||||
", brokerId=" + brokerId +
|
||||
", lastUpdateTimestamp=" + lastUpdateTimestamp +
|
||||
", epoch=" + epoch +
|
||||
", maxOffset=" + maxOffset +
|
||||
", confirmOffset=" + confirmOffset +
|
||||
'}';
|
||||
}
|
||||
|
||||
public String getBrokerName() {
|
||||
@@ -83,6 +81,10 @@ public class BrokerLiveInfo {
|
||||
return heartbeatTimeoutMillis;
|
||||
}
|
||||
|
||||
public void setHeartbeatTimeoutMillis(long heartbeatTimeoutMillis) {
|
||||
this.heartbeatTimeoutMillis = heartbeatTimeoutMillis;
|
||||
}
|
||||
|
||||
public Channel getChannel() {
|
||||
return channel;
|
||||
}
|
||||
@@ -138,4 +140,5 @@ public class BrokerLiveInfo {
|
||||
public long getConfirmOffset() {
|
||||
return confirmOffset;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -48,6 +48,8 @@ import org.apache.rocketmq.remoting.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.remoting.protocol.header.NotifyBrokerRoleChangedRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.controller.GetReplicaInfoResponseHeader;
|
||||
|
||||
public class ControllerManager {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.CONTROLLER_LOGGER_NAME);
|
||||
@@ -64,7 +66,7 @@ public class ControllerManager {
|
||||
private BlockingQueue<Runnable> controllerRequestThreadPoolQueue;
|
||||
|
||||
public ControllerManager(ControllerConfig controllerConfig, NettyServerConfig nettyServerConfig,
|
||||
NettyClientConfig nettyClientConfig) {
|
||||
NettyClientConfig nettyClientConfig) {
|
||||
this.controllerConfig = controllerConfig;
|
||||
this.nettyServerConfig = nettyServerConfig;
|
||||
this.nettyClientConfig = nettyClientConfig;
|
||||
@@ -77,12 +79,12 @@ public class ControllerManager {
|
||||
public boolean initialize() {
|
||||
this.controllerRequestThreadPoolQueue = new LinkedBlockingQueue<>(this.controllerConfig.getControllerRequestThreadPoolQueueCapacity());
|
||||
this.controllerRequestExecutor = new ThreadPoolExecutor(
|
||||
this.controllerConfig.getControllerThreadPoolNums(),
|
||||
this.controllerConfig.getControllerThreadPoolNums(),
|
||||
1000 * 60,
|
||||
TimeUnit.MILLISECONDS,
|
||||
this.controllerRequestThreadPoolQueue,
|
||||
new ThreadFactoryImpl("ControllerRequestExecutorThread_")) {
|
||||
this.controllerConfig.getControllerThreadPoolNums(),
|
||||
this.controllerConfig.getControllerThreadPoolNums(),
|
||||
1000 * 60,
|
||||
TimeUnit.MILLISECONDS,
|
||||
this.controllerRequestThreadPoolQueue,
|
||||
new ThreadFactoryImpl("ControllerRequestExecutorThread_")) {
|
||||
@Override
|
||||
protected <T> RunnableFuture<T> newTaskFor(final Runnable runnable, final T value) {
|
||||
return new FutureTaskExt<>(runnable, value);
|
||||
@@ -96,8 +98,8 @@ public class ControllerManager {
|
||||
throw new IllegalArgumentException("Attribute value controllerDLegerSelfId of ControllerConfig is null or empty");
|
||||
}
|
||||
this.controller = new DLedgerController(this.controllerConfig, this.heartbeatManager::isBrokerActive,
|
||||
this.nettyServerConfig, this.nettyClientConfig, this.brokerHousekeepingService,
|
||||
new DefaultElectPolicy(this.heartbeatManager::isBrokerActive, this.heartbeatManager::getBrokerLiveInfo));
|
||||
this.nettyServerConfig, this.nettyClientConfig, this.brokerHousekeepingService,
|
||||
new DefaultElectPolicy(this.heartbeatManager::isBrokerActive, this.heartbeatManager::getBrokerLiveInfo));
|
||||
|
||||
// Register broker inactive listener
|
||||
this.heartbeatManager.addBrokerLifecycleListener(this::onBrokerInactive);
|
||||
@@ -106,34 +108,39 @@ public class ControllerManager {
|
||||
}
|
||||
|
||||
/**
|
||||
* When the heartbeatManager detects the "Broker is not active",
|
||||
* we call this method to elect a master and do something else.
|
||||
* When the heartbeatManager detects the "Broker is not active", we call this method to elect a master and do
|
||||
* something else.
|
||||
*
|
||||
* @param clusterName The cluster name of this inactive broker
|
||||
* @param brokerName The inactive broker name
|
||||
* @param brokerAddress The inactive broker address(ip)
|
||||
* @param brokerId The inactive broker id
|
||||
*/
|
||||
private void onBrokerInactive(String clusterName, String brokerName, String brokerAddress, long brokerId) {
|
||||
if (brokerId == MixAll.MASTER_ID) {
|
||||
if (controller.isLeaderState()) {
|
||||
final CompletableFuture<RemotingCommand> future = controller.electMaster(new ElectMasterRequestHeader(brokerName));
|
||||
try {
|
||||
final RemotingCommand response = future.get(5, TimeUnit.SECONDS);
|
||||
final ElectMasterResponseHeader responseHeader = (ElectMasterResponseHeader) response.readCustomHeader();
|
||||
if (responseHeader != null) {
|
||||
log.info("Broker {}'s master {} shutdown, elect a new master done, result:{}", brokerName, brokerAddress, responseHeader);
|
||||
if (StringUtils.isNotEmpty(responseHeader.getNewMasterAddress())) {
|
||||
heartbeatManager.changeBrokerMetadata(clusterName, responseHeader.getNewMasterAddress(), MixAll.MASTER_ID);
|
||||
}
|
||||
if (controllerConfig.isNotifyBrokerRoleChanged()) {
|
||||
notifyBrokerRoleChanged(responseHeader, clusterName);
|
||||
}
|
||||
}
|
||||
} catch (Exception ignored) {
|
||||
if (controller.isLeaderState()) {
|
||||
try {
|
||||
final CompletableFuture<RemotingCommand> replicaInfoFuture = controller.getReplicaInfo(new GetReplicaInfoRequestHeader(brokerName, brokerAddress));
|
||||
final RemotingCommand replicaInfoResponse = replicaInfoFuture.get(5, TimeUnit.SECONDS);
|
||||
final GetReplicaInfoResponseHeader replicaInfoResponseHeader = (GetReplicaInfoResponseHeader) replicaInfoResponse.readCustomHeader();
|
||||
// Not master broker offline
|
||||
if (!replicaInfoResponseHeader.getMasterAddress().equals(brokerAddress)) {
|
||||
log.warn("The {} broker with IP address {} shutdown", brokerName, brokerAddress);
|
||||
return;
|
||||
}
|
||||
} else {
|
||||
log.info("Broker{}' master shutdown", brokerName);
|
||||
final CompletableFuture<RemotingCommand> electMasterFuture = controller.electMaster(new ElectMasterRequestHeader(brokerName));
|
||||
final RemotingCommand electMasterResponse = electMasterFuture.get(5, TimeUnit.SECONDS);
|
||||
final ElectMasterResponseHeader responseHeader = (ElectMasterResponseHeader) electMasterResponse.readCustomHeader();
|
||||
if (responseHeader != null) {
|
||||
log.info("Broker {}'s master {} shutdown, elect a new master done, result:{}", brokerName, brokerAddress, responseHeader);
|
||||
if (controllerConfig.isNotifyBrokerRoleChanged()) {
|
||||
notifyBrokerRoleChanged(responseHeader, clusterName);
|
||||
}
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.error("", e);
|
||||
}
|
||||
} else {
|
||||
log.info("The {} broker with IP address {} shutdown", brokerName, brokerAddress);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -161,11 +168,11 @@ public class ControllerManager {
|
||||
}
|
||||
|
||||
public void doNotifyBrokerRoleChanged(final String brokerAddr, final Long brokerId,
|
||||
final ElectMasterResponseHeader responseHeader) {
|
||||
final ElectMasterResponseHeader responseHeader) {
|
||||
if (StringUtils.isNoneEmpty(brokerAddr)) {
|
||||
log.info("Try notify broker {} with id {} that role changed, responseHeader:{}", brokerAddr, brokerId, responseHeader);
|
||||
final NotifyBrokerRoleChangedRequestHeader requestHeader = new NotifyBrokerRoleChangedRequestHeader(responseHeader.getNewMasterAddress(),
|
||||
responseHeader.getMasterEpoch(), responseHeader.getSyncStateSetEpoch(), brokerId);
|
||||
responseHeader.getMasterEpoch(), responseHeader.getSyncStateSetEpoch(), brokerId);
|
||||
final RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.NOTIFY_BROKER_ROLE_CHANGED, requestHeader);
|
||||
try {
|
||||
this.remotingClient.invokeOneway(brokerAddr, request, 3000);
|
||||
|
||||
+28
-41
@@ -99,53 +99,40 @@ public class DefaultBrokerHeartbeatManager implements BrokerHeartbeatManager {
|
||||
this.brokerLifecycleListeners.add(listener);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void registerBroker(String clusterName, String brokerName, String brokerAddr,
|
||||
long brokerId, Long timeoutMillis, Channel channel, Integer epoch, Long maxOffset, final Integer electionPriority) {
|
||||
final BrokerAddrInfo addrInfo = new BrokerAddrInfo(clusterName, brokerAddr);
|
||||
final BrokerLiveInfo prevBrokerLiveInfo = this.brokerLiveTable.put(addrInfo,
|
||||
new BrokerLiveInfo(brokerName,
|
||||
brokerAddr,
|
||||
brokerId,
|
||||
System.currentTimeMillis(),
|
||||
timeoutMillis == null ? DEFAULT_BROKER_CHANNEL_EXPIRED_TIME : timeoutMillis,
|
||||
channel,
|
||||
epoch == null ? -1 : epoch,
|
||||
maxOffset == null ? -1 : maxOffset,
|
||||
electionPriority == null ? Integer.MAX_VALUE : electionPriority));
|
||||
if (prevBrokerLiveInfo == null) {
|
||||
log.info("new broker registered, {}, brokerId:{}", addrInfo, brokerId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void changeBrokerMetadata(String clusterName, String brokerAddr, Long brokerId) {
|
||||
@Override public void onBrokerHeartbeat(String clusterName, String brokerName, String brokerAddr, Long brokerId,
|
||||
Long timeoutMillis, Channel channel, Integer epoch, Long maxOffset, Long confirmOffset, Integer electionPriority) {
|
||||
BrokerAddrInfo addrInfo = new BrokerAddrInfo(clusterName, brokerAddr);
|
||||
BrokerLiveInfo prev = this.brokerLiveTable.get(addrInfo);
|
||||
if (prev != null) {
|
||||
prev.setBrokerId(brokerId);
|
||||
log.info("Change broker {}'s brokerId to {}", brokerAddr, brokerId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onBrokerHeartbeat(String clusterName, String brokerAddr, Integer epoch, Long maxOffset,
|
||||
Long confirmOffset) {
|
||||
BrokerAddrInfo addrInfo = new BrokerAddrInfo(clusterName, brokerAddr);
|
||||
BrokerLiveInfo prev = this.brokerLiveTable.get(addrInfo);
|
||||
if (null == prev) {
|
||||
return;
|
||||
}
|
||||
int realEpoch = Optional.ofNullable(epoch).orElse(-1);
|
||||
long realBrokerId = Optional.ofNullable(brokerId).orElse(-1L);
|
||||
long realMaxOffset = Optional.ofNullable(maxOffset).orElse(-1L);
|
||||
long realConfirmOffset = Optional.ofNullable(confirmOffset).orElse(-1L);
|
||||
|
||||
prev.setLastUpdateTimestamp(System.currentTimeMillis());
|
||||
if (realEpoch > prev.getEpoch() || realEpoch == prev.getEpoch() && realMaxOffset > prev.getMaxOffset()) {
|
||||
prev.setEpoch(realEpoch);
|
||||
prev.setMaxOffset(realMaxOffset);
|
||||
prev.setConfirmOffset(realConfirmOffset);
|
||||
long realTimeoutMillis = Optional.ofNullable(timeoutMillis).orElse(DEFAULT_BROKER_CHANNEL_EXPIRED_TIME);
|
||||
int realElectionPriority = Optional.ofNullable(electionPriority).orElse(Integer.MAX_VALUE);
|
||||
if (null == prev) {
|
||||
this.brokerLiveTable.put(addrInfo,
|
||||
new BrokerLiveInfo(brokerName,
|
||||
brokerAddr,
|
||||
realBrokerId,
|
||||
System.currentTimeMillis(),
|
||||
realTimeoutMillis,
|
||||
channel,
|
||||
realEpoch,
|
||||
realMaxOffset,
|
||||
realElectionPriority));
|
||||
log.info("new broker registered, {}, brokerId:{}", addrInfo, realBrokerId);
|
||||
} else {
|
||||
prev.setLastUpdateTimestamp(System.currentTimeMillis());
|
||||
prev.setHeartbeatTimeoutMillis(realTimeoutMillis);
|
||||
prev.setElectionPriority(realElectionPriority);
|
||||
prev.setBrokerId(realBrokerId);
|
||||
if (realEpoch > prev.getEpoch() || realEpoch == prev.getEpoch() && realMaxOffset > prev.getMaxOffset()) {
|
||||
prev.setEpoch(realEpoch);
|
||||
prev.setMaxOffset(realMaxOffset);
|
||||
prev.setConfirmOffset(realConfirmOffset);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+5
-9
@@ -22,7 +22,6 @@ import java.util.List;
|
||||
import java.util.Properties;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.controller.BrokerHeartbeatManager;
|
||||
@@ -95,9 +94,6 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
final ElectMasterResponseHeader responseHeader = (ElectMasterResponseHeader) response.readCustomHeader();
|
||||
|
||||
if (null != responseHeader) {
|
||||
if (StringUtils.isNotEmpty(responseHeader.getNewMasterAddress())) {
|
||||
heartbeatManager.changeBrokerMetadata(electMasterRequest.getClusterName(), responseHeader.getNewMasterAddress(), MixAll.MASTER_ID);
|
||||
}
|
||||
if (this.controllerManager.getControllerConfig().isNotifyBrokerRoleChanged()) {
|
||||
this.controllerManager.notifyBrokerRoleChanged(responseHeader, electMasterRequest.getClusterName());
|
||||
}
|
||||
@@ -113,9 +109,9 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
final RemotingCommand response = future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS);
|
||||
final RegisterBrokerToControllerResponseHeader responseHeader = (RegisterBrokerToControllerResponseHeader) response.readCustomHeader();
|
||||
if (responseHeader != null && responseHeader.getBrokerId() >= 0) {
|
||||
this.heartbeatManager.registerBroker(controllerRequest.getClusterName(), controllerRequest.getBrokerName(), controllerRequest.getBrokerAddress(),
|
||||
responseHeader.getBrokerId(), controllerRequest.getHeartbeatTimeoutMillis(), ctx.channel(),
|
||||
controllerRequest.getEpoch(), controllerRequest.getMaxOffset(), controllerRequest.getElectionPriority());
|
||||
this.heartbeatManager.onBrokerHeartbeat(controllerRequest.getClusterName(), controllerRequest.getBrokerName(), controllerRequest.getBrokerAddress(),
|
||||
responseHeader.getBrokerId(), controllerRequest.getHeartbeatTimeoutMillis(), ctx.channel(),
|
||||
controllerRequest.getEpoch(), controllerRequest.getMaxOffset(), controllerRequest.getConfirmOffset(), controllerRequest.getElectionPriority());
|
||||
}
|
||||
return response;
|
||||
}
|
||||
@@ -134,8 +130,8 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
}
|
||||
case BROKER_HEARTBEAT: {
|
||||
final BrokerHeartbeatRequestHeader requestHeader = (BrokerHeartbeatRequestHeader) request.decodeCommandCustomHeader(BrokerHeartbeatRequestHeader.class);
|
||||
this.heartbeatManager.onBrokerHeartbeat(requestHeader.getClusterName(), requestHeader.getBrokerAddr(),
|
||||
requestHeader.getEpoch(), requestHeader.getMaxOffset(), requestHeader.getConfirmOffset());
|
||||
this.heartbeatManager.onBrokerHeartbeat(requestHeader.getClusterName(), requestHeader.getBrokerName(), requestHeader.getBrokerAddr(), requestHeader.getBrokerId(),
|
||||
requestHeader.getHeartbeatTimeoutMills(), ctx.channel(), requestHeader.getEpoch(), requestHeader.getMaxOffset(), requestHeader.getConfirmOffset(), requestHeader.getElectionPriority());
|
||||
return RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "Heart beat success");
|
||||
}
|
||||
case CONTROLLER_GET_SYNC_STATE_DATA: {
|
||||
|
||||
+3
-5
@@ -128,9 +128,7 @@ public class ControllerManagerTest {
|
||||
final String brokerName, final String address, final RemotingClient client,
|
||||
final long heartbeatTimeoutMillis) throws Exception {
|
||||
|
||||
final RegisterBrokerToControllerRequestHeader requestHeader = new RegisterBrokerToControllerRequestHeader(clusterName, brokerName, address);
|
||||
// Timeout = 3000
|
||||
requestHeader.setHeartbeatTimeoutMillis(heartbeatTimeoutMillis);
|
||||
final RegisterBrokerToControllerRequestHeader requestHeader = new RegisterBrokerToControllerRequestHeader(clusterName, brokerName, address, heartbeatTimeoutMillis, 1, 1000L, 0);
|
||||
final RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONTROLLER_REGISTER_BROKER, requestHeader);
|
||||
final RemotingCommand response = client.invokeSync(controllerAddress, request, 3000);
|
||||
assert response != null;
|
||||
@@ -173,8 +171,8 @@ public class ControllerManagerTest {
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}, 0, 2000L, TimeUnit.MILLISECONDS);
|
||||
Boolean flag = await().atMost(Duration.ofSeconds(5)).until(() -> {
|
||||
}, 0, 1000L, TimeUnit.MILLISECONDS);
|
||||
Boolean flag = await().atMost(Duration.ofSeconds(10)).until(() -> {
|
||||
final GetReplicaInfoRequestHeader requestHeader = new GetReplicaInfoRequestHeader("broker1");
|
||||
final RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONTROLLER_GET_REPLICA_INFO, requestHeader);
|
||||
final RemotingCommand response = this.remotingClient1.invokeSync(leaderAddr, request, 3000);
|
||||
|
||||
+2
-1
@@ -43,7 +43,8 @@ public class DefaultBrokerHeartbeatManagerTest {
|
||||
this.heartbeatManager.addBrokerLifecycleListener((clusterName, brokerName, brokerAddress, brokerId) -> {
|
||||
latch.countDown();
|
||||
});
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:7000", 1L, 3000L, null, 1, 1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:7000", 1L,3000L, null,
|
||||
1, 1L,-1L, 0);
|
||||
assertTrue(latch.await(5000, TimeUnit.MILLISECONDS));
|
||||
this.heartbeatManager.shutdown();
|
||||
}
|
||||
|
||||
+24
-24
@@ -150,39 +150,39 @@ public class ReplicasInfoManagerTest {
|
||||
}
|
||||
|
||||
public void mockHeartbeatDataMasterStillAlive() {
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9000", 1L, 10000000000L, null,
|
||||
1, 3L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9000", 1L, 10000000000L, null,
|
||||
1, 1L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, -1L, 0);
|
||||
}
|
||||
|
||||
public void mockHeartbeatDataHigherEpoch() {
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
0, 3L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
0, 3L, -1L, 0);
|
||||
}
|
||||
|
||||
public void mockHeartbeatDataHigherOffset() {
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, 0);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 2L, -1L, 0);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, -1L, 0);
|
||||
}
|
||||
|
||||
public void mockHeartbeatDataHigherPriority() {
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, 3);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 3L, 2);
|
||||
this.heartbeatManager.registerBroker("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, 1);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9000", 1L, -10000L, null,
|
||||
1, 3L, -1L, 3);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9001", 1L, 10000000000L, null,
|
||||
1, 3L, -1L, 2);
|
||||
this.heartbeatManager.onBrokerHeartbeat("cluster1", "broker1", "127.0.0.1:9002", 1L, 10000000000L, null,
|
||||
1, 3L, -1L, 1);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user