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);