From e5c538a25e7239c5bbd7c37b3d759cae832f387f Mon Sep 17 00:00:00 2001 From: RongtongJin Date: Mon, 4 Jul 2022 16:45:01 +0800 Subject: [PATCH] When the size of syncStateSet is less than minInSyncReplicas, put message returns PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH --- .../org/apache/rocketmq/store/CommitLog.java | 37 ++++++++--------- .../store/config/MessageStoreConfig.java | 4 +- .../rocketmq/store/ha/DefaultHAService.java | 4 +- .../store/ha/GroupTransferService.java | 1 + .../apache/rocketmq/store/ha/HAService.java | 5 ++- .../ha/autoswitch/AutoSwitchHAService.java | 41 +++++++++++++++++++ .../rocketmq/store/ha/HAServerTest.java | 10 ++--- .../container/PullMultipleReplicasIT.java | 2 +- .../test/container/SlaveBrokerIT.java | 2 +- 9 files changed, 75 insertions(+), 31 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java index 75213b83a7..f1acff1fee 100644 --- a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java +++ b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java @@ -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); 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 63e35c5990..7c7989f50c 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 @@ -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; diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java index ea27288bea..32945c6d02 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/DefaultHAService.java @@ -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++; diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java b/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java index f84bdbf052..cd3e50adf4 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/GroupTransferService.java @@ -96,6 +96,7 @@ public class GroupTransferService extends ServiceThread { transferOK = true; break; } + // Include master int ackNums = 1; for (HAConnection conn : haService.getConnectionList()) { diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/HAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/HAService.java index 4006ff774b..c321f0a58f 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/HAService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/HAService.java @@ -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 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 775ce9b6b5..1c0a471ff9 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 @@ -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; } diff --git a/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java b/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java index 5304bec467..a8ce5179dc 100644 --- a/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/ha/HAServerTest.java @@ -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() { @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 diff --git a/test/src/test/java/org/apache/rocketmq/test/container/PullMultipleReplicasIT.java b/test/src/test/java/org/apache/rocketmq/test/container/PullMultipleReplicasIT.java index 94e056aeda..02578b1599 100644 --- a/test/src/test/java/org/apache/rocketmq/test/container/PullMultipleReplicasIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/container/PullMultipleReplicasIT.java @@ -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; diff --git a/test/src/test/java/org/apache/rocketmq/test/container/SlaveBrokerIT.java b/test/src/test/java/org/apache/rocketmq/test/container/SlaveBrokerIT.java index e6449c9ee0..1071f3495d 100644 --- a/test/src/test/java/org/apache/rocketmq/test/container/SlaveBrokerIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/container/SlaveBrokerIT.java @@ -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); }