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