mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-29 04:13:00 +08:00
[Summer of code] Support async learner in controller mode (#4367)
* merge branch support_async_learner
* use isSlave() to replace BrokerRole == Slave
* mark asyncLearner
* code review
* Revert "use isSlave() to replace BrokerRole == Slave"
This reverts commit 6599f97f44.
* review
* remove asyncLeaner role
* code review
* code review
This commit is contained in:
@@ -306,6 +306,8 @@ public class MessageStoreConfig {
|
||||
*/
|
||||
private boolean syncFromLastFile = false;
|
||||
|
||||
private boolean isAsyncLearner = false;
|
||||
|
||||
public boolean isDebugLockEnable() {
|
||||
return debugLockEnable;
|
||||
}
|
||||
@@ -1322,4 +1324,12 @@ public class MessageStoreConfig {
|
||||
public void setScheduleAsyncDeliverMaxResendNum2Blocked(int scheduleAsyncDeliverMaxResendNum2Blocked) {
|
||||
this.scheduleAsyncDeliverMaxResendNum2Blocked = scheduleAsyncDeliverMaxResendNum2Blocked;
|
||||
}
|
||||
|
||||
public boolean isAsyncLearner() {
|
||||
return isAsyncLearner;
|
||||
}
|
||||
|
||||
public void setAsyncLearner(boolean asyncLearner) {
|
||||
this.isAsyncLearner = asyncLearner;
|
||||
}
|
||||
}
|
||||
|
||||
+14
-19
@@ -44,9 +44,10 @@ import org.apache.rocketmq.store.ha.io.HAWriter;
|
||||
public class AutoSwitchHAClient extends ServiceThread implements HAClient {
|
||||
|
||||
/**
|
||||
* Handshake header buffer size. Schema: state ordinal + flag(isSyncFromLastFile) + slaveId + slaveAddressLength.
|
||||
* Handshake header buffer size. Schema: state ordinal + Two flags + slaveAddressLength
|
||||
* Flag: isSyncFromLastFile(short), isAsyncLearner(short)... we can add more flags in the future if needed
|
||||
*/
|
||||
public static final int HANDSHAKE_HEADER_SIZE = 4 + 4 + 8 + 4;
|
||||
public static final int HANDSHAKE_HEADER_SIZE = 4 + 4 + 4;
|
||||
|
||||
/**
|
||||
* Header + slaveAddress.
|
||||
@@ -97,10 +98,6 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
|
||||
*/
|
||||
private volatile long confirmOffset;
|
||||
|
||||
public static final int SYNC_FROM_LAST_FILE = -1;
|
||||
|
||||
public static final int SYNC_FROM_FIRST_FILE = -2;
|
||||
|
||||
public AutoSwitchHAClient(AutoSwitchHAService haService, DefaultMessageStore defaultMessageStore,
|
||||
EpochFileCache epochCache) throws IOException {
|
||||
this.haService = haService;
|
||||
@@ -256,13 +253,11 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
|
||||
// Original state
|
||||
this.handshakeHeaderBuffer.putInt(HAConnectionState.HANDSHAKE.ordinal());
|
||||
// IsSyncFromLastFile
|
||||
if (this.haService.getDefaultMessageStore().getMessageStoreConfig().isSyncFromLastFile()) {
|
||||
this.handshakeHeaderBuffer.putInt(SYNC_FROM_LAST_FILE);
|
||||
} else {
|
||||
this.handshakeHeaderBuffer.putInt(SYNC_FROM_FIRST_FILE);
|
||||
}
|
||||
// Slave Id
|
||||
this.handshakeHeaderBuffer.putLong(this.slaveId.get());
|
||||
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;
|
||||
this.handshakeHeaderBuffer.putShort(isAsyncLearner);
|
||||
// Address length
|
||||
this.handshakeHeaderBuffer.putInt(this.localAddress == null ? 0 : this.localAddress.length());
|
||||
// Slave address
|
||||
@@ -443,12 +438,12 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
|
||||
int diff = byteBufferRead.position() - AutoSwitchHAClient.this.processPosition;
|
||||
if (diff >= AutoSwitchHAConnection.MSG_HEADER_SIZE) {
|
||||
int processPosition = AutoSwitchHAClient.this.processPosition;
|
||||
int masterState = byteBufferRead.getInt(processPosition);
|
||||
int bodySize = byteBufferRead.getInt(processPosition + 4);
|
||||
long masterOffset = byteBufferRead.getLong(processPosition + 4 + 4);
|
||||
int masterEpoch = byteBufferRead.getInt(processPosition + 4 + 4 + 8);
|
||||
long masterEpochStartOffset = byteBufferRead.getLong(processPosition + 4 + 4 + 8 + 4);
|
||||
long confirmOffset = byteBufferRead.getLong(processPosition + 4 + 4 + 8 + 4 + 8);
|
||||
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);
|
||||
int masterEpoch = byteBufferRead.getInt(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 20);
|
||||
long masterEpochStartOffset = byteBufferRead.getLong(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 16);
|
||||
long confirmOffset = byteBufferRead.getLong(processPosition + AutoSwitchHAConnection.MSG_HEADER_SIZE - 8);
|
||||
|
||||
if (masterState != AutoSwitchHAClient.this.currentState.ordinal()) {
|
||||
AutoSwitchHAClient.this.processPosition += AutoSwitchHAConnection.MSG_HEADER_SIZE + bodySize;
|
||||
|
||||
+34
-22
@@ -64,6 +64,7 @@ public class AutoSwitchHAConnection implements HAConnection {
|
||||
private volatile int currentTransferEpoch = -1;
|
||||
private volatile long currentTransferEpochEndOffset = 0;
|
||||
private volatile boolean isSyncFromLastFile = false;
|
||||
private volatile boolean isAsyncLearner = false;
|
||||
private volatile long slaveId = -1;
|
||||
private volatile String slaveAddress;
|
||||
|
||||
@@ -172,13 +173,21 @@ public class AutoSwitchHAConnection implements HAConnection {
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isAsyncLearner() {
|
||||
return isAsyncLearner;
|
||||
}
|
||||
|
||||
public boolean isSyncFromLastFile() {
|
||||
return isSyncFromLastFile;
|
||||
}
|
||||
|
||||
private synchronized void updateLastTransferInfo() {
|
||||
this.lastMasterMaxOffset = this.haService.getDefaultMessageStore().getMaxPhyOffset();
|
||||
this.lastTransferTimeMs = System.currentTimeMillis();
|
||||
}
|
||||
|
||||
private synchronized void maybeExpandInSyncStateSet(long slaveMaxOffset) {
|
||||
if (slaveMaxOffset >= this.lastMasterMaxOffset) {
|
||||
if (!this.isAsyncLearner && slaveMaxOffset >= this.lastMasterMaxOffset) {
|
||||
this.lastCatchUpTimeMs = Math.max(this.lastTransferTimeMs, this.lastCatchUpTimeMs);
|
||||
this.haService.maybeExpandInSyncStateSet(this.slaveAddress, slaveMaxOffset);
|
||||
}
|
||||
@@ -277,30 +286,33 @@ public class AutoSwitchHAConnection implements HAConnection {
|
||||
switch (slaveState) {
|
||||
case HANDSHAKE:
|
||||
if (diff >= AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE) {
|
||||
isSlaveSendHandshake = true;
|
||||
// Flag(SyncFromLastFile)
|
||||
long syncFromLastFileFlag = byteBufferRead.getInt(readPosition + 4);
|
||||
if (syncFromLastFileFlag == AutoSwitchHAClient.SYNC_FROM_LAST_FILE) {
|
||||
AutoSwitchHAConnection.this.isSyncFromLastFile = true;
|
||||
LOGGER.info("Slave request sync from lastFile");
|
||||
}
|
||||
// SlaveId
|
||||
AutoSwitchHAConnection.this.slaveId = byteBufferRead.getLong(readPosition + 8);
|
||||
// AddressLength
|
||||
int addressLength = byteBufferRead.getInt(readPosition + 16);
|
||||
// Address
|
||||
if (diff >= AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE + addressLength) {
|
||||
final byte[] addressData = new byte[addressLength];
|
||||
byteBufferRead.position(readPosition + 20);
|
||||
byteBufferRead.get(addressData);
|
||||
AutoSwitchHAConnection.this.slaveAddress = new String(addressData);
|
||||
} else {
|
||||
AutoSwitchHAConnection.this.slaveAddress = "";
|
||||
int addressLength = byteBufferRead.getInt(readPosition + AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE - 4);
|
||||
if (diff < AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE + addressLength) {
|
||||
break;
|
||||
}
|
||||
LOGGER.info("Receive slave handshake, syncFromLastFile:{}, slaveId:{}, slaveAddress:{}",
|
||||
AutoSwitchHAConnection.this.isSyncFromLastFile, AutoSwitchHAConnection.this.slaveId, AutoSwitchHAConnection.this.slaveAddress);
|
||||
// Flag(isSyncFromLastFile)
|
||||
short syncFromLastFileFlag = byteBufferRead.getShort(readPosition + AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE - 8);
|
||||
if (syncFromLastFileFlag == 1) {
|
||||
AutoSwitchHAConnection.this.isSyncFromLastFile = true;
|
||||
}
|
||||
// Flag(isAsyncLearner role)
|
||||
short isAsyncLearner = byteBufferRead.getShort(readPosition + AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE - 6);
|
||||
if (isAsyncLearner == 1) {
|
||||
AutoSwitchHAConnection.this.isAsyncLearner = true;
|
||||
}
|
||||
// Address
|
||||
final byte[] addressData = new byte[addressLength];
|
||||
byteBufferRead.position(readPosition + AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE);
|
||||
byteBufferRead.get(addressData);
|
||||
AutoSwitchHAConnection.this.slaveAddress = new String(addressData);
|
||||
|
||||
isSlaveSendHandshake = true;
|
||||
byteBufferRead.position(readSocketPos);
|
||||
ReadSocketService.this.processPosition += AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE + addressLength;
|
||||
LOGGER.info("Receive slave handshake, slaveId:{}, slaveAddress:{}, isSyncFromLastFile:{}, isAsyncLearner:{}",
|
||||
AutoSwitchHAConnection.this.slaveId, AutoSwitchHAConnection.this.slaveAddress,
|
||||
AutoSwitchHAConnection.this.isSyncFromLastFile, AutoSwitchHAConnection.this.isAsyncLearner);
|
||||
}
|
||||
break;
|
||||
case TRANSFER:
|
||||
@@ -312,10 +324,10 @@ public class AutoSwitchHAConnection implements HAConnection {
|
||||
if (slaveRequestOffset < 0) {
|
||||
slaveRequestOffset = slaveMaxOffset;
|
||||
}
|
||||
LOGGER.info("slave[" + clientAddress + "] request offset " + slaveMaxOffset);
|
||||
byteBufferRead.position(readSocketPos);
|
||||
maybeExpandInSyncStateSet(slaveMaxOffset);
|
||||
AutoSwitchHAConnection.this.haService.notifyTransferSome(AutoSwitchHAConnection.this.slaveAckOffset);
|
||||
LOGGER.info("slave[" + clientAddress + "] request offset " + slaveMaxOffset);
|
||||
}
|
||||
break;
|
||||
default:
|
||||
|
||||
@@ -45,6 +45,7 @@ import org.junit.Test;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
public class AutoSwitchHATest {
|
||||
@@ -184,6 +185,33 @@ public class AutoSwitchHATest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAsyncLearnerBrokerRole() throws Exception {
|
||||
init(defaultMappedFileSize);
|
||||
((AutoSwitchHAService) this.messageStore1.getHaService()).setLocalAddress("127.0.0.1:8000");
|
||||
((AutoSwitchHAService) this.messageStore2.getHaService()).setLocalAddress("127.0.0.1:8001");
|
||||
|
||||
storeConfig1.setBrokerRole(BrokerRole.SYNC_MASTER);
|
||||
storeConfig2.setBrokerRole(BrokerRole.SLAVE);
|
||||
storeConfig2.setAsyncLearner(true);
|
||||
messageStore1.getHaService().changeToMaster(1);
|
||||
messageStore2.getHaService().changeToSlave("", 1, 2L);
|
||||
messageStore2.getHaService().updateHaMasterAddress(store1HaAddress);
|
||||
Thread.sleep(6000);
|
||||
|
||||
// Put message on master
|
||||
for (int i = 0; i < 10; i++) {
|
||||
messageStore1.putMessage(buildMessage());
|
||||
}
|
||||
Thread.sleep(200);
|
||||
|
||||
checkMessage(messageStore2, 10, 0);
|
||||
|
||||
Thread.sleep(1000);
|
||||
final Set<String> syncStateSet = ((AutoSwitchHAService) this.messageStore1.getHaService()).getSyncStateSet();
|
||||
assertFalse(syncStateSet.contains("127.0.0.1:8001"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOptionAllAckInSyncStateSet() throws Exception {
|
||||
init(defaultMappedFileSize, true);
|
||||
|
||||
Reference in New Issue
Block a user