This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 248e3869 fix(cluster): keep write/read queue counts in sync on partial
updates (#1075)
248e3869 is described below
commit 248e3869565ea8a2611ad9d92a6df8e884155ef6
Author: yyqdbngt <[email protected]>
AuthorDate: Thu Aug 6 16:18:03 2026 +0800
fix(cluster): keep write/read queue counts in sync on partial updates
(#1075)
Co-authored-by: yyqdbngt <[email protected]>
---
.../rocketmq/studio/cluster/broker/ClusterService.java | 13 ++++++++-----
.../studio/cluster/broker/ClusterServiceTest.java | 16 ++++++++++++++++
2 files changed, 24 insertions(+), 5 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
index 16a8e155..a8099614 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
@@ -154,11 +154,14 @@ public class ClusterService {
if (command.getFileReservedTime() != null) {
config.setFileReservedTime(command.getFileReservedTime());
}
- if (command.getWriteQueueNums() != null) {
- config.setWriteQueueNums(command.getWriteQueueNums());
- }
- if (command.getReadQueueNums() != null) {
- config.setReadQueueNums(command.getReadQueueNums());
+ // The broker exposes a single defaultTopicQueueNums property, so a
partial update of only
+ // one of the write/read queues must mirror the value onto the other
to keep the stored
+ // config consistent with what the broker actually applies.
+ if (command.getWriteQueueNums() != null || command.getReadQueueNums()
!= null) {
+ int queueNums = command.getWriteQueueNums() != null
+ ? command.getWriteQueueNums() : command.getReadQueueNums();
+ config.setWriteQueueNums(queueNums);
+ config.setReadQueueNums(queueNums);
}
if (command.getBrokerPermission() != null) {
config.setBrokerPermission(command.getBrokerPermission());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
index f9ecfe66..6100d4a9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
@@ -558,6 +558,22 @@ class ClusterServiceTest {
.satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(404));
}
+ @Test
+ void updateClusterConfigShouldMirrorSingleQueueNumsUpdate() {
+
when(clusterProvider.refreshClusterDetail("cluster-1")).thenReturn(sampleCluster);
+ UpdateConfigDTO command = UpdateConfigDTO.builder()
+ .id("cluster-1")
+ .writeQueueNums(12)
+ .build();
+
+ ClusterConfigUpdateResultVO result =
clusterService.updateClusterConfig(command);
+
+ // The broker exposes a single defaultTopicQueueNums property, so the
stored write and read
+ // queue counts must stay in sync after a partial update.
+
assertThat(result.getCluster().getConfig().getWriteQueueNums()).isEqualTo(12);
+
assertThat(result.getCluster().getConfig().getReadQueueNums()).isEqualTo(12);
+ }
+
@Test
void updateClusterConfigShouldUseLiveClusterWhenItIsNotPersisted() {
when(clusterProvider.refreshClusterDetail("cluster-1")).thenReturn(sampleCluster);