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 1c8542efa fix(group): preflight settings across brokers (#4542)
1c8542efa is described below

commit 1c8542efa7d86c0fa12b62b0b4e718372a660ed5
Author: zmuxuny <[email protected]>
AuthorDate: Mon Sep 21 17:30:32 2026 +0800

    fix(group): preflight settings across brokers (#4542)
    
    updateConsumerGroupSettings read, mutated and wrote each master Broker in a 
single pass, so a later missing or unavailable Broker config failed the request 
after earlier Brokers had already been changed. The update is now two-phase: 
every master's SubscriptionGroupConfig is read first (any failure aborts with 
zero writes), then the mutations are applied and written back.
    
    Fixes #4541
---
 .../provider/apache/RocketMQAdminClientImpl.java   |  8 ++++-
 .../apache/RocketMQAdminClientImplTest.java        | 36 ++++++++++++++++++++++
 2 files changed, 43 insertions(+), 1 deletion(-)

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 6257eb2ec..1a0cfd522 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
@@ -659,12 +659,18 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
                     throw new BusinessException(502, "No broker available to 
update consumer group settings");
                 }
                 totalBrokers = brokerAddrs.size();
-                SubscriptionGroupConfig applied = null;
+                Map<String, SubscriptionGroupConfig> configsByBroker = new 
LinkedHashMap<>();
                 for (String brokerAddr : brokerAddrs) {
                     SubscriptionGroupConfig config = 
admin.examineSubscriptionGroupConfig(brokerAddr, name);
                     if (config == null) {
                         throw new BusinessException(404, "Consumer group not 
found: " + name);
                     }
+                    configsByBroker.put(brokerAddr, config);
+                }
+                SubscriptionGroupConfig applied = null;
+                for (Map.Entry<String, SubscriptionGroupConfig> entry : 
configsByBroker.entrySet()) {
+                    String brokerAddr = entry.getKey();
+                    SubscriptionGroupConfig config = entry.getValue();
                     config.setRetryQueueNums(command.retryQueueNums());
                     config.setRetryMaxTimes(command.retryMaxTimes());
                     if (command.consumeEnable() != null) {
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 91a5735ad..d644b47c6 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
@@ -1153,6 +1153,42 @@ class RocketMQAdminClientImplTest {
         verifyNoInteractions(groupMapper);
     }
 
+    @Test
+    void updateConsumerGroupSettingsShouldReadAllBrokersBeforeWritingTest() 
throws Exception {
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithTwoMasters());
+        SubscriptionGroupConfig config = new SubscriptionGroupConfig();
+        config.setGroupName("cg-orders");
+        when(adminExt.examineSubscriptionGroupConfig(anyString(), 
eq("cg-orders")))
+                .thenReturn(config)
+                .thenThrow(new IllegalStateException("broker unavailable"));
+
+        assertThatThrownBy(() -> adminClient.updateConsumerGroupSettings(null, 
"cg-orders",
+                new ConsumerGroupSettingsCommand(2, 8, null, null, null)))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("broker unavailable");
+
+        verify(adminExt, 
never()).createAndUpdateSubscriptionGroupConfig(anyString(), any());
+        verifyNoInteractions(groupMapper);
+    }
+
+    @Test
+    void 
updateConsumerGroupSettingsShouldNotWriteWhenLaterBrokerConfigIsMissingTest() 
throws Exception {
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithTwoMasters());
+        SubscriptionGroupConfig config = new SubscriptionGroupConfig();
+        config.setGroupName("cg-orders");
+        when(adminExt.examineSubscriptionGroupConfig(anyString(), 
eq("cg-orders")))
+                .thenReturn(config)
+                .thenReturn(null);
+
+        assertThatThrownBy(() -> adminClient.updateConsumerGroupSettings(null, 
"cg-orders",
+                new ConsumerGroupSettingsCommand(2, 8, null, null, null)))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("Consumer group not found");
+
+        verify(adminExt, 
never()).createAndUpdateSubscriptionGroupConfig(anyString(), any());
+        verifyNoInteractions(groupMapper);
+    }
+
     @Test
     void updateConsumerGroupSettingsPreservesBrokerConfiguration() throws 
Exception {
         TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new 
MybatisConfiguration(), ""), RmqGroup.class);

Reply via email to