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 8d0c91595 fix(topic): honor explicit queue counts on update (#2354)
8d0c91595 is described below

commit 8d0c9159556d39b2a224e1d752e8352bef3c8cfe
Author: yyqdbngt <[email protected]>
AuthorDate: Fri Aug 21 16:40:55 2026 +0800

    fix(topic): honor explicit queue counts on update (#2354)
---
 .../provider/apache/RocketMQAdminClientImpl.java   | 16 ++++++-------
 .../apache/RocketMQAdminClientImplTest.java        | 28 ++++++++++++++++++++++
 2 files changed, 36 insertions(+), 8 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
index a333897f8..3190569f9 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
@@ -317,14 +317,14 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
                 // Preserve the existing queue counts when the update request 
does not change them,
                 // matching the perm semantics below; defaulting to 8 would 
silently resize the
                 // topic on partial updates (e.g. perm or remark only).
-                int writeQueues = existing != null && 
existing.getWriteQueueNums() != null
-                                && existing.getWriteQueueNums() > 0
-                        ? existing.getWriteQueueNums()
-                        : topic.getWriteQueues() > 0 ? topic.getWriteQueues() 
: 8;
-                int readQueues = existing != null && 
existing.getReadQueueNums() != null
-                                && existing.getReadQueueNums() > 0
-                        ? existing.getReadQueueNums()
-                        : topic.getReadQueues() > 0 ? topic.getReadQueues() : 
8;
+                int writeQueues = topic.getWriteQueues() > 0
+                        ? topic.getWriteQueues()
+                        : existing != null && existing.getWriteQueueNums() != 
null
+                                && existing.getWriteQueueNums() > 0 ? 
existing.getWriteQueueNums() : 8;
+                int readQueues = topic.getReadQueues() > 0
+                        ? topic.getReadQueues()
+                        : existing != null && existing.getReadQueueNums() != 
null
+                                && existing.getReadQueueNums() > 0 ? 
existing.getReadQueueNums() : 8;
                 TopicPerm effectivePerm = topic.getPerm() != null
                         ? topic.getPerm()
                         : existing == null ? TopicPerm.RW : 
fromRocketMQPerm(existing.getPerm());
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
index 04cbdc210..b19ee8a1e 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
@@ -479,6 +479,34 @@ class RocketMQAdminClientImplTest {
         assertThat(existing.getReadQueueNums()).isEqualTo(16);
     }
 
+    @Test
+    void updateTopicAppliesExplicitQueueCounts() throws Exception {
+        TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new 
MybatisConfiguration(), ""), RmqTopic.class);
+        RmqTopic existing = new RmqTopic();
+        existing.setWriteQueueNums(8);
+        existing.setReadQueueNums(8);
+        existing.setPerm(6);
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithMaster());
+        when(topicMapper.selectOne(any())).thenReturn(existing);
+        doNothing().when(adminExt).createAndUpdateTopicConfig(anyString(), 
any(TopicConfig.class));
+
+        TopicVO topic = new TopicVO();
+        topic.setName("orders");
+        topic.setWriteQueues(16);
+        topic.setReadQueues(12);
+
+        TopicVO updated = adminClient.updateTopic(topic);
+
+        ArgumentCaptor<TopicConfig> topicConfigCaptor = 
ArgumentCaptor.forClass(TopicConfig.class);
+        verify(adminExt).createAndUpdateTopicConfig(anyString(), 
topicConfigCaptor.capture());
+        
assertThat(topicConfigCaptor.getValue().getWriteQueueNums()).isEqualTo(16);
+        
assertThat(topicConfigCaptor.getValue().getReadQueueNums()).isEqualTo(12);
+        assertThat(existing.getWriteQueueNums()).isEqualTo(16);
+        assertThat(existing.getReadQueueNums()).isEqualTo(12);
+        assertThat(updated.getWriteQueues()).isEqualTo(16);
+        assertThat(updated.getReadQueues()).isEqualTo(12);
+    }
+
     @Test
     void 
topicDeleteUsesSelectedInstanceAndScopesMetadataToClusterAndInstance() throws 
Exception {
         TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new 
MybatisConfiguration(), ""), RmqTopic.class);

Reply via email to