From 97e6cecbf9d6df1c548b110363b2928a7348da66 Mon Sep 17 00:00:00 2001
From: hzh0425 <642256541@qq.com>
Date: Thu, 26 May 2022 20:44:34 +0800
Subject: [PATCH] [Summer of code] Fix some bugs in controller mode.
* modify pom for using dledger
* rename option
* 1.send heartbeat to haclient when get epochEntry failed to hold connection.
2.update lastCatchupTimeMs in time.
* fix some bugs
---
.../apache/rocketmq/common/BrokerConfig.java | 6 ++--
.../common/namesrv/ControllerConfig.java | 10 +++---
controller/pom.xml | 11 ++++++-
namesrv/pom.xml | 5 ---
.../rocketmq/namesrv/NamesrvController.java | 8 ++---
.../namesrv/routeinfo/RouteInfoManager.java | 2 +-
pom.xml | 5 +++
store/pom.xml | 1 -
.../store/config/MessageStoreConfig.java | 8 ++---
.../ha/autoswitch/AutoSwitchHAClient.java | 17 ++++------
.../ha/autoswitch/AutoSwitchHAConnection.java | 32 ++++++++++++-------
.../ha/autoswitch/AutoSwitchHAService.java | 9 +++---
.../autoswitchrole/AutoSwitchRoleBase.java | 2 +-
13 files changed, 64 insertions(+), 52 deletions(-)
diff --git a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java
index 26d05e1d7c..a181a8ee3e 100644
--- a/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java
+++ b/common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java
@@ -303,7 +303,7 @@ public class BrokerConfig extends BrokerIdentity {
/**
* Whether the controller is deployed independently
*/
- private boolean isControllerDeployedStandAlone = false;
+ private boolean controllerDeployedStandAlone = false;
/**
* If isControllerDeployedStandAlone = false, controllerAddr should be the addresses of all name-srv which running the controller instance.
@@ -1302,11 +1302,11 @@ public class BrokerConfig extends BrokerIdentity {
}
public boolean isControllerDeployedStandAlone() {
- return isControllerDeployedStandAlone;
+ return controllerDeployedStandAlone;
}
public void setControllerDeployedStandAlone(boolean controllerDeployedStandAlone) {
- isControllerDeployedStandAlone = controllerDeployedStandAlone;
+ this.controllerDeployedStandAlone = controllerDeployedStandAlone;
}
public String getControllerAddr() {
diff --git a/common/src/main/java/org/apache/rocketmq/common/namesrv/ControllerConfig.java b/common/src/main/java/org/apache/rocketmq/common/namesrv/ControllerConfig.java
index 1d73c9d026..ca2a423eb1 100644
--- a/common/src/main/java/org/apache/rocketmq/common/namesrv/ControllerConfig.java
+++ b/common/src/main/java/org/apache/rocketmq/common/namesrv/ControllerConfig.java
@@ -27,7 +27,7 @@ public class ControllerConfig {
/**
* Is startup the controller in this name-srv
*/
- private boolean isStartupController = false;
+ private boolean enableStartupController = false;
/**
* Interval of periodic scanning for non-active broker;
@@ -76,12 +76,12 @@ public class ControllerConfig {
this.configStorePath = configStorePath;
}
- public boolean isStartupController() {
- return isStartupController;
+ public boolean isEnableStartupController() {
+ return enableStartupController;
}
- public void setStartupController(boolean startupController) {
- isStartupController = startupController;
+ public void setEnableStartupController(boolean enableStartupController) {
+ this.enableStartupController = enableStartupController;
}
public long getScanNotActiveBrokerInterval() {
diff --git a/controller/pom.xml b/controller/pom.xml
index 920ea0de2f..acddd82414 100644
--- a/controller/pom.xml
+++ b/controller/pom.xml
@@ -30,7 +30,16 @@
io.openmessaging.storage
dledger
- 0.2.5
+
+
+ org.apache.rocketmq
+ rocketmq-remoting
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
${project.groupId}
diff --git a/namesrv/pom.xml b/namesrv/pom.xml
index 1ce091429e..6461c3fae6 100644
--- a/namesrv/pom.xml
+++ b/namesrv/pom.xml
@@ -28,11 +28,6 @@
rocketmq-namesrv ${project.version}
-
- io.openmessaging.storage
- dledger
- 0.2.5
-
${project.groupId}
rocketmq-controller
diff --git a/namesrv/src/main/java/org/apache/rocketmq/namesrv/NamesrvController.java b/namesrv/src/main/java/org/apache/rocketmq/namesrv/NamesrvController.java
index 0b5e39d64a..3f2e0f7f46 100644
--- a/namesrv/src/main/java/org/apache/rocketmq/namesrv/NamesrvController.java
+++ b/namesrv/src/main/java/org/apache/rocketmq/namesrv/NamesrvController.java
@@ -110,7 +110,7 @@ public class NamesrvController {
this.kvConfigManager = new KVConfigManager(this);
this.brokerHousekeepingService = new BrokerHousekeepingService(this);
this.routeInfoManager = new RouteInfoManager(namesrvConfig, this);
- if (controllerConfig.isStartupController()) {
+ if (controllerConfig.isEnableStartupController()) {
this.controller = new DledgerController(controllerConfig, this.routeInfoManager::isBrokerAlive);
this.routeInfoManager.setController(this.controller);
}
@@ -155,7 +155,7 @@ public class NamesrvController {
}
};
- if (this.controllerConfig.isStartupController()) {
+ if (this.controllerConfig.isEnableStartupController()) {
this.controllerRequestThreadPoolQueue = new LinkedBlockingQueue<>(this.controllerConfig.getControllerRequestThreadPoolQueueCapacity());
this.controllerRequestExecutor = new ThreadPoolExecutor(
this.controllerConfig.getControllerThreadPoolNums(),
@@ -269,7 +269,7 @@ public class NamesrvController {
this.remotingServer.registerDefaultProcessor(new DefaultRequestProcessor(this), this.defaultExecutor);
- if (controllerConfig.isStartupController()) {
+ if (controllerConfig.isEnableStartupController()) {
final ControllerRequestProcessor controllerRequestProcessor = new ControllerRequestProcessor(this.controller);
this.remotingServer.registerProcessor(RequestCode.CONTROLLER_ALTER_SYNC_STATE_SET, controllerRequestProcessor, this.controllerRequestExecutor);
this.remotingServer.registerProcessor(RequestCode.CONTROLLER_ELECT_MASTER, controllerRequestProcessor, this.controllerRequestExecutor);
@@ -290,7 +290,7 @@ public class NamesrvController {
this.routeInfoManager.start();
- if (this.controllerConfig.isStartupController()) {
+ if (this.controllerConfig.isEnableStartupController()) {
this.controller.startup();
}
}
diff --git a/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java b/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java
index cb9c113edb..f1e89f59ca 100644
--- a/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java
+++ b/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java
@@ -600,7 +600,7 @@ public class RouteInfoManager {
}
// Check whether we need to elect a new master
- if (this.namesrvController != null && this.namesrvController.getControllerConfig().isStartupController() && this.controller != null) {
+ if (this.namesrvController != null && this.namesrvController.getControllerConfig().isEnableStartupController() && this.controller != null) {
if (unRegisterRequest.getBrokerId() == 0) {
this.controller.electMaster(new ElectMasterRequestHeader(unRegisterRequest.getBrokerName()));
}
diff --git a/pom.xml b/pom.xml
index 8bfef09867..9c72fffb00 100644
--- a/pom.xml
+++ b/pom.xml
@@ -534,6 +534,11 @@
rocketmq-example
${project.version}
+
+ io.openmessaging.storage
+ dledger
+ 0.2.5
+
org.slf4j
slf4j-api
diff --git a/store/pom.xml b/store/pom.xml
index 55880524ff..6ce81d6e85 100644
--- a/store/pom.xml
+++ b/store/pom.xml
@@ -31,7 +31,6 @@
io.openmessaging.storage
dledger
- 0.2.3
org.apache.rocketmq
diff --git a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
index 80ac3f79d8..6843537dc5 100644
--- a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
+++ b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
@@ -36,7 +36,7 @@ public class MessageStoreConfig {
//The directory in which the commitlog is kept
@ImportantField
private String storePathEpochFile = System.getProperty("user.home") + File.separator + "store"
- + File.separator + "commitlog" + File.separator + "epochFileCheckpoint";
+ + File.separator + "epochFileCheckpoint";
private String readOnlyCommitLogStorePaths = null;
@@ -306,7 +306,7 @@ public class MessageStoreConfig {
*/
private boolean syncFromLastFile = false;
- private boolean isAsyncLearner = false;
+ private boolean asyncLearner = false;
public boolean isDebugLockEnable() {
return debugLockEnable;
@@ -1326,10 +1326,10 @@ public class MessageStoreConfig {
}
public boolean isAsyncLearner() {
- return isAsyncLearner;
+ return asyncLearner;
}
public void setAsyncLearner(boolean asyncLearner) {
- this.isAsyncLearner = asyncLearner;
+ this.asyncLearner = asyncLearner;
}
}
diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
index f784d3873b..01272d47aa 100644
--- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
+++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
@@ -156,14 +156,14 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
@Override public void updateMasterAddress(String newAddress) {
String currentAddr = this.masterAddress.get();
- if (masterAddress.compareAndSet(currentAddr, newAddress)) {
+ if (!StringUtils.equals(newAddress, currentAddr) && masterAddress.compareAndSet(currentAddr, newAddress)) {
LOGGER.info("update master address, OLD: " + currentAddr + " NEW: " + newAddress);
}
}
@Override public void updateHaMasterAddress(String newAddress) {
String currentAddr = this.masterHaAddress.get();
- if (masterHaAddress.compareAndSet(currentAddr, newAddress)) {
+ if (!StringUtils.equals(newAddress, currentAddr) && masterHaAddress.compareAndSet(currentAddr, newAddress)) {
LOGGER.info("update master ha address, OLD: " + currentAddr + " NEW: " + newAddress);
wakeup();
}
@@ -253,10 +253,10 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
// Original state
this.handshakeHeaderBuffer.putInt(HAConnectionState.HANDSHAKE.ordinal());
// IsSyncFromLastFile
- short isSyncFromLastFile = this.haService.getDefaultMessageStore().getMessageStoreConfig().isSyncFromLastFile() ? (short)1 : (short) 0;
+ short isSyncFromLastFile = this.haService.getDefaultMessageStore().getMessageStoreConfig().isSyncFromLastFile() ? (short) 1 : (short) 0;
this.handshakeHeaderBuffer.putShort(isSyncFromLastFile);
// IsAsyncLearner role
- short isAsyncLearner = this.haService.getDefaultMessageStore().getMessageStoreConfig().isAsyncLearner() ? (short)1 : (short) 0;
+ short isAsyncLearner = this.haService.getDefaultMessageStore().getMessageStoreConfig().isAsyncLearner() ? (short) 1 : (short) 0;
this.handshakeHeaderBuffer.putShort(isAsyncLearner);
// Address length
this.handshakeHeaderBuffer.putInt(this.localAddress == null ? 0 : this.localAddress.length());
@@ -437,7 +437,7 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
while (true) {
int diff = byteBufferRead.position() - AutoSwitchHAClient.this.processPosition;
if (diff >= AutoSwitchHAConnection.MSG_HEADER_SIZE) {
- int processPosition = AutoSwitchHAClient.this.processPosition;
+ int processPosition = AutoSwitchHAClient.this.processPosition;
int masterState = byteBufferRead.getInt(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 36);
int bodySize = byteBufferRead.getInt(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 32);
long masterOffset = byteBufferRead.getLong(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 28);
@@ -452,8 +452,6 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
masterState, AutoSwitchHAClient.this.currentState, bodySize, masterOffset, masterEpoch, masterEpochStartOffset, confirmOffset);
return true;
}
- LOGGER.info("Receive master msg, masterState:{}, bodySize:{}, offset:{}, masterEpoch:{}, masterEpochStartOffset:{}, confirmOffset:{}",
- HAConnectionState.values()[masterState], bodySize, masterOffset, masterEpoch, masterEpochStartOffset, confirmOffset);
if (diff >= (AutoSwitchHAConnection.MSG_HEADER_SIZE + bodySize)) {
switch (AutoSwitchHAClient.this.currentState) {
@@ -501,10 +499,7 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
AutoSwitchHAClient.this.confirmOffset = Math.min(confirmOffset, messageStore.getMaxPhyOffset());
if (bodySize > 0) {
- final DefaultMessageStore messageStore = AutoSwitchHAClient.this.messageStore;
- if (messageStore.appendToCommitLog(masterOffset, bodyData, 0, bodyData.length)) {
- LOGGER.info("Slave append master log success, from {}, size {}, epoch:{}", masterOffset, bodySize, masterEpoch);
- }
+ AutoSwitchHAClient.this.messageStore.appendToCommitLog(masterOffset, bodyData, 0, bodyData.length);
}
if (!reportSlaveMaxOffset()) {
diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
index ef03dc583a..a253b4cb1d 100644
--- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
+++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
@@ -503,10 +503,14 @@ public class AutoSwitchHAConnection implements HAConnection {
// Normal transfer method
private void buildTransferHeaderBuffer(long nextOffset, int bodySize) {
- final EpochEntry entry = AutoSwitchHAConnection.this.epochCache.getEntry(AutoSwitchHAConnection.this.currentTransferEpoch);
+ EpochEntry entry = AutoSwitchHAConnection.this.epochCache.getEntry(AutoSwitchHAConnection.this.currentTransferEpoch);
if (entry == null) {
LOGGER.error("Failed to find epochEntry with epoch {} when build msg header", AutoSwitchHAConnection.this.currentTransferEpoch);
- return;
+ if (bodySize > 0) {
+ return;
+ }
+ // Maybe it's used for heartbeat
+ entry = AutoSwitchHAConnection.this.epochCache.firstEntry();
}
// Build Header
this.byteBufferHeader.position(0);
@@ -529,20 +533,23 @@ public class AutoSwitchHAConnection implements HAConnection {
currentState, bodySize, nextOffset, entry.getEpoch(), entry.getStartOffset(), confirmOffset);
}
+ private boolean sendHeartbeatIfNeeded() throws Exception {
+ long interval = haService.getDefaultMessageStore().getSystemClock().now() - this.lastWriteTimestamp;
+ if (interval > haService.getDefaultMessageStore().getMessageStoreConfig().getHaSendHeartbeatInterval()) {
+ buildTransferHeaderBuffer(this.nextTransferFromWhere, 0);
+ return this.transferData(0);
+ }
+ return true;
+ }
+
private void transferToSlave() throws Exception {
if (this.lastWriteOver) {
long interval =
haService.getDefaultMessageStore().getSystemClock().now() - this.lastWriteTimestamp;
- if (interval > haService.getDefaultMessageStore().getMessageStoreConfig()
- .getHaSendHeartbeatInterval()) {
-
- buildTransferHeaderBuffer(this.nextTransferFromWhere, 0);
-
- this.lastWriteOver = this.transferData(0);
- if (!this.lastWriteOver) {
- return;
- }
+ this.lastWriteOver = sendHeartbeatIfNeeded();
+ if (!this.lastWriteOver) {
+ return;
}
} else {
// maxTransferSize == -1 means to continue transfer remaining data.
@@ -596,6 +603,8 @@ public class AutoSwitchHAConnection implements HAConnection {
this.lastWriteOver = this.transferData(size);
} else {
+ // If size == 0, we should update the lastCatchupTimeMs
+ AutoSwitchHAConnection.this.lastCatchUpTimeMs = System.currentTimeMillis();
haService.getWaitNotifyObject().allWaitForRunning(100);
}
}
@@ -657,6 +666,7 @@ public class AutoSwitchHAConnection implements HAConnection {
EpochEntry epochEntry = AutoSwitchHAConnection.this.epochCache.findEpochEntryByOffset(this.nextTransferFromWhere);
if (epochEntry == null) {
LOGGER.error("Failed to find an epochEntry to match slaveRequestOffset {}", this.nextTransferFromWhere);
+ sendHeartbeatIfNeeded();
waitForRunning(500);
break;
}
diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java
index 98b480d914..9417600827 100644
--- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java
+++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java
@@ -35,7 +35,6 @@ import org.apache.rocketmq.logging.InternalLoggerFactory;
import org.apache.rocketmq.store.DefaultMessageStore;
import org.apache.rocketmq.store.DispatchRequest;
import org.apache.rocketmq.store.SelectMappedBufferResult;
-import org.apache.rocketmq.store.config.BrokerRole;
import org.apache.rocketmq.store.ha.DefaultHAService;
import org.apache.rocketmq.store.ha.GroupTransferService;
import org.apache.rocketmq.store.ha.HAClient;
@@ -64,9 +63,6 @@ public class AutoSwitchHAService extends DefaultHAService {
this.defaultMessageStore = defaultMessageStore;
this.acceptSocketService = new AutoSwitchAcceptSocketService(defaultMessageStore.getMessageStoreConfig().getHaListenPort());
this.groupTransferService = new GroupTransferService(this, defaultMessageStore);
- if (this.defaultMessageStore.getMessageStoreConfig().getBrokerRole() == BrokerRole.SLAVE) {
- this.haClient = new AutoSwitchHAClient(this, defaultMessageStore, this.epochCache);
- }
this.haConnectionStateNotificationService = new HAConnectionStateNotificationService(this, defaultMessageStore);
}
@@ -80,7 +76,7 @@ public class AutoSwitchHAService extends DefaultHAService {
@Override public boolean changeToMaster(int masterEpoch) {
final int lastEpoch = this.epochCache.lastEpoch();
- if (masterEpoch <= lastEpoch) {
+ if (masterEpoch < lastEpoch) {
return false;
}
destroyConnections();
@@ -172,6 +168,9 @@ public class AutoSwitchHAService extends DefaultHAService {
final AutoSwitchHAConnection connection = (AutoSwitchHAConnection) haConnection;
final String slaveAddress = connection.getSlaveAddress();
if (currentSyncStateSet.contains(slaveAddress)) {
+ if (this.defaultMessageStore.getMaxPhyOffset() == connection.getSlaveAckOffset()) {
+ continue;
+ }
if ((System.currentTimeMillis() - connection.getLastCatchUpTimeMs()) > haMaxTimeSlaveNotCatchup) {
newSyncStateSet.remove(slaveAddress);
}
diff --git a/test/src/test/java/org/apache/rocketmq/test/autoswitchrole/AutoSwitchRoleBase.java b/test/src/test/java/org/apache/rocketmq/test/autoswitchrole/AutoSwitchRoleBase.java
index 0703c65a75..067e6424d0 100644
--- a/test/src/test/java/org/apache/rocketmq/test/autoswitchrole/AutoSwitchRoleBase.java
+++ b/test/src/test/java/org/apache/rocketmq/test/autoswitchrole/AutoSwitchRoleBase.java
@@ -123,7 +123,7 @@ public class AutoSwitchRoleBase {
protected ControllerConfig buildControllerConfig(final String id, final String peers) {
final ControllerConfig config = new ControllerConfig();
- config.setStartupController(true);
+ config.setEnableStartupController(true);
config.setControllerDLegerGroup("group1");
config.setControllerDLegerPeers(peers);
config.setControllerDLegerSelfId(id);