When the size of syncStateSet is less than minInSyncReplicas, put message returns PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH

This commit is contained in:
RongtongJin
2022-07-04 16:45:01 +08:00
parent 505cb2af9b
commit e5c538a25e
9 changed files with 75 additions and 31 deletions
@@ -803,12 +803,17 @@ public class CommitLog implements Swappable {
int needAckNums = 1;
if (needHandleHA) {
if (this.defaultMessageStore.getBrokerConfig().isEnableControllerMode() && this.defaultMessageStore.getMessageStoreConfig().isAllAckInSyncStateSet()) {
// -1 means all ack in SyncStateSet
needAckNums = MixAll.ALL_ACK_IN_SYNC_STATE_SET;
if (this.defaultMessageStore.getBrokerConfig().isEnableControllerMode()) {
if (this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset) < this.defaultMessageStore.getMessageStoreConfig().getMinInSyncReplicas()) {
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, null));
}
if (this.defaultMessageStore.getMessageStoreConfig().isAllAckInSyncStateSet()) {
// -1 means all ack in SyncStateSet
needAckNums = MixAll.ALL_ACK_IN_SYNC_STATE_SET;
}
} else {
int inSyncReplicas = Math.min(this.defaultMessageStore.getAliveReplicaNumInGroup(),
this.defaultMessageStore.getHaService().inSyncSlaveNums(currOffset) + 1);
this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset));
needAckNums = calcNeedAckNums(inSyncReplicas);
if (needAckNums > inSyncReplicas) {
// Tell the producer, don't have enough slaves to handle the send request
@@ -955,12 +960,17 @@ public class CommitLog implements Swappable {
boolean needHandleHA = needHandleHA(messageExtBatch);
if (needHandleHA) {
if (this.defaultMessageStore.getBrokerConfig().isEnableControllerMode() && this.defaultMessageStore.getMessageStoreConfig().isAllAckInSyncStateSet()) {
// -1 means all ack in SyncStateSet
needAckNums = MixAll.ALL_ACK_IN_SYNC_STATE_SET;
if (this.defaultMessageStore.getBrokerConfig().isEnableControllerMode()) {
if (this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset) < this.defaultMessageStore.getMessageStoreConfig().getMinInSyncReplicas()) {
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, null));
}
if (this.defaultMessageStore.getMessageStoreConfig().isAllAckInSyncStateSet()) {
// -1 means all ack in SyncStateSet
needAckNums = MixAll.ALL_ACK_IN_SYNC_STATE_SET;
}
} else {
int inSyncReplicas = Math.min(this.defaultMessageStore.getAliveReplicaNumInGroup(),
this.defaultMessageStore.getHaService().inSyncSlaveNums(currOffset) + 1);
this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset));
needAckNums = calcNeedAckNums(inSyncReplicas);
if (needAckNums > inSyncReplicas) {
// Tell the producer, don't have enough slaves to handle the send request
@@ -1117,17 +1127,6 @@ public class CommitLog implements Swappable {
HAService haService = this.defaultMessageStore.getHaService();
long nextOffset = result.getWroteOffset() + result.getWroteBytes();
// NOTE: Plus the master replicas
// int inSyncReplicas = haService.inSyncSlaveNums(nextOffset) + 1;
// if (needAckNums > inSyncReplicas) {
// /*
// * Tell the producer, don't have enough slaves to handle the send request.
// * NOTE: this may cause msg duplicate
// */
// putMessageResult.setPutMessageStatus(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH);
// return CompletableFuture.completedFuture(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH);
// }
// Wait enough acks from different slaves
GroupCommitRequest request = new GroupCommitRequest(nextOffset, this.defaultMessageStore.getMessageStoreConfig().getSlaveTimeout(), needAckNums);
@@ -236,12 +236,14 @@ public class MessageStoreConfig {
* Each message must be written successfully to at least in-sync replicas.
* The master broker is considered one of the in-sync replicas, and it's included in the count of total.
* If a master broker is ASYNC_MASTER, inSyncReplicas will be ignored.
* If enableControllerMode is true and ackAckInSyncStateSet is true, inSyncReplicas will be ignored.
*/
@ImportantField
private int inSyncReplicas = 1;
/**
* Will be worked in auto multiple replicas mode, to provide minimum in-sync replicas.
* It is still valid in controller mode.
*/
@ImportantField
private int minInSyncReplicas = 1;
@@ -300,7 +302,7 @@ public class MessageStoreConfig {
private long maxSlaveResendLength = 256 * 1024 * 1024;
/**
* Whether sync from lastFile when a new broker replicas join the master.
* Whether sync from lastFile when a new broker replicas(no data) join the master.
*/
private boolean syncFromLastFile = false;
@@ -178,8 +178,8 @@ public class DefaultHAService implements HAService {
return push2SlaveMaxOffset;
}
public int inSyncSlaveNums(final long masterPutWhere) {
int inSyncNums = 0;
public int inSyncReplicasNums(final long masterPutWhere) {
int inSyncNums = 1;
for (HAConnection conn : this.connectionList) {
if (this.isInSyncSlave(masterPutWhere, conn)) {
inSyncNums++;
@@ -96,6 +96,7 @@ public class GroupTransferService extends ServiceThread {
transferOK = true;
break;
}
// Include master
int ackNums = 1;
for (HAConnection conn : haService.getConnectionList()) {
@@ -76,12 +76,13 @@ public interface HAService {
void updateHaMasterAddress(String newAddr);
/**
* Returns the number of slaves those commit log are not far behind the master.
* Returns the number of replicas those commit log are not far behind the master. It includes master itself.
* Returns syncStateSet size if HAService instanceof AutoSwitchService
*
* @return the number of slaves
* @see MessageStoreConfig#getHaMaxGapNotInSync()
*/
int inSyncSlaveNums(long masterPutWhere);
int inSyncReplicasNums(long masterPutWhere);
/**
* Get connection count
@@ -32,6 +32,7 @@ import java.util.function.Consumer;
import org.apache.rocketmq.common.EpochEntry;
import org.apache.rocketmq.common.ThreadFactoryImpl;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.common.protocol.body.HARuntimeInfo;
import org.apache.rocketmq.logging.InternalLogger;
import org.apache.rocketmq.logging.InternalLoggerFactory;
import org.apache.rocketmq.store.DefaultMessageStore;
@@ -253,6 +254,46 @@ public class AutoSwitchHAService extends DefaultHAService {
}
}
public int inSyncReplicasNums(final long masterPutWhere) {
return syncStateSet.size();
}
@Override public HARuntimeInfo getRuntimeInfo(long masterPutWhere) {
HARuntimeInfo info = new HARuntimeInfo();
if (BrokerRole.SLAVE.equals(this.getDefaultMessageStore().getMessageStoreConfig().getBrokerRole())) {
info.setMaster(false);
info.getHaClientRuntimeInfo().setMasterAddr(this.haClient.getHaMasterAddress());
info.getHaClientRuntimeInfo().setMaxOffset(this.getDefaultMessageStore().getMaxPhyOffset());
info.getHaClientRuntimeInfo().setLastReadTimestamp(this.haClient.getLastReadTimestamp());
info.getHaClientRuntimeInfo().setLastWriteTimestamp(this.haClient.getLastWriteTimestamp());
info.getHaClientRuntimeInfo().setTransferredByteInSecond(this.haClient.getTransferredByteInSecond());
info.getHaClientRuntimeInfo().setMasterFlushOffset(this.defaultMessageStore.getMasterFlushedOffset());
} else {
info.setMaster(true);
info.setMasterCommitLogMaxOffset(masterPutWhere);
for (HAConnection conn : this.connectionList) {
HARuntimeInfo.HAConnectionRuntimeInfo cInfo = new HARuntimeInfo.HAConnectionRuntimeInfo();
long slaveAckOffset = conn.getSlaveAckOffset();
cInfo.setSlaveAckOffset(slaveAckOffset);
cInfo.setDiff(masterPutWhere - slaveAckOffset);
cInfo.setAddr(conn.getClientAddress().substring(1));
cInfo.setTransferredByteInSecond(conn.getTransferredByteInSecond());
cInfo.setTransferFromWhere(conn.getTransferFromWhere());
cInfo.setInSync(syncStateSet.contains(((AutoSwitchHAConnection) conn).getSlaveAddress()));
info.getHaConnectionInfo().add(cInfo);
}
info.setInSyncSlaveNums(syncStateSet.size() - 1);
}
return info;
}
public void updateConfirmOffset(long confirmOffset) {
this.confirmOffset = confirmOffset;
}
@@ -114,7 +114,7 @@ public class HAServerTest {
}
@Test
public void inSyncSlaveNums() throws IOException {
public void inSyncReplicasNums() throws IOException {
DefaultMessageStore messageStore = mockMessageStore();
doReturn(123L).when(messageStore).getMaxPhyOffset();
doReturn(123L).when(messageStore).getMasterFlushedOffset();
@@ -140,13 +140,13 @@ public class HAServerTest {
await().atMost(Duration.ofMinutes(1)).until(new Callable<Boolean>() {
@Override
public Boolean call() throws Exception {
return HAServerTest.this.haService.inSyncSlaveNums(haSlaveFallbehindMax) == 4;
return HAServerTest.this.haService.inSyncReplicasNums(haSlaveFallbehindMax) == 5;
}
});
assertThat(HAServerTest.this.haService.inSyncSlaveNums(123L + haSlaveFallbehindMax)).isEqualTo(2);
assertThat(HAServerTest.this.haService.inSyncSlaveNums(124L + haSlaveFallbehindMax)).isEqualTo(1);
assertThat(HAServerTest.this.haService.inSyncSlaveNums(125L + haSlaveFallbehindMax)).isEqualTo(0);
assertThat(HAServerTest.this.haService.inSyncReplicasNums(123L + haSlaveFallbehindMax)).isEqualTo(3);
assertThat(HAServerTest.this.haService.inSyncReplicasNums(124L + haSlaveFallbehindMax)).isEqualTo(2);
assertThat(HAServerTest.this.haService.inSyncReplicasNums(125L + haSlaveFallbehindMax)).isEqualTo(1);
}
@Test
@@ -166,7 +166,7 @@ public class PullMultipleReplicasIT extends ContainerIntegrationTestBase {
await().atMost(Duration.ofSeconds(60)).until(() -> {
DefaultMessageStore messageStore = (DefaultMessageStore) master3With3Replicas.getMessageStore();
return messageStore.getHaService().inSyncSlaveNums(messageStore.getMaxPhyOffset()) == 2;
return messageStore.getHaService().inSyncReplicasNums(messageStore.getMaxPhyOffset()) == 3;
});
InnerSalveBrokerController slaveBroker = null;
@@ -110,7 +110,7 @@ public class SlaveBrokerIT extends ContainerIntegrationTestBase {
.until(() -> ((DefaultMessageStore) master3With3Replicas.getMessageStore()).getHaService().getConnectionCount().get() == 2);
await().atMost(100, TimeUnit.SECONDS)
.until(() -> ((DefaultMessageStore) master3With3Replicas.getMessageStore()).getHaService().inSyncSlaveNums(0) == 2);
.until(() -> ((DefaultMessageStore) master3With3Replicas.getMessageStore()).getHaService().inSyncReplicasNums(0) == 3);
Thread.sleep(1000 * 101);
}