From 96d883286880b2ca16af677140b4adede85af9f7 Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 16:28:57 +0800 Subject: [PATCH 01/10] fix --- .../apache/rocketmq/store/ConsumeQueue.java | 14 +- .../rocketmq/store/DefaultMessageStore.java | 21 ++- .../store/MessageStoreStateMachine.java | 109 ++++++++++++++ .../store/MessageStoreStateMachineTest.java | 137 ++++++++++++++++++ 4 files changed, 275 insertions(+), 6 deletions(-) create mode 100644 store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java create mode 100644 store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index 79e5368216b..8a3fade7bde 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -61,7 +61,7 @@ public class ConsumeQueue implements ConsumeQueueInterface, FileQueueLifeCycle { public static final int MSG_TAG_OFFSET_INDEX = 12; private static final Logger LOG_ERROR = LoggerFactory.getLogger(LoggerName.STORE_ERROR_LOGGER_NAME); - private final MessageStore messageStore; + private final DefaultMessageStore messageStore; private final ConsumeQueueStore consumeQueueStore; private final MappedFileQueue mappedFileQueue; @@ -80,12 +80,12 @@ public class ConsumeQueue implements ConsumeQueueInterface, FileQueueLifeCycle { private ConsumeQueueExt consumeQueueExt = null; public ConsumeQueue(final String topic, final int queueId, final String storePath, final int mappedFileSize, - final MessageStore messageStore) { + final DefaultMessageStore messageStore) { this(topic, queueId, storePath, mappedFileSize, messageStore, (ConsumeQueueStore) messageStore.getQueueStore()); } public ConsumeQueue(final String topic, final int queueId, final String storePath, final int mappedFileSize, - final MessageStore messageStore, final ConsumeQueueStore consumeQueueStore) { + final DefaultMessageStore messageStore, final ConsumeQueueStore consumeQueueStore) { this.storePath = storePath; this.mappedFileSize = mappedFileSize; this.messageStore = messageStore; @@ -792,7 +792,13 @@ private boolean putMessagePositionInfo(final long offset, final int size, final final long cqOffset) { if (offset + size <= this.getMaxPhysicOffset()) { - log.warn("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", this.getMaxPhysicOffset(), offset); + if (messageStore.getStateMachine().getCurrentState().isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK)) { + log.warn("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", + this.getMaxPhysicOffset(), offset); + } else { + log.debug("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", + this.getMaxPhysicOffset(), offset); + } return true; } diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index ea0a814caf2..f8c3949a472 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -204,6 +204,8 @@ public class DefaultMessageStore implements MessageStore { // this is a unmodifiableMap private final ConcurrentMap topicConfigTable; + private final MessageStoreStateMachine stateMachine; + private final ScheduledExecutorService scheduledCleanQueueExecutorService = ThreadUtils.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("StoreCleanQueueScheduledThread")); @@ -250,6 +252,8 @@ public DefaultMessageStore(final MessageStoreConfig messageStoreConfig, final Br lockFile = new RandomAccessFile(file, "rw"); parseDelayLevel(); + + stateMachine = new MessageStoreStateMachine(LOGGER); } public ConsumeQueueStoreInterface createConsumeQueueStore() { @@ -304,17 +308,20 @@ public boolean load() { // load Commit Log result = this.commitLog.load(); - + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMMITLOG_OK); // load Consume Queue result = result && this.consumeQueueStore.load(); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_CONSUME_QUEUE_OK); if (messageStoreConfig.isEnableCompaction()) { result = result && this.compactionService.load(lastExitOK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMPACTION_OK); } if (result) { loadCheckPoint(); result = this.indexService.load(lastExitOK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_INDEX_OK); this.recover(lastExitOK); LOGGER.info("message store recover end, and the max phy offset = {}", this.getMaxPhyOffset()); } @@ -329,6 +336,7 @@ public boolean load() { if (!result) { this.allocateMappedFileService.shutdown(); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.ERROR); } return result; @@ -346,20 +354,23 @@ private void recover(final boolean lastExitOK) throws RocksDBException { // recover consume queue long recoverConsumeQueueStart = System.currentTimeMillis(); this.consumeQueueStore.recover(this.brokerConfig.isRecoverConcurrently()); - long dispatchFromPhyOffset = this.consumeQueueStore.getDispatchFromPhyOffset(); long recoverConsumeQueueEnd = System.currentTimeMillis(); + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_CONSUME_QUEUE_OK); // recover commitlog + long dispatchFromPhyOffset = this.consumeQueueStore.getDispatchFromPhyOffset(); if (lastExitOK) { this.commitLog.recoverNormally(dispatchFromPhyOffset); } else { this.commitLog.recoverAbnormally(dispatchFromPhyOffset); } + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK); // recover consume offset table long recoverCommitLogEnd = System.currentTimeMillis(); this.recoverTopicQueueTable(); long recoverConsumeOffsetEnd = System.currentTimeMillis(); + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_TOPIC_QUEUE_TABLE_OK); LOGGER.info("message store recover total cost: {} ms, " + "recoverConsumeQueue: {} ms, recoverCommitLog: {} ms, recoverOffsetTable: {} ms", @@ -411,6 +422,8 @@ public void start() throws Exception { this.addScheduleTask(); this.perfs.start(); this.shutdown = false; + + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RUNNING); } private void doRecheckReputOffsetFromCq() throws InterruptedException { @@ -3001,4 +3014,8 @@ public ScheduledExecutorService getScheduledCleanQueueExecutorService() { public void destroyConsumeQueueStore(boolean loadAfterDestroy) { consumeQueueStore.destroy(loadAfterDestroy); } + + public MessageStoreStateMachine getStateMachine() { + return stateMachine; + } } diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java new file mode 100644 index 00000000000..bd1f0eae7e1 --- /dev/null +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -0,0 +1,109 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.store; + +import org.apache.rocketmq.common.constant.LoggerName; +import org.apache.rocketmq.logging.org.slf4j.Logger; +import org.apache.rocketmq.logging.org.slf4j.LoggerFactory; + +public class MessageStoreStateMachine { + protected final Logger log; + + private MessageStoreState currentState; + private long lastStateChangeTimestamp; + private final long startTimestamp; + + public enum MessageStoreState { + INIT(0), + + LOAD_COMMITLOG_OK(10), + LOAD_CONSUME_QUEUE_OK(11), + LOAD_COMPACTION_OK(12), + LOAD_INDEX_OK(13), + + RECOVER_CONSUME_QUEUE_OK(20), + RECOVER_COMMITLOG_OK(21), + RECOVER_TOPIC_QUEUE_TABLE_OK(22), + + RUNNING(30), + SHUTDOWN(40), + + ERROR(Integer.MAX_VALUE); + + final int order; + + MessageStoreState(int order) { + this.order = order; + } + + public int getOrder() { + return order; + } + + public boolean isBefore(MessageStoreState storeState) { + return this.order < storeState.order; + } + + public boolean isAfter(MessageStoreState storeState) { + return this.order > storeState.order; + } + } + + + public MessageStoreStateMachine(Logger log) { + this.log = log == null ? LoggerFactory.getLogger(LoggerName.STORE_LOGGER_NAME) : log; + this.currentState = MessageStoreState.INIT; + this.startTimestamp = System.currentTimeMillis(); + this.lastStateChangeTimestamp = startTimestamp; + logStateChange(null, currentState); + } + + public void transitTo(MessageStoreState newState) { + if (!newState.isAfter(currentState)) { + throw new IllegalStateException( + String.format("Invalid state transition from %s to %s. Can only move forward.", + currentState, newState) + ); + } + + logStateChange(currentState, newState); + this.currentState = newState; + this.lastStateChangeTimestamp = System.currentTimeMillis(); + } + + private void logStateChange(MessageStoreState fromState, MessageStoreState toState) { + if (fromState == null) { + log.info("MessageStore initialized, state={}", toState); + } else { + log.info("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} " + + "ms", fromState, toState, getCurrentStateRunningTimeMs(), getTotalRunningTimeMs()); + } + } + + public MessageStoreState getCurrentState() { + return currentState; + } + + public long getTotalRunningTimeMs() { + return System.currentTimeMillis() - startTimestamp; + } + + public long getCurrentStateRunningTimeMs() { + return System.currentTimeMillis() - lastStateChangeTimestamp; + } +} diff --git a/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java new file mode 100644 index 00000000000..500cbc482fb --- /dev/null +++ b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java @@ -0,0 +1,137 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.store; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.anyLong; +import static org.mockito.Mockito.verify; + +import org.apache.rocketmq.logging.org.slf4j.Logger; +import org.apache.rocketmq.store.MessageStoreStateMachine.MessageStoreState; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +class MessageStoreStateMachineTest { + + private Logger mockLogger; + private MessageStoreStateMachine stateMachine; + + @BeforeEach + void setUp() { + // Mock Logger + mockLogger = Mockito.mock(Logger.class); + + // Initialize StateMachine + stateMachine = new MessageStoreStateMachine(mockLogger); + } + + /** + * Test the constructor of MessageStoreStateMachine. + */ + @Test + void testConstructor() { + // Verify initial state + assertEquals(MessageStoreState.INIT, stateMachine.getCurrentState()); + + // Verify logger was called for initialization + verify(mockLogger).info("MessageStore initialized, state={}", MessageStoreState.INIT); + } + + /** + * Test valid state transition in transitTo method. + */ + @Test + void testValidStateTransition() { + // Perform a valid state transition + stateMachine.transitTo(MessageStoreState.LOAD_COMMITLOG_OK); + + // Verify the current state is updated + assertEquals(MessageStoreState.LOAD_COMMITLOG_OK, stateMachine.getCurrentState()); + + // Verify logger was called for state transition + verify(mockLogger).info( + eq("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} ms"), + eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), anyLong(), anyLong() + ); + } + + /** + * Test invalid state transition in transitTo method. + */ + @Test + void testInvalidStateTransition() { + // Perform an invalid state transition + Exception exception = assertThrows(IllegalStateException.class, () -> { + stateMachine.transitTo(MessageStoreState.INIT); + }); + + // Verify the exception message + String expectedMessage = "Invalid state transition from INIT to INIT. Can only move forward."; + assertEquals(expectedMessage, exception.getMessage()); + } + + /** + * Test getCurrentState method. + */ + @Test + void testGetCurrentState() { + // Verify the current state + assertEquals(MessageStoreState.INIT, stateMachine.getCurrentState()); + } + + /** + * Test getTotalRunningTimeMs method. + */ + @Test + void testGetTotalRunningTimeMs() { + // Sleep for a short duration to simulate elapsed time + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + // Verify the total running time is approximately correct + long totalTime = stateMachine.getTotalRunningTimeMs(); + assertTrue(totalTime >= 100 && totalTime < 200); + } + + /** + * Test getCurrentStateRunningTimeMs method. + */ + @Test + void testGetCurrentStateRunningTimeMs() { + // Perform a state transition + stateMachine.transitTo(MessageStoreState.LOAD_COMMITLOG_OK); + + // Sleep for a short duration to simulate elapsed time + try { + Thread.sleep(100); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + + // Verify the current state running time is approximately correct + long currentStateTime = stateMachine.getCurrentStateRunningTimeMs(); + assertTrue(currentStateTime >= 100 && currentStateTime < 200); + } +} From dafcd584cf930f08f07e77583e04e7b5ef2a1ed6 Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 16:47:34 +0800 Subject: [PATCH 02/10] fix --- .../main/java/org/apache/rocketmq/store/ConsumeQueue.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index 8a3fade7bde..ae1ffff2767 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -792,7 +792,13 @@ private boolean putMessagePositionInfo(final long offset, final int size, final final long cqOffset) { if (offset + size <= this.getMaxPhysicOffset()) { - if (messageStore.getStateMachine().getCurrentState().isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK)) { + MessageStoreStateMachine stateMachine = messageStore.getStateMachine(); + MessageStoreStateMachine.MessageStoreState messageStoreState = stateMachine.getCurrentState(); + + // During the recovery process after broker crashes, this logs will cause the scrolling of valid logs. + // So only print the warning log after RECOVER_COMMITLOG_OK or the current state running time < 3 seconds. + if (messageStoreState.isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK) || + stateMachine.getCurrentStateRunningTimeMs() < 3000) { log.warn("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", this.getMaxPhysicOffset(), offset); } else { From e1d42cecce474dd5c843371e5c869c358e503430 Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 16:48:01 +0800 Subject: [PATCH 03/10] fix --- .../src/main/java/org/apache/rocketmq/store/ConsumeQueue.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index ae1ffff2767..134846872c8 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -21,6 +21,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.BoundaryType; import org.apache.rocketmq.common.MixAll; @@ -798,7 +799,7 @@ private boolean putMessagePositionInfo(final long offset, final int size, final // During the recovery process after broker crashes, this logs will cause the scrolling of valid logs. // So only print the warning log after RECOVER_COMMITLOG_OK or the current state running time < 3 seconds. if (messageStoreState.isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK) || - stateMachine.getCurrentStateRunningTimeMs() < 3000) { + stateMachine.getCurrentStateRunningTimeMs() < TimeUnit.SECONDS.toMillis(3)) { log.warn("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", this.getMaxPhysicOffset(), offset); } else { From cce736d6864bfbe067ad8bbada81581a8ff0658d Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 17:49:12 +0800 Subject: [PATCH 04/10] fix --- .../rocketmq/store/DefaultMessageStore.java | 10 ++++++---- .../store/MessageStoreStateMachine.java | 19 +++++++++++++------ .../store/DefaultMessageStoreTest.java | 16 ++++++++-------- .../store/MessageStoreStateMachineTest.java | 16 +++++++++++++++- 4 files changed, 42 insertions(+), 19 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index f8c3949a472..5926d1f3f27 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -308,20 +308,20 @@ public boolean load() { // load Commit Log result = this.commitLog.load(); - stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMMITLOG_OK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMMITLOG_OK, result); // load Consume Queue result = result && this.consumeQueueStore.load(); - stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_CONSUME_QUEUE_OK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_CONSUME_QUEUE_OK, result); if (messageStoreConfig.isEnableCompaction()) { result = result && this.compactionService.load(lastExitOK); - stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMPACTION_OK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_COMPACTION_OK, result); } if (result) { loadCheckPoint(); result = this.indexService.load(lastExitOK); - stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_INDEX_OK); + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_INDEX_OK, result); this.recover(lastExitOK); LOGGER.info("message store recover end, and the max phy offset = {}", this.getMaxPhyOffset()); } @@ -483,6 +483,7 @@ private void doRecheckReputOffsetFromCq() throws InterruptedException { public void shutdown() { if (!this.shutdown) { this.shutdown = true; + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.SHUTDOWN_BEGIN); this.scheduledExecutorService.shutdown(); this.scheduledCleanQueueExecutorService.shutdown(); @@ -514,6 +515,7 @@ public void shutdown() { if (this.runningFlags.isWriteable() && dispatchBehindBytes() == 0) { this.deleteFile(StorePathConfigHelper.getAbortFile(this.messageStoreConfig.getStorePathRootDir())); shutDownNormal = true; + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.SHUTDOWN_OK); } else { LOGGER.warn("the store may be wrong, so shutdown abnormally, and keep abort file."); } diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java index bd1f0eae7e1..6cd09781a43 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -41,7 +41,9 @@ public enum MessageStoreState { RECOVER_TOPIC_QUEUE_TABLE_OK(22), RUNNING(30), - SHUTDOWN(40), + + SHUTDOWN_BEGIN(40), + SHUTDOWN_OK(41), ERROR(Integer.MAX_VALUE); @@ -74,21 +76,26 @@ public MessageStoreStateMachine(Logger log) { } public void transitTo(MessageStoreState newState) { + transitTo(newState, true); + } + + public void transitTo(MessageStoreState newState, boolean success) { if (!newState.isAfter(currentState)) { throw new IllegalStateException( String.format("Invalid state transition from %s to %s. Can only move forward.", currentState, newState) ); } - - logStateChange(currentState, newState); - this.currentState = newState; - this.lastStateChangeTimestamp = System.currentTimeMillis(); + if (success) { + logStateChange(currentState, newState); + this.currentState = newState; + this.lastStateChangeTimestamp = System.currentTimeMillis(); + } } private void logStateChange(MessageStoreState fromState, MessageStoreState toState) { if (fromState == null) { - log.info("MessageStore initialized, state={}", toState); + log.info("MessageStore state initialized, state={}", toState); } else { log.info("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} " + "ms", fromState, toState, getCurrentStateRunningTimeMs(), getTotalRunningTimeMs()); diff --git a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java index eee38e0a8f4..ac25ac5430b 100644 --- a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java @@ -17,6 +17,10 @@ package org.apache.rocketmq.store; +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + import com.google.common.collect.Sets; import java.io.File; import java.io.RandomAccessFile; @@ -35,21 +39,21 @@ import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Properties; import java.util.Random; import java.util.UUID; -import java.util.Properties; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import org.apache.rocketmq.common.BrokerConfig; +import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.UtilAll; import org.apache.rocketmq.common.message.MessageBatch; import org.apache.rocketmq.common.message.MessageDecoder; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.message.MessageExtBatch; import org.apache.rocketmq.common.message.MessageExtBrokerInner; -import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.store.config.BrokerRole; import org.apache.rocketmq.store.config.FlushDiskType; import org.apache.rocketmq.store.config.MessageStoreConfig; @@ -66,10 +70,6 @@ import org.junit.runner.RunWith; import org.mockito.junit.MockitoJUnitRunner; -import static org.assertj.core.api.Assertions.assertThat; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; - @RunWith(MockitoJUnitRunner.class) public class DefaultMessageStoreTest { private final String storeMessage = "Once, there was a chance for me!"; @@ -911,7 +911,7 @@ public void testDeleteTopics() { String topicName = "topic-" + i; for (int j = 0; j < 4; j++) { ConsumeQueue consumeQueue = new ConsumeQueue(topicName, j, messageStoreConfig.getStorePathRootDir(), - messageStoreConfig.getMappedFileSizeConsumeQueue(), messageStore); + messageStoreConfig.getMappedFileSizeConsumeQueue(), (DefaultMessageStore) messageStore); cqTable.put(j, consumeQueue); } consumeQueueTable.put(topicName, cqTable); @@ -933,7 +933,7 @@ public void testCleanUnusedTopic() { String topicName = "topic-" + i; for (int j = 0; j < 4; j++) { ConsumeQueue consumeQueue = new ConsumeQueue(topicName, j, messageStoreConfig.getStorePathRootDir(), - messageStoreConfig.getMappedFileSizeConsumeQueue(), messageStore); + messageStoreConfig.getMappedFileSizeConsumeQueue(), (DefaultMessageStore) messageStore); cqTable.put(j, consumeQueue); } consumeQueueTable.put(topicName, cqTable); diff --git a/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java index 500cbc482fb..4675a2b6603 100644 --- a/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java @@ -20,6 +20,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.anyLong; import static org.mockito.Mockito.verify; @@ -53,7 +54,7 @@ void testConstructor() { assertEquals(MessageStoreState.INIT, stateMachine.getCurrentState()); // Verify logger was called for initialization - verify(mockLogger).info("MessageStore initialized, state={}", MessageStoreState.INIT); + verify(mockLogger).info(anyString(), eq(MessageStoreState.INIT)); } /** @@ -74,6 +75,19 @@ void testValidStateTransition() { ); } + /** + * Test fail state transition in transitTo method. + */ + @Test + void testValidFailStateTransition() { + stateMachine.transitTo(MessageStoreState.LOAD_COMMITLOG_OK, false); + assertEquals(MessageStoreState.INIT, stateMachine.getCurrentState()); + verify(mockLogger, Mockito.times(0)).info( + eq("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} ms"), + eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), anyLong(), anyLong() + ); + } + /** * Test invalid state transition in transitTo method. */ From 8f78ec9de9c844e6ff8cabc66d531ded5176dac5 Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 19:07:44 +0800 Subject: [PATCH 05/10] fix --- .../java/org/apache/rocketmq/store/ConsumeQueue.java | 12 ++---------- .../rocketmq/store/config/MessageStoreConfig.java | 10 ++++++++++ 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index 134846872c8..e33c331e0a2 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -21,7 +21,6 @@ import java.util.Collections; import java.util.List; import java.util.Map; -import java.util.concurrent.TimeUnit; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.BoundaryType; import org.apache.rocketmq.common.MixAll; @@ -793,18 +792,11 @@ private boolean putMessagePositionInfo(final long offset, final int size, final final long cqOffset) { if (offset + size <= this.getMaxPhysicOffset()) { - MessageStoreStateMachine stateMachine = messageStore.getStateMachine(); - MessageStoreStateMachine.MessageStoreState messageStoreState = stateMachine.getCurrentState(); - // During the recovery process after broker crashes, this logs will cause the scrolling of valid logs. - // So only print the warning log after RECOVER_COMMITLOG_OK or the current state running time < 3 seconds. - if (messageStoreState.isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK) || - stateMachine.getCurrentStateRunningTimeMs() < TimeUnit.SECONDS.toMillis(3)) { + if (messageStore.getStateMachine().getCurrentState().isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK) || + messageStore.getMessageStoreConfig().isEnableLogConsumeQueueRepeatedlyBuildWhenRecover()) { log.warn("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", this.getMaxPhysicOffset(), offset); - } else { - log.debug("Maybe try to build consume queue repeatedly maxPhysicOffset={} phyOffset={}", - this.getMaxPhysicOffset(), offset); } return true; } 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 28ab74eb353..60f6a90381c 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 @@ -483,6 +483,8 @@ public void setRocksdbCompressionType(String compressionType) { **/ private boolean useABSLock = false; + private boolean enableLogConsumeQueueRepeatedlyBuildWhenRecover = false; + public boolean isRocksdbCQDoubleWriteEnable() { return rocksdbCQDoubleWriteEnable; } @@ -2001,4 +2003,12 @@ public int getCombineCQMaxExtraSearchCommitLogFiles() { public void setCombineCQMaxExtraSearchCommitLogFiles(int combineCQMaxExtraSearchCommitLogFiles) { this.combineCQMaxExtraSearchCommitLogFiles = combineCQMaxExtraSearchCommitLogFiles; } + + public boolean isEnableLogConsumeQueueRepeatedlyBuildWhenRecover() { + return enableLogConsumeQueueRepeatedlyBuildWhenRecover; + } + + public void setEnableLogConsumeQueueRepeatedlyBuildWhenRecover(boolean enableLogConsumeQueueRepeatedlyBuildWhenRecover) { + this.enableLogConsumeQueueRepeatedlyBuildWhenRecover = enableLogConsumeQueueRepeatedlyBuildWhenRecover; + } } From 58dc9ef33403aaa96b206909c0ec23c151189e2e Mon Sep 17 00:00:00 2001 From: qianye Date: Wed, 9 Jul 2025 19:28:04 +0800 Subject: [PATCH 06/10] fix --- .../org/apache/rocketmq/store/MessageStoreStateMachine.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java index 6cd09781a43..d3a4071f883 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -80,7 +80,7 @@ public void transitTo(MessageStoreState newState) { } public void transitTo(MessageStoreState newState, boolean success) { - if (!newState.isAfter(currentState)) { + if (!MessageStoreState.ERROR.equals(currentState) && !newState.isAfter(currentState)) { throw new IllegalStateException( String.format("Invalid state transition from %s to %s. Can only move forward.", currentState, newState) From 3c7bb5969a8d0cf0be83f36434bac0b3825d029f Mon Sep 17 00:00:00 2001 From: qianye Date: Thu, 10 Jul 2025 10:15:35 +0800 Subject: [PATCH 07/10] fix --- .../store/MessageStoreStateMachine.java | 18 +++++++++++------- .../store/MessageStoreStateMachineTest.java | 12 ++++-------- 2 files changed, 15 insertions(+), 15 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java index d3a4071f883..312c6156ef4 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -72,7 +72,7 @@ public MessageStoreStateMachine(Logger log) { this.currentState = MessageStoreState.INIT; this.startTimestamp = System.currentTimeMillis(); this.lastStateChangeTimestamp = startTimestamp; - logStateChange(null, currentState); + logStateChange(null, currentState, true); } public void transitTo(MessageStoreState newState) { @@ -86,19 +86,23 @@ public void transitTo(MessageStoreState newState, boolean success) { currentState, newState) ); } + + logStateChange(currentState, newState, success); if (success) { - logStateChange(currentState, newState); this.currentState = newState; this.lastStateChangeTimestamp = System.currentTimeMillis(); } } - private void logStateChange(MessageStoreState fromState, MessageStoreState toState) { - if (fromState == null) { - log.info("MessageStore state initialized, state={}", toState); + private void logStateChange(MessageStoreState fromState, MessageStoreState toState, boolean success) { + if (fromState == null && success) { + log.info("MessageStoreState initialized, state={}", toState); + } else if (success) { + log.info("MessageStoreState transition from {} to {}; Time in previous state={}ms, Total time={}ms", + fromState, toState, getCurrentStateRunningTimeMs(), getTotalRunningTimeMs()); } else { - log.info("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} " - + "ms", fromState, toState, getCurrentStateRunningTimeMs(), getTotalRunningTimeMs()); + log.warn("MessageStoreState transition from {} to {} failed; Time in previous state={}ms, Total " + + "time={}ms", fromState, toState, getCurrentStateRunningTimeMs(), getTotalRunningTimeMs()); } } diff --git a/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java index 4675a2b6603..b6f424147d2 100644 --- a/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/MessageStoreStateMachineTest.java @@ -69,10 +69,8 @@ void testValidStateTransition() { assertEquals(MessageStoreState.LOAD_COMMITLOG_OK, stateMachine.getCurrentState()); // Verify logger was called for state transition - verify(mockLogger).info( - eq("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} ms"), - eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), anyLong(), anyLong() - ); + verify(mockLogger).info(anyString(), eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), + anyLong(), anyLong()); } /** @@ -82,10 +80,8 @@ void testValidStateTransition() { void testValidFailStateTransition() { stateMachine.transitTo(MessageStoreState.LOAD_COMMITLOG_OK, false); assertEquals(MessageStoreState.INIT, stateMachine.getCurrentState()); - verify(mockLogger, Mockito.times(0)).info( - eq("MessageStore state transition from {} to {}; Time in previous state: {} ms, Total time: {} ms"), - eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), anyLong(), anyLong() - ); + verify(mockLogger).warn(anyString(), eq(MessageStoreState.INIT), eq(MessageStoreState.LOAD_COMMITLOG_OK), + anyLong(), anyLong()); } /** From 3b5f07e5c83ced35723d9d6345835919d9148458 Mon Sep 17 00:00:00 2001 From: qianye Date: Thu, 10 Jul 2025 10:26:58 +0800 Subject: [PATCH 08/10] fix --- .../rocketmq/store/DefaultMessageStore.java | 3 ++- .../store/MessageStoreStateMachine.java | 18 ++++++++++-------- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index 5926d1f3f27..d37dbb81070 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -300,7 +300,7 @@ public boolean parseDelayLevel() { @Override public boolean load() { boolean result = true; - + stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.LOAD_BEGIN); try { boolean lastExitOK = !this.isTempFileExist(); LOGGER.info("last shutdown {}, store path root dir: {}", @@ -351,6 +351,7 @@ public void loadCheckPoint() throws IOException { } private void recover(final boolean lastExitOK) throws RocksDBException { + this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_BEGIN); // recover consume queue long recoverConsumeQueueStart = System.currentTimeMillis(); this.consumeQueueStore.recover(this.brokerConfig.isRecoverConcurrently()); diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java index 312c6156ef4..46d3839d298 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -31,14 +31,16 @@ public class MessageStoreStateMachine { public enum MessageStoreState { INIT(0), - LOAD_COMMITLOG_OK(10), - LOAD_CONSUME_QUEUE_OK(11), - LOAD_COMPACTION_OK(12), - LOAD_INDEX_OK(13), - - RECOVER_CONSUME_QUEUE_OK(20), - RECOVER_COMMITLOG_OK(21), - RECOVER_TOPIC_QUEUE_TABLE_OK(22), + LOAD_BEGIN(10), + LOAD_COMMITLOG_OK(11), + LOAD_CONSUME_QUEUE_OK(12), + LOAD_COMPACTION_OK(13), + LOAD_INDEX_OK(14), + + RECOVER_BEGIN(20), + RECOVER_CONSUME_QUEUE_OK(21), + RECOVER_COMMITLOG_OK(22), + RECOVER_TOPIC_QUEUE_TABLE_OK(23), RUNNING(30), From 2080b51592347a36d99a6c48b23e11b316582961 Mon Sep 17 00:00:00 2001 From: qianye Date: Thu, 10 Jul 2025 14:07:51 +0800 Subject: [PATCH 09/10] fix --- .../org/apache/rocketmq/store/DefaultMessageStore.java | 10 ---------- .../rocketmq/store/MessageStoreStateMachine.java | 6 ++---- 2 files changed, 2 insertions(+), 14 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index d37dbb81070..99eaa4b43c9 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -336,7 +336,6 @@ public boolean load() { if (!result) { this.allocateMappedFileService.shutdown(); - stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.ERROR); } return result; @@ -353,9 +352,7 @@ public void loadCheckPoint() throws IOException { private void recover(final boolean lastExitOK) throws RocksDBException { this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_BEGIN); // recover consume queue - long recoverConsumeQueueStart = System.currentTimeMillis(); this.consumeQueueStore.recover(this.brokerConfig.isRecoverConcurrently()); - long recoverConsumeQueueEnd = System.currentTimeMillis(); this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_CONSUME_QUEUE_OK); // recover commitlog @@ -368,15 +365,8 @@ private void recover(final boolean lastExitOK) throws RocksDBException { this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK); // recover consume offset table - long recoverCommitLogEnd = System.currentTimeMillis(); this.recoverTopicQueueTable(); - long recoverConsumeOffsetEnd = System.currentTimeMillis(); this.stateMachine.transitTo(MessageStoreStateMachine.MessageStoreState.RECOVER_TOPIC_QUEUE_TABLE_OK); - - LOGGER.info("message store recover total cost: {} ms, " + - "recoverConsumeQueue: {} ms, recoverCommitLog: {} ms, recoverOffsetTable: {} ms", - recoverConsumeOffsetEnd - recoverConsumeQueueStart, recoverConsumeQueueEnd - recoverConsumeQueueStart, - recoverCommitLogEnd - recoverConsumeQueueEnd, recoverConsumeOffsetEnd - recoverCommitLogEnd); } /** diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java index 46d3839d298..e5e07676dc7 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStoreStateMachine.java @@ -45,9 +45,7 @@ public enum MessageStoreState { RUNNING(30), SHUTDOWN_BEGIN(40), - SHUTDOWN_OK(41), - - ERROR(Integer.MAX_VALUE); + SHUTDOWN_OK(41); final int order; @@ -82,7 +80,7 @@ public void transitTo(MessageStoreState newState) { } public void transitTo(MessageStoreState newState, boolean success) { - if (!MessageStoreState.ERROR.equals(currentState) && !newState.isAfter(currentState)) { + if (!newState.isAfter(currentState)) { throw new IllegalStateException( String.format("Invalid state transition from %s to %s. Can only move forward.", currentState, newState) From feb674145dd21e347e73d294493e5951aaacad30 Mon Sep 17 00:00:00 2001 From: qianye Date: Thu, 10 Jul 2025 15:19:00 +0800 Subject: [PATCH 10/10] fix --- .../main/java/org/apache/rocketmq/store/ConsumeQueue.java | 6 +++--- .../main/java/org/apache/rocketmq/store/MessageStore.java | 2 ++ .../rocketmq/store/plugin/AbstractPluginMessageStore.java | 6 ++++++ 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index e33c331e0a2..2850299b7d6 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -61,7 +61,7 @@ public class ConsumeQueue implements ConsumeQueueInterface, FileQueueLifeCycle { public static final int MSG_TAG_OFFSET_INDEX = 12; private static final Logger LOG_ERROR = LoggerFactory.getLogger(LoggerName.STORE_ERROR_LOGGER_NAME); - private final DefaultMessageStore messageStore; + private final MessageStore messageStore; private final ConsumeQueueStore consumeQueueStore; private final MappedFileQueue mappedFileQueue; @@ -80,12 +80,12 @@ public class ConsumeQueue implements ConsumeQueueInterface, FileQueueLifeCycle { private ConsumeQueueExt consumeQueueExt = null; public ConsumeQueue(final String topic, final int queueId, final String storePath, final int mappedFileSize, - final DefaultMessageStore messageStore) { + final MessageStore messageStore) { this(topic, queueId, storePath, mappedFileSize, messageStore, (ConsumeQueueStore) messageStore.getQueueStore()); } public ConsumeQueue(final String topic, final int queueId, final String storePath, final int mappedFileSize, - final DefaultMessageStore messageStore, final ConsumeQueueStore consumeQueueStore) { + final MessageStore messageStore, final ConsumeQueueStore consumeQueueStore) { this.storePath = storePath; this.mappedFileSize = mappedFileSize; this.messageStore = messageStore; diff --git a/store/src/main/java/org/apache/rocketmq/store/MessageStore.java b/store/src/main/java/org/apache/rocketmq/store/MessageStore.java index 9c9a556f6d3..52c2de33fd3 100644 --- a/store/src/main/java/org/apache/rocketmq/store/MessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/MessageStore.java @@ -983,4 +983,6 @@ DispatchRequest checkMessageAndReturnSize(final ByteBuffer byteBuffer, final boo * notify message arrive if necessary */ void notifyMessageArriveIfNecessary(DispatchRequest dispatchRequest); + + MessageStoreStateMachine getStateMachine(); } diff --git a/store/src/main/java/org/apache/rocketmq/store/plugin/AbstractPluginMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/plugin/AbstractPluginMessageStore.java index 74c67d01621..19ace1c8e3a 100644 --- a/store/src/main/java/org/apache/rocketmq/store/plugin/AbstractPluginMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/plugin/AbstractPluginMessageStore.java @@ -42,6 +42,7 @@ import org.apache.rocketmq.store.GetMessageResult; import org.apache.rocketmq.store.MessageFilter; import org.apache.rocketmq.store.MessageStore; +import org.apache.rocketmq.store.MessageStoreStateMachine; import org.apache.rocketmq.store.PutMessageResult; import org.apache.rocketmq.store.QueryMessageResult; import org.apache.rocketmq.store.RunningFlags; @@ -660,4 +661,9 @@ public void notifyMessageArriveIfNecessary(DispatchRequest dispatchRequest) { public MessageStore getNext() { return next; } + + @Override + public MessageStoreStateMachine getStateMachine() { + return next.getStateMachine(); + } }