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 a4b9d97c3 feat(consumer): manage consumption switches in group
settings (#3991)
a4b9d97c3 is described below
commit a4b9d97c3b2acedb293b1287a0546c78397bf894
Author: 烤化の初雪 <[email protected]>
AuthorDate: Mon Sep 7 18:57:23 2026 +0800
feat(consumer): manage consumption switches in group settings (#3991)
Extend the Apache consumer group settings editor beyond retryQueueNums
and retryMaxTimes with the SubscriptionGroupConfig consumption switches
(consumeEnable, consumeMessageOrderly, consumeBroadcastEnable). The
settings endpoints already read the effective broker config, update
every master broker, preserve unrelated fields, and audit the change;
omitted switches preserve the current broker values.
The classic dashboard edits these fields through the same broker API;
in Studio they required external tooling. Cloud instances stay
unsupported for this editor, matching the existing Apache-only scope.
Related to #3990.
Co-authored-by: unbridled-41
<[email protected]>
---
.../instance/group/ConsumerGroupController.java | 5 +-
...gsVO.java => ConsumerGroupSettingsCommand.java} | 24 +++----
.../instance/group/ConsumerGroupSettingsVO.java | 3 +
.../group/UpdateConsumerGroupSettingsDTO.java | 3 +
.../studio/instance/topic/MetadataService.java | 7 +-
.../studio/provider/apache/AdminClient.java | 5 +-
.../provider/apache/RocketMQAdminClientImpl.java | 59 ++++++++++++----
.../group/ConsumerGroupControllerTest.java | 31 ++++++++-
.../apache/RocketMQAdminClientImplTest.java | 80 +++++++++++++++++++++-
web/src/api/metadata.ts | 3 +
.../pages/instance/__tests__/ConsumerPage.test.tsx | 48 +++++++++++++
web/src/pages/instance/consumer.tsx | 62 ++++++++++++++++-
12 files changed, 294 insertions(+), 36 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
index 99d727c27..1eb357f07 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
@@ -77,8 +77,11 @@ public class ConsumerGroupController {
@PostMapping("/settings")
public Result<ConsumerGroupSettingsVO> updateConsumerGroupSettings(
@Valid @RequestBody UpdateConsumerGroupSettingsDTO request) {
+ ConsumerGroupSettingsCommand command = new
ConsumerGroupSettingsCommand(
+ request.getRetryQueueNums(), request.getRetryMaxTimes(),
request.getConsumeEnable(),
+ request.getConsumeMessageOrderly(),
request.getConsumeBroadcastEnable());
return
Result.ok(metadataService.updateConsumerGroupSettings(request.getInstanceId(),
request.getName(),
- request.getRetryQueueNums(), request.getRetryMaxTimes()));
+ command));
}
@GetMapping("/{name}/refresh")
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsCommand.java
similarity index 64%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsCommand.java
index 7cbff0e5e..e397c6054 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsCommand.java
@@ -16,17 +16,15 @@
*/
package org.apache.rocketmq.studio.instance.group;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
-
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class ConsumerGroupSettingsVO {
- private String groupName;
- private int retryQueueNums;
- private int retryMaxTimes;
+/**
+ * Carries the consumer-group settings update fields, replacing the long
positional parameter list
+ * on the update path. {@code retryQueueNums}/{@code retryMaxTimes} are always
applied; a null
+ * consumption switch means "preserve the current broker value".
+ */
+public record ConsumerGroupSettingsCommand(
+ int retryQueueNums,
+ int retryMaxTimes,
+ Boolean consumeEnable,
+ Boolean consumeMessageOrderly,
+ Boolean consumeBroadcastEnable) {
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
index 7cbff0e5e..954703ab9 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupSettingsVO.java
@@ -29,4 +29,7 @@ public class ConsumerGroupSettingsVO {
private String groupName;
private int retryQueueNums;
private int retryMaxTimes;
+ private boolean consumeEnable;
+ private boolean consumeMessageOrderly;
+ private boolean consumeBroadcastEnable;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/UpdateConsumerGroupSettingsDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/UpdateConsumerGroupSettingsDTO.java
index 0c3cf330f..56036204d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/UpdateConsumerGroupSettingsDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/UpdateConsumerGroupSettingsDTO.java
@@ -33,4 +33,7 @@ public class UpdateConsumerGroupSettingsDTO {
@NotNull(message = "retryMaxTimes is required")
@Positive(message = "retryMaxTimes must be positive")
private Integer retryMaxTimes;
+ private Boolean consumeEnable;
+ private Boolean consumeMessageOrderly;
+ private Boolean consumeBroadcastEnable;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index fd7176880..04d347c8b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -32,6 +32,7 @@ import
org.apache.rocketmq.studio.instance.group.CreateConsumerGroupDTO;
import org.apache.rocketmq.studio.instance.group.ImportConsumerGroupsResultVO;
import org.springframework.util.StringUtils;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsCommand;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetPreviewVO;
@@ -291,12 +292,12 @@ public class MetadataService {
return adminClient.getConsumerGroupSettings(instanceId,
requireName(name, "consumer group name"));
}
- public ConsumerGroupSettingsVO updateConsumerGroupSettings(String
instanceId, String name, int retryQueueNums,
- int
retryMaxTimes) {
+ public ConsumerGroupSettingsVO updateConsumerGroupSettings(String
instanceId, String name,
+
ConsumerGroupSettingsCommand command) {
instanceId = normalizeInstanceId(instanceId);
requireApacheInstance(instanceId);
String groupName = requireName(name, "consumer group name");
- return adminClient.updateConsumerGroupSettings(instanceId, groupName,
retryQueueNums, retryMaxTimes);
+ return adminClient.updateConsumerGroupSettings(instanceId, groupName,
command);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/AdminClient.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/AdminClient.java
index 1c4f91c29..747f38ad5 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/AdminClient.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/AdminClient.java
@@ -20,6 +20,7 @@ import
org.apache.rocketmq.studio.instance.topic.SendMessageVO;
import org.apache.rocketmq.studio.instance.topic.SendMessageDTO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsCommand;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsVO;
import org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetPreviewVO;
@@ -33,8 +34,8 @@ public interface AdminClient {
SendMessageVO sendMessage(SendMessageDTO request);
ConsumerGroupVO createConsumerGroup(ConsumerGroupVO group);
ConsumerGroupSettingsVO getConsumerGroupSettings(String instanceId, String
name);
- ConsumerGroupSettingsVO updateConsumerGroupSettings(String instanceId,
String name, int retryQueueNums,
- int retryMaxTimes);
+ ConsumerGroupSettingsVO updateConsumerGroupSettings(String instanceId,
String name,
+
ConsumerGroupSettingsCommand command);
void deleteConsumerGroup(String instanceId, String name);
ResetConsumerOffsetPreviewVO previewResetOffset(String instanceId, String
name, long timestamp, String topic);
void resetOffset(String instanceId, String name, long timestamp, String
topic);
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 43478ff22..1baa1ee3a 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
@@ -38,6 +38,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsCommand;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsVO;
import org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetPreviewVO;
import
org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetQueuePreviewVO;
@@ -520,7 +521,9 @@ public class RocketMQAdminClientImpl implements AdminClient
{
throw new BusinessException(404, "Consumer group not
found: " + name);
}
return
ConsumerGroupSettingsVO.builder().groupName(name).retryQueueNums(config.getRetryQueueNums())
- .retryMaxTimes(config.getRetryMaxTimes()).build();
+
.retryMaxTimes(config.getRetryMaxTimes()).consumeEnable(config.isConsumeEnable())
+
.consumeMessageOrderly(config.isConsumeMessageOrderly())
+
.consumeBroadcastEnable(config.isConsumeBroadcastEnable()).build();
} catch (BusinessException exception) {
throw exception;
} catch (Exception exception) {
@@ -530,47 +533,79 @@ public class RocketMQAdminClientImpl implements
AdminClient {
}
@Override
- public ConsumerGroupSettingsVO updateConsumerGroupSettings(String
instanceId, String name, int retryQueueNums,
- int
retryMaxTimes) {
+ public ConsumerGroupSettingsVO updateConsumerGroupSettings(String
instanceId, String name,
+
ConsumerGroupSettingsCommand command) {
return executeForInstance(instanceId, admin -> {
+ int totalBrokers = 0;
+ int updatedBrokers = 0;
try {
String clusterName = getClusterName(admin);
Set<String> brokerAddrs =
getMasterBrokerAddrsForCluster(admin, clusterName);
if (brokerAddrs.isEmpty()) {
throw new BusinessException(502, "No broker available to
update consumer group settings");
}
+ totalBrokers = brokerAddrs.size();
+ SubscriptionGroupConfig applied = null;
for (String brokerAddr : brokerAddrs) {
SubscriptionGroupConfig config =
admin.examineSubscriptionGroupConfig(brokerAddr, name);
if (config == null) {
throw new BusinessException(404, "Consumer group not
found: " + name);
}
- config.setRetryQueueNums(retryQueueNums);
- config.setRetryMaxTimes(retryMaxTimes);
+ config.setRetryQueueNums(command.retryQueueNums());
+ config.setRetryMaxTimes(command.retryMaxTimes());
+ if (command.consumeEnable() != null) {
+ config.setConsumeEnable(command.consumeEnable());
+ }
+ if (command.consumeMessageOrderly() != null) {
+
config.setConsumeMessageOrderly(command.consumeMessageOrderly());
+ }
+ if (command.consumeBroadcastEnable() != null) {
+
config.setConsumeBroadcastEnable(command.consumeBroadcastEnable());
+ }
admin.createAndUpdateSubscriptionGroupConfig(brokerAddr,
config);
+ updatedBrokers++;
+ if (applied == null) {
+ applied = config;
+ }
}
RmqGroup group = groupMapper.selectOne(new
LambdaQueryWrapper<RmqGroup>()
.eq(RmqGroup::getClusterId, clusterName)
.eq(RmqGroup::getInstanceId, metadataScope(instanceId))
.eq(RmqGroup::getName, name));
if (group != null) {
- group.setMaxRetry(retryMaxTimes);
+ group.setMaxRetry(command.retryMaxTimes());
group.setGmtModified(LocalDateTime.now());
groupMapper.updateById(group);
}
- recordAudit("UPDATE_GROUP_SETTINGS", name, "retryQueueNums=" +
retryQueueNums
- + ", retryMaxTimes=" + retryMaxTimes, "SUCCESS");
- return
ConsumerGroupSettingsVO.builder().groupName(name).retryQueueNums(retryQueueNums)
- .retryMaxTimes(retryMaxTimes).build();
+ recordAudit("UPDATE_GROUP_SETTINGS", name, "retryQueueNums=" +
command.retryQueueNums()
+ + ", retryMaxTimes=" + command.retryMaxTimes()
+ + ", consumeEnable=" +
describeSwitch(command.consumeEnable())
+ + ", consumeMessageOrderly=" +
describeSwitch(command.consumeMessageOrderly())
+ + ", consumeBroadcastEnable=" +
describeSwitch(command.consumeBroadcastEnable())
+ + ", brokersUpdated=" + updatedBrokers + "/" +
totalBrokers, "SUCCESS");
+ return ConsumerGroupSettingsVO.builder().groupName(name)
+ .retryQueueNums(applied.getRetryQueueNums())
+ .retryMaxTimes(applied.getRetryMaxTimes())
+ .consumeEnable(applied.isConsumeEnable())
+
.consumeMessageOrderly(applied.isConsumeMessageOrderly())
+
.consumeBroadcastEnable(applied.isConsumeBroadcastEnable())
+ .build();
} catch (BusinessException exception) {
- recordAudit("UPDATE_GROUP_SETTINGS", name,
exception.getMessage(), "FAILED");
+ recordAudit("UPDATE_GROUP_SETTINGS", name, "updated " +
updatedBrokers + "/" + totalBrokers
+ + " brokers before failure: " +
exception.getMessage(), "FAILED");
throw exception;
} catch (Exception exception) {
- recordAudit("UPDATE_GROUP_SETTINGS", name,
exception.getMessage(), "FAILED");
+ recordAudit("UPDATE_GROUP_SETTINGS", name, "updated " +
updatedBrokers + "/" + totalBrokers
+ + " brokers before failure: " +
exception.getMessage(), "FAILED");
throw classifyBrokerFailure(exception, "update consumer group
settings");
}
});
}
+ private static String describeSwitch(Boolean value) {
+ return value == null ? "preserved" : value.toString();
+ }
+
private ConsumerGroupVO createConsumerGroup(MQAdminExt admin,
ConsumerGroupVO group) {
String groupName = group.getName();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
index 736dd3136..23cd80a99 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
@@ -238,16 +238,41 @@ class ConsumerGroupControllerTest {
void consumerGroupSettingsUpdateShouldValidateAndDelegate() throws
Exception {
Map<String, Object> body = Map.of("instanceId", "instance-a", "name",
"cg-orders",
"retryQueueNums", 2, "retryMaxTimes", 8);
- when(metadataService.updateConsumerGroupSettings("instance-a",
"cg-orders", 2, 8))
+ when(metadataService.updateConsumerGroupSettings("instance-a",
"cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, null, null, null)))
.thenReturn(ConsumerGroupSettingsVO.builder().groupName("cg-orders").retryQueueNums(2)
- .retryMaxTimes(8).build());
+
.retryMaxTimes(8).consumeEnable(true).consumeMessageOrderly(false)
+ .consumeBroadcastEnable(true).build());
mockMvc.perform(post("/api/groups/settings").contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(body)))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.retryMaxTimes").value(8));
- verify(metadataService).updateConsumerGroupSettings("instance-a",
"cg-orders", 2, 8);
+ verify(metadataService).updateConsumerGroupSettings("instance-a",
"cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, null, null, null));
+ }
+
+ @Test
+ void consumerGroupSettingsUpdateShouldForwardConsumptionSwitches() throws
Exception {
+ Map<String, Object> body = Map.of("instanceId", "instance-a", "name",
"cg-orders",
+ "retryQueueNums", 2, "retryMaxTimes", 8, "consumeEnable",
false,
+ "consumeMessageOrderly", true);
+ when(metadataService.updateConsumerGroupSettings("instance-a",
"cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, false, true, null)))
+
.thenReturn(ConsumerGroupSettingsVO.builder().groupName("cg-orders").retryQueueNums(2)
+
.retryMaxTimes(8).consumeEnable(false).consumeMessageOrderly(true)
+ .consumeBroadcastEnable(false).build());
+
+
mockMvc.perform(post("/api/groups/settings").contentType(MediaType.APPLICATION_JSON)
+ .content(objectMapper.writeValueAsString(body)))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.consumeEnable").value(false))
+
.andExpect(jsonPath("$.data.consumeMessageOrderly").value(true))
+
.andExpect(jsonPath("$.data.consumeBroadcastEnable").value(false));
+
+ verify(metadataService).updateConsumerGroupSettings("instance-a",
"cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, false, true, null));
}
@Test
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 72be6c2ba..4e0993def 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
@@ -38,6 +38,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsCommand;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupSettingsVO;
import org.apache.rocketmq.studio.instance.group.ConsumerInstanceVO;
import org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetPreviewVO;
@@ -781,7 +782,8 @@ class RocketMQAdminClientImplTest {
return action.apply(selectedAdmin);
});
- ConsumerGroupSettingsVO settings =
adminClient.updateConsumerGroupSettings("instance-a", "cg-orders", 2, 8);
+ ConsumerGroupSettingsVO settings =
adminClient.updateConsumerGroupSettings("instance-a", "cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, null, null, null));
assertThat(settings.getRetryQueueNums()).isEqualTo(2);
assertThat(settings.getRetryMaxTimes()).isEqualTo(8);
@@ -793,6 +795,82 @@ class RocketMQAdminClientImplTest {
assertThat(captor.getValue().getRetryMaxTimes()).isEqualTo(8);
}
+ @Test
+ void updateConsumerGroupSettingsAppliesConsumptionSwitches() throws
Exception {
+ TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new
MybatisConfiguration(), ""), RmqGroup.class);
+ DefaultMQAdminExt selectedAdmin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ ClusterInfo clusterInfo = new ClusterInfo();
+ clusterInfo.setClusterAddrTable(new HashMap<>(Map.of("cluster-1", new
HashSet<>(List.of("broker-1")))));
+ BrokerData brokerData = new BrokerData();
+ brokerData.setBrokerName("broker-1");
+ brokerData.setBrokerAddrs(new HashMap<>(Map.of(0L, "10.0.0.1:10911")));
+ clusterInfo.setBrokerAddrTable(new HashMap<>(Map.of("broker-1",
brokerData)));
+ SubscriptionGroupConfig config = new SubscriptionGroupConfig();
+ config.setGroupName("cg-orders");
+ config.setConsumeEnable(true);
+ config.setConsumeMessageOrderly(false);
+ config.setConsumeBroadcastEnable(false);
+ config.setRetryQueueNums(1);
+ config.setRetryMaxTimes(16);
+ when(selectedAdmin.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(selectedAdmin.examineSubscriptionGroupConfig("10.0.0.1:10911",
"cg-orders")).thenReturn(config);
+ when(groupMapper.selectOne(any())).thenReturn(null);
+
doNothing().when(selectedAdmin).createAndUpdateSubscriptionGroupConfig(anyString(),
any());
+
when(runtimeAdminClientResolver.execute(org.mockito.ArgumentMatchers.eq("instance-a"),
any()))
+ .thenAnswer(invocation -> {
+ MqAdminExtFactory.AdminAction<?> action =
invocation.getArgument(1);
+ return action.apply(selectedAdmin);
+ });
+
+ ConsumerGroupSettingsVO settings =
adminClient.updateConsumerGroupSettings("instance-a", "cg-orders",
+ new ConsumerGroupSettingsCommand(2, 8, false, true, null));
+
+ assertThat(settings.isConsumeEnable()).isFalse();
+ assertThat(settings.isConsumeMessageOrderly()).isTrue();
+ assertThat(settings.isConsumeBroadcastEnable()).isFalse();
+ ArgumentCaptor<SubscriptionGroupConfig> captor =
ArgumentCaptor.forClass(SubscriptionGroupConfig.class);
+ verify(selectedAdmin).createAndUpdateSubscriptionGroupConfig(
+ org.mockito.ArgumentMatchers.eq("10.0.0.1:10911"),
captor.capture());
+ assertThat(captor.getValue().isConsumeEnable()).isFalse();
+ assertThat(captor.getValue().isConsumeMessageOrderly()).isTrue();
+ assertThat(captor.getValue().isConsumeBroadcastEnable()).isFalse();
+ assertThat(captor.getValue().getRetryQueueNums()).isEqualTo(2);
+ assertThat(captor.getValue().getRetryMaxTimes()).isEqualTo(8);
+ }
+
+ @Test
+ void getConsumerGroupSettingsReturnsConsumptionSwitches() throws Exception
{
+ DefaultMQAdminExt selectedAdmin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ ClusterInfo clusterInfo = new ClusterInfo();
+ clusterInfo.setClusterAddrTable(new HashMap<>(Map.of("cluster-1", new
HashSet<>(List.of("broker-1")))));
+ BrokerData brokerData = new BrokerData();
+ brokerData.setBrokerName("broker-1");
+ brokerData.setBrokerAddrs(new HashMap<>(Map.of(0L, "10.0.0.1:10911")));
+ clusterInfo.setBrokerAddrTable(new HashMap<>(Map.of("broker-1",
brokerData)));
+ SubscriptionGroupConfig config = new SubscriptionGroupConfig();
+ config.setGroupName("cg-orders");
+ config.setConsumeEnable(false);
+ config.setConsumeMessageOrderly(true);
+ config.setConsumeBroadcastEnable(false);
+ config.setRetryQueueNums(1);
+ config.setRetryMaxTimes(16);
+ when(selectedAdmin.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(selectedAdmin.examineSubscriptionGroupConfig("10.0.0.1:10911",
"cg-orders")).thenReturn(config);
+
when(runtimeAdminClientResolver.execute(org.mockito.ArgumentMatchers.eq("instance-a"),
any()))
+ .thenAnswer(invocation -> {
+ MqAdminExtFactory.AdminAction<?> action =
invocation.getArgument(1);
+ return action.apply(selectedAdmin);
+ });
+
+ ConsumerGroupSettingsVO settings =
adminClient.getConsumerGroupSettings("instance-a", "cg-orders");
+
+ assertThat(settings.isConsumeEnable()).isFalse();
+ assertThat(settings.isConsumeMessageOrderly()).isTrue();
+ assertThat(settings.isConsumeBroadcastEnable()).isFalse();
+ assertThat(settings.getRetryQueueNums()).isEqualTo(1);
+ assertThat(settings.getRetryMaxTimes()).isEqualTo(16);
+ }
+
@Test
void deleteConsumerGroupUsesSelectedInstanceAdmin() throws Exception {
TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new
MybatisConfiguration(), ""), RmqGroup.class);
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index 185fe3271..60f93f20e 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -145,6 +145,9 @@ export interface ConsumerGroupSettings {
groupName: string;
retryQueueNums: number;
retryMaxTimes: number;
+ consumeEnable?: boolean;
+ consumeMessageOrderly?: boolean;
+ consumeBroadcastEnable?: boolean;
}
export interface QueueProgress {
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index aea364109..45743d7e0 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -1356,6 +1356,54 @@ describe('Consumer page', () => {
expect(await screen.findByText('消费组配置已保存')).toBeInTheDocument();
});
+ it('edits consumption switches from the detail modal settings tab', async ()
=> {
+ vi.mocked(consumerService.getConsumerGroupSettings).mockResolvedValue({
+ groupName: 'remote-cg',
+ retryQueueNums: 1,
+ retryMaxTimes: 16,
+ consumeEnable: false,
+ consumeMessageOrderly: true,
+ consumeBroadcastEnable: false,
+ });
+ vi.mocked(consumerService.updateConsumerGroupSettings).mockResolvedValue({
+ groupName: 'remote-cg',
+ retryQueueNums: 1,
+ retryMaxTimes: 16,
+ consumeEnable: true,
+ consumeMessageOrderly: true,
+ consumeBroadcastEnable: false,
+ });
+ const user = userEvent.setup();
+ renderWithProviders(<ConsumerPage />);
+
+ const row = await screen.findByRole('row', { name: /remote-cg/ });
+ await user.click(within(row).getByRole('button', { name: /详\s*情/ }));
+
+ const dialog = await screen.findByRole('dialog', { name: /remote-cg/ });
+ await user.click(within(dialog).getByRole('tab', { name: /配\s*置/ }));
+
+ const consumeSwitch = await within(dialog).findByRole('switch', { name:
'启用消费' });
+ expect(consumeSwitch).not.toBeChecked();
+ expect(within(dialog).getByRole('switch', { name: '顺序消费' })).toBeChecked();
+ expect(within(dialog).getByRole('switch', { name: '广播消费'
})).not.toBeChecked();
+
+ await user.click(consumeSwitch);
+ await user.click(within(dialog).getByRole('button', { name: /保\s*存/ }));
+
+ await waitFor(() => {
+
expect(consumerService.updateConsumerGroupSettings).toHaveBeenCalledWith({
+ instanceId: 'instance-1',
+ name: 'remote-cg',
+ retryQueueNums: 1,
+ retryMaxTimes: 16,
+ consumeEnable: true,
+ consumeMessageOrderly: true,
+ consumeBroadcastEnable: false,
+ });
+ });
+ expect((await screen.findAllByText('消费组配置已保存')).length).toBeGreaterThan(0);
+ });
+
it('renders an unknown (-1) lag as unavailable in the table and the lag
detail', async () => {
const user = userEvent.setup();
vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index 34036b1ff..08eb5c523 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -42,6 +42,7 @@ import {
Tooltip,
Spin,
Progress,
+ Switch,
message,
} from 'antd';
import {
@@ -273,7 +274,13 @@ const ConsumerPageContent = ({
const [settingsGroup, setSettingsGroup] = useState<ConsumerGroup |
null>(null);
const [settingsLoading, setSettingsLoading] = useState(false);
const [settingsSubmitting, setSettingsSubmitting] = useState(false);
- const [settingsForm] = Form.useForm<{ retryQueueNums: number; retryMaxTimes:
number }>();
+ const [settingsForm] = Form.useForm<{
+ retryQueueNums: number;
+ retryMaxTimes: number;
+ consumeEnable?: boolean;
+ consumeMessageOrderly?: boolean;
+ consumeBroadcastEnable?: boolean;
+ }>();
const [createModalOpen, setCreateModalOpen] = useState(false);
const [form] = Form.useForm();
const [dataTypeValue, setDataTypeValue] = useState<string |
undefined>(undefined);
@@ -312,6 +319,12 @@ const ConsumerPageContent = ({
const groupRequestIdRef = useRef(0);
const stackRequestIdRef = useRef(0);
const settingsRequestIdRef = useRef(0);
+ // Consumption switches as loaded from the broker, used to detect high-risk
changes
+ // (disabling consumption / toggling ordered consumption) that need a
confirm before saving.
+ const originalSettingsRef = useRef<{
+ consumeEnable?: boolean;
+ consumeMessageOrderly?: boolean;
+ } | null>(null);
const [autoRefresh, setAutoRefresh] = useState(false);
const silentRefreshRef = useRef(false);
@@ -496,6 +509,10 @@ const ConsumerPageContent = ({
const settings = await getConsumerGroupSettings(group.name,
selectedInstanceId);
if (requestId === settingsRequestIdRef.current) {
settingsForm.setFieldsValue(settings);
+ originalSettingsRef.current = {
+ consumeEnable: settings.consumeEnable,
+ consumeMessageOrderly: settings.consumeMessageOrderly,
+ };
}
} catch {
if (requestId === settingsRequestIdRef.current) {
@@ -522,6 +539,32 @@ const ConsumerPageContent = ({
const saveSettings = async () => {
if (!settingsGroup || !selectedInstanceId) return;
const values = await settingsForm.validateFields();
+ const original = originalSettingsRef.current;
+ const risks: string[] = [];
+ if (original && values.consumeEnable === false && original.consumeEnable
!== false) {
+ risks.push('关闭「启用消费」会立即停止该消费组的消息消费,可能导致消息堆积');
+ }
+ if (
+ original &&
+ values.consumeMessageOrderly !== undefined &&
+ values.consumeMessageOrderly !== original.consumeMessageOrderly
+ ) {
+ risks.push('切换「顺序消费」会改变该消费组的消费语义,可能影响消息顺序与吞吐');
+ }
+ if (risks.length > 0) {
+ const confirmed = await new Promise<boolean>((resolve) => {
+ Modal.confirm({
+ title: '确认修改高危消费配置?',
+ content: `${risks.join(';')}。`,
+ okText: '确认修改',
+ okButtonProps: { danger: true },
+ cancelText: '取消',
+ onOk: () => resolve(true),
+ onCancel: () => resolve(false),
+ });
+ });
+ if (!confirmed) return;
+ }
setSettingsSubmitting(true);
try {
const saved = await updateConsumerGroupSettings({
@@ -2050,6 +2093,23 @@ const ConsumerPageContent = ({
>
<InputNumber min={1} max={128} style={{ width: '100%'
}} />
</Form.Item>
+ <Form.Item label="启用消费" name="consumeEnable"
valuePropName="checked">
+ <Switch />
+ </Form.Item>
+ <Form.Item
+ label="顺序消费"
+ name="consumeMessageOrderly"
+ valuePropName="checked"
+ >
+ <Switch />
+ </Form.Item>
+ <Form.Item
+ label="广播消费"
+ name="consumeBroadcastEnable"
+ valuePropName="checked"
+ >
+ <Switch />
+ </Form.Item>
<Form.Item style={{ marginBottom: 0 }}>
<Button
type="primary"