This is an automated email from the ASF dual-hosted git repository.
lwclover pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git
The following commit(s) were added to refs/heads/develop by this push:
new 706142718f [ISSUE #9916] Remove getTimerCheckpoint() from
BrokerController (#9905)
706142718f is described below
commit 706142718f2e49e8fe22239e1159633eeb231d81
Author: rongtong <[email protected]>
AuthorDate: Mon Sep 7 17:33:09 2026 +0800
[ISSUE #9916] Remove getTimerCheckpoint() from BrokerController (#9905)
* refactor: remove getTimerCheckpoint() from BrokerController
* Update SlaveSynchronizeTest.java
fix UnnecessaryStubbing
---------
Co-authored-by: RongtongJin <[email protected]>
Co-authored-by: sun <[email protected]>
---
.../main/java/org/apache/rocketmq/broker/BrokerController.java | 4 ----
.../org/apache/rocketmq/broker/BrokerPreOnlineService.java | 10 +++++-----
.../apache/rocketmq/broker/processor/AdminBrokerProcessor.java | 2 +-
.../org/apache/rocketmq/broker/slave/SlaveSynchronize.java | 8 ++++----
.../rocketmq/broker/processor/AdminBrokerProcessorTest.java | 5 +++--
.../org/apache/rocketmq/broker/slave/SlaveSynchronizeTest.java | 4 +++-
.../apache/rocketmq/test/container/GetMetadataReverseIT.java | 2 +-
7 files changed, 17 insertions(+), 18 deletions(-)
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java
b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java
index 91f281e5e1..a105c71376 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java
@@ -2766,10 +2766,6 @@ public class BrokerController {
return this.isIsolated;
}
- public TimerCheckpoint getTimerCheckpoint() {
- return timerCheckpoint;
- }
-
public TopicRouteInfoManager getTopicRouteInfoManager() {
return this.topicRouteInfoManager;
}
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java
b/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java
index de2ccb2939..1629838cc2 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/BrokerPreOnlineService.java
@@ -182,12 +182,12 @@ public class BrokerPreOnlineService extends ServiceThread
{
}
}
- if (null != this.brokerController.getTimerCheckpoint() &&
this.brokerController.getTimerCheckpoint().getDataVersion().compare(timerCheckpoint.getDataVersion())
<= 0) {
+ if (null !=
this.brokerController.getTimerMessageStore().getTimerCheckpoint() &&
this.brokerController.getTimerMessageStore().getTimerCheckpoint().getDataVersion().compare(timerCheckpoint.getDataVersion())
<= 0) {
LOGGER.info("{}'s timerCheckpoint data version is larger than
master broker, {}'s timerCheckpoint will be used.", brokerAddr, brokerAddr);
-
this.brokerController.getTimerCheckpoint().setLastReadTimeMs(timerCheckpoint.getLastReadTimeMs());
-
this.brokerController.getTimerCheckpoint().setMasterTimerQueueOffset(timerCheckpoint.getMasterTimerQueueOffset());
-
this.brokerController.getTimerCheckpoint().getDataVersion().assignNewOne(timerCheckpoint.getDataVersion());
- this.brokerController.getTimerCheckpoint().flush();
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().setLastReadTimeMs(timerCheckpoint.getLastReadTimeMs());
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().setMasterTimerQueueOffset(timerCheckpoint.getMasterTimerQueueOffset());
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().getDataVersion().assignNewOne(timerCheckpoint.getDataVersion());
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().flush();
}
for (BrokerAttachedPlugin brokerAttachedPlugin :
brokerController.getBrokerAttachedPlugins()) {
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
index de90396064..602a8efce0 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
@@ -971,7 +971,7 @@ public class AdminBrokerProcessor implements
NettyRequestProcessor {
private RemotingCommand getTimerCheckPoint(ChannelHandlerContext ctx,
RemotingCommand request) {
final RemotingCommand response =
RemotingCommand.createResponseCommand(ResponseCode.SYSTEM_ERROR, "Unknown");
- TimerCheckpoint timerCheckpoint =
this.brokerController.getTimerCheckpoint();
+ TimerCheckpoint timerCheckpoint =
this.brokerController.getTimerMessageStore().getTimerCheckpoint();
if (null == timerCheckpoint) {
LOGGER.error("AdminBrokerProcessor#getTimerCheckPoint: checkpoint
is null, caller={}", ctx.channel().remoteAddress());
response.setCode(ResponseCode.SYSTEM_ERROR);
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/slave/SlaveSynchronize.java
b/broker/src/main/java/org/apache/rocketmq/broker/slave/SlaveSynchronize.java
index 78f78216b5..355e20953f 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/slave/SlaveSynchronize.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/slave/SlaveSynchronize.java
@@ -236,10 +236,10 @@ public class SlaveSynchronize {
if (null !=
brokerController.getMessageStore().getTimerMessageStore() &&
!brokerController.getTimerMessageStore().isShouldRunningDequeue()) {
TimerCheckpoint checkpoint =
this.brokerController.getBrokerOuterAPI().getTimerCheckPoint(masterAddrBak);
- if (null != this.brokerController.getTimerCheckpoint()) {
-
this.brokerController.getTimerCheckpoint().setLastReadTimeMs(checkpoint.getLastReadTimeMs());
-
this.brokerController.getTimerCheckpoint().setMasterTimerQueueOffset(checkpoint.getMasterTimerQueueOffset());
-
this.brokerController.getTimerCheckpoint().getDataVersion().assignNewOne(checkpoint.getDataVersion());
+ if (null !=
this.brokerController.getTimerMessageStore().getTimerCheckpoint()) {
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().setLastReadTimeMs(checkpoint.getLastReadTimeMs());
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().setMasterTimerQueueOffset(checkpoint.getMasterTimerQueueOffset());
+
this.brokerController.getTimerMessageStore().getTimerCheckpoint().getDataVersion().assignNewOne(checkpoint.getDataVersion());
}
}
} catch (Exception e) {
diff --git
a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
index ddadfae410..006979ce86 100644
---
a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
+++
b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
@@ -1561,13 +1561,14 @@ public class AdminBrokerProcessorTest {
@Test
public void testGetTimeCheckPoint() throws RemotingCommandException {
- when(this.brokerController.getTimerCheckpoint()).thenReturn(null);
+
when(this.brokerController.getTimerMessageStore()).thenReturn(timerMessageStore);
+ when(this.timerMessageStore.getTimerCheckpoint()).thenReturn(null);
RemotingCommand request =
RemotingCommand.createRequestCommand(RequestCode.GET_TIMER_CHECK_POINT, null);
RemotingCommand response =
adminBrokerProcessor.processRequest(handlerContext, request);
assertThat(response.getCode()).isEqualTo(ResponseCode.SYSTEM_ERROR);
assertThat(response.getRemark()).isEqualTo("The checkpoint is null");
- when(this.brokerController.getTimerCheckpoint()).thenReturn(new
TimerCheckpoint());
+ when(this.timerMessageStore.getTimerCheckpoint()).thenReturn(new
TimerCheckpoint());
response = adminBrokerProcessor.processRequest(handlerContext,
request);
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
}
diff --git
a/broker/src/test/java/org/apache/rocketmq/broker/slave/SlaveSynchronizeTest.java
b/broker/src/test/java/org/apache/rocketmq/broker/slave/SlaveSynchronizeTest.java
index 192448b897..073f57a77a 100644
---
a/broker/src/test/java/org/apache/rocketmq/broker/slave/SlaveSynchronizeTest.java
+++
b/broker/src/test/java/org/apache/rocketmq/broker/slave/SlaveSynchronizeTest.java
@@ -114,7 +114,9 @@ public class SlaveSynchronizeTest {
when(brokerController.getQueryAssignmentProcessor()).thenReturn(queryAssignmentProcessor);
when(brokerController.getMessageStore()).thenReturn(messageStore);
when(brokerController.getTimerMessageStore()).thenReturn(timerMessageStore);
-
when(brokerController.getTimerCheckpoint()).thenReturn(timerCheckpoint);
+
when(timerMessageStore.getTimerCheckpoint()).thenReturn(timerCheckpoint);
+ // when(topicConfigManager.getDataVersion()).thenReturn(new
DataVersion());
+ // when(topicConfigManager.getTopicConfigTable()).thenReturn(new
ConcurrentHashMap<>());
when(brokerController.getConsumerOffsetManager()).thenReturn(consumerOffsetManager);
when(consumerOffsetManager.getOffsetTable()).thenReturn(new
ConcurrentHashMap<>());
when(consumerOffsetManager.getDataVersion()).thenReturn(new
DataVersion());
diff --git
a/test/src/test/java/org/apache/rocketmq/test/container/GetMetadataReverseIT.java
b/test/src/test/java/org/apache/rocketmq/test/container/GetMetadataReverseIT.java
index b9bb7b2e1e..4ebf9a85ac 100644
---
a/test/src/test/java/org/apache/rocketmq/test/container/GetMetadataReverseIT.java
+++
b/test/src/test/java/org/apache/rocketmq/test/container/GetMetadataReverseIT.java
@@ -282,7 +282,7 @@ public class GetMetadataReverseIT extends
ContainerIntegrationTestBase {
awaitUntilSlaveOK();
- await().atMost(Duration.ofMinutes(1)).until(() ->
master1With3Replicas.getTimerCheckpoint().getMasterTimerQueueOffset() >=
MESSAGE_COUNT);
+ await().atMost(Duration.ofMinutes(1)).until(() ->
master1With3Replicas.getTimerMessageStore().getTimerCheckpoint().getMasterTimerQueueOffset()
>= MESSAGE_COUNT);
pushConsumer.shutdown();
}