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 369a7da84 feat(group): support consumer group bulk import and export
(#2638)
369a7da84 is described below
commit 369a7da845cad12237514985edac70933eb64249
Author: coder999o <[email protected]>
AuthorDate: Mon Aug 31 21:00:39 2026 +0800
feat(group): support consumer group bulk import and export (#2638)
---
.claude/skills/pr-review/SKILL.md | 265 ---------------------
.../rocketmq/studio/common/util/CsvUtil.java | 48 ++++
.../instance/group/ConsumerGroupController.java | 30 +++
.../instance/group/ImportConsumerGroupsDTO.java | 37 +++
.../group/ImportConsumerGroupsResultVO.java | 52 ++++
.../studio/instance/topic/MetadataService.java | 100 ++++++++
.../rocketmq/studio/ops/audit/AuditService.java | 21 +-
.../rocketmq/studio/common/util/CsvUtilTest.java | 38 +++
.../group/ConsumerGroupControllerTest.java | 54 +++++
.../studio/instance/topic/MetadataServiceTest.java | 119 +++++++++
web/src/api/metadata.test.ts | 61 +++++
web/src/api/metadata.ts | 34 ++-
.../pages/instance/__tests__/ConsumerPage.test.tsx | 98 ++++++--
web/src/pages/instance/consumer.tsx | 69 +++---
web/src/services/consumerService.ts | 66 +++++
15 files changed, 736 insertions(+), 356 deletions(-)
diff --git a/.claude/skills/pr-review/SKILL.md
b/.claude/skills/pr-review/SKILL.md
deleted file mode 100644
index c92aea686..000000000
--- a/.claude/skills/pr-review/SKILL.md
+++ /dev/null
@@ -1,265 +0,0 @@
----
-name: pr-review
-description: RocketMQ Studio 的 PR 评审助手。输入一个 GitHub PR 链接,自动用 gh 拉取 PR 及其关联
Issue,切出本地分支,编译前端与后端,用 docker compose 拉起整个项目(不改端口),并产出结构化的 PR 分析总结。当用户提到 评审
PR、review PR、看一下这个 PR、PR 分析、拉 PR 编译、检查 PR、pr review、审查合并请求 等场景时触发;即使用户只贴了一个
GitHub PR 链接,也应触发此 skill。
----
-
-# RocketMQ Studio PR 评审
-
-对 RocketMQ Studio(`apache/rocketmq-dashboard`,分支 `rocketmq-studio`)的一个 GitHub
PR 做端到端评审:拉取 PR 与关联 Issue → 本地切分支 → 编译前后端 → docker compose 拉起 → 输出结构化分析总结。
-
-## 前置条件
-
-- 已安装并登录 `gh`(`gh auth status` 确认)。
-- 已安装 `docker`(含 `docker compose`)、`node`(>=20)、`npm`、`mvn`(JDK 21)。
-- 当前工作目录为项目根目录(含 `server/`、`web/`、`deploy/`)。
-- 后端构建需要 JDK 21(`JAVA_HOME` 指向 JDK 21,或使用 `mvn -B -ntp` 配合系统 JDK 21)。
-
-> 端口约定(**禁止修改**):前端 6789(Nginx)、后端 8888(Spring Boot)、NameServer 9876、Broker
10911、Proxy 8080/8081。
-
-## 输入
-
-一个 PR 链接,例如:`https://github.com/apache/rocketmq-dashboard/pull/123`
-
-从链接中解析出 PR 编号 `<PR>`(URL 最后一段数字)。
-
-## 标准流程(Pipeline)
-
-评审按以下 8 个阶段顺序执行。每个阶段独立产出结果,任一阶段失败不阻断后续步骤,但需在总结中标记 ❌。
-
-```
-Stage 1 拉取元信息 gh pr view → JSON + diff
-Stage 2 标题规范检查 正则校验 [Studio] type: description
-Stage 3 关联 Issue gh issue view(若有 Closes/Fixes 引用)
-Stage 4 切出评审分支 gh pr checkout → pr-review-<PR>
-Stage 5 预检 & 修复 Dockerfile style/ 目录修复(已知问题)
-Stage 6 编译后端 mvn package -DskipTests(含 checkstyle)
-Stage 7 编译前端 npm ci && npm run build(tsc + vite)
-Stage 8 Docker 部署 docker compose up -d --build + 健康检查
-```
-
-**完成后**:复原工作区(切回原分支 + stash pop)。
-
----
-
-## 各阶段详细步骤
-
-### Stage 1: 拉取 PR 元信息
-
-```bash
-gh pr view <PR> --repo apache/rocketmq-dashboard \
- --json
number,title,author,state,baseRefName,headRefName,url,body,additions,deletions,changedFiles,labels,commits,files,closingIssuesReferences
-```
-
-重点关注:
-- `baseRefName` 是否为 `rocketmq-studio`(目标分支应为它,否则在总结中提示)。
-- `files` / `changedFiles` / `additions` / `deletions`:变更范围与体量。
-- `closingIssuesReferences`:PR 声明会关闭的 Issue(`Closes #N`)。
-
-同时拉取 diff 供后续分析:
-
-```bash
-gh pr diff <PR> --repo apache/rocketmq-dashboard > /tmp/pr-<PR>.diff
-```
-
-### Stage 2: 检查 PR 标题是否规范
-
-规范来源:`docs/contributing.md`「提交规范」+ `README.md`「Commit Format」+ 本仓库既有 PR/commit
历史,基于 [Conventional Commits](https://www.conventionalcommits.org/) 并带项目前缀。
-
-**格式**:`[Studio] <type>: <description>`
-
-- **`[Studio]`** 为项目前缀(必须,大小写敏感,首字母大写 `Studio`)。
-- **type** 必须为以下之一(小写):`feat`(新功能)、`fix`(修复
Bug)、`docs`(文档)、`refactor`(重构)、`test`(测试)、`chore`(构建/工具)、`perf`(性能)。
-- `[Studio]` 与 `type` 之间有一个空格;type 后紧跟半角冒号 `:` 加一个空格,再接 **description**。
-- description 用英文小写祈使句,简洁描述改动,结尾不加句号。
-- squash 合并后 GitHub 会在标题末尾追加 ` (#PR号)`,属正常现象,检查时应先剥离该后缀。
-
-校验正则(先去掉可能存在的 ` (#N)` 尾巴):
-
-```bash
-TITLE=$(gh pr view <PR> --repo apache/rocketmq-dashboard --json title -q
.title)
-CLEAN=$(printf '%s' "$TITLE" | sed -E 's/ \(#[0-9]+\)$//')
-if printf '%s' "$CLEAN" | grep -Eq '^\[Studio\]
(feat|fix|docs|refactor|test|chore|perf): .+'; then
- echo "标题规范 ✅: $TITLE"
-else
- echo "标题不规范 ❌: $TITLE"
-fi
-```
-
-对照参考(本仓库历史):
-- `[Studio] fix: validate audit query and cleanup parameters` ✅
-- `[Studio] fix: connect K8s certificate page to backend APIs` ✅
-- `[Studio] feat: add i18n language context with useLanguage alias` ✅
-- `[Studio][fix] Connect K8s certificate page to backend APIs` ❌(type
应在冒号前,不应使用 `[fix]` 方括号)
-- `[studio] feat: extend translation keys` ❌(`studio` 应为 `Studio`,首字母大写)
-
-不规范时在总结中明确指出问题并给出建议标题。
-
-### Stage 3: 拉取关联 Issue(若有)
-
-若 `closingIssuesReferences` 非空,或 PR 正文中出现 `#N` / `Closes #N` / `Fixes #N`
引用,逐个拉取:
-
-```bash
-gh issue view <ISSUE> --repo apache/rocketmq-dashboard \
- --json number,title,state,body,labels,url
-```
-
-用于判断 PR 是否真正解决了 Issue 描述的问题(需求对齐度)。
-
-### Stage 4: 本地切出评审分支
-
-不要污染当前分支。用 `gh` 直接 checkout PR 分支(会自动创建本地分支):
-
-```bash
-# 记录当前分支以便复原
-ORIGINAL_BRANCH=$(git rev-parse --abbrev-ref HEAD)
-git stash push -u -m "pr-review-stash" 2>/dev/null || true
-
-gh pr checkout <PR> --repo apache/rocketmq-dashboard --branch pr-review-<PR>
-```
-
-若因权限/fork 无法直接 checkout,退化为手动 fetch:
-
-```bash
-git fetch origin pull/<PR>/head:pr-review-<PR>
-git checkout pr-review-<PR>
-```
-
-### Stage 5: 预检 & 修复(Dockerfile style/ 目录)
-
-**已知问题**:`server/Dockerfile` 在基准分支上缺少 `COPY style ./style`,导致 Maven checkstyle
插件在构建阶段找不到 `style/rmq_checkstyle.xml` 而报错,docker compose 无法启动。
-
-每次评审前检查并修复:
-
-```bash
-if ! grep -q 'COPY style' server/Dockerfile; then
- # 在 "COPY src ./src" 后插入 "COPY style ./style"
- sed -i '/^COPY src \.\/src$/a COPY style ./style' server/Dockerfile
- echo "✅ 已修复 Dockerfile: 添加 COPY style ./style"
-else
- echo "✅ Dockerfile 已包含 style/ 目录复制"
-fi
-```
-
-> 此修复仅用于本地评审,不影响 PR diff。若 PR 本身已修复此问题,此步为 no-op。
-
-### Stage 6: 编译后端
-
-严格对齐 `server/Dockerfile` 的构建方式(`mvn package -DskipTests`,含 checkstyle 校验):
-
-```bash
-cd server
-mvn -B -ntp clean package -DskipTests
-cd ..
-```
-
-- 编译失败 → 记录报错(编译错误 / checkstyle 违规),标记后端为 ❌,仍继续后续步骤并如实汇报。
-- 通过 → 标记 ✅。
-
-### Stage 7: 编译前端
-
-对齐 `web/Dockerfile`(`npm ci && npm run build`,即 `tsc -b && vite build`):
-
-```bash
-cd web
-npm ci
-npm run build
-cd ..
-```
-
-- 出现 TypeScript 类型错误或构建失败 → 标记前端 ❌,记录关键报错。
-- 通过 → 标记 ✅。
-
-### Stage 8: Docker 部署
-
-使用项目自带 compose 文件,**不修改任何端口**:
-
-```bash
-cd deploy
-docker compose down 2>/dev/null || true
-docker compose up -d --build
-cd ..
-```
-
-启动后做健康检查(不改端口):
-
-```bash
-# 前端(Nginx)
-curl -fsS -o /dev/null -w "web=%{http_code}\n" http://127.0.0.1:6789/ || echo
"web 未就绪"
-# 后端 actuator(经前端 nginx 反代或容器内),若暴露则直接探测
-curl -fsS -o /dev/null -w "api=%{http_code}\n"
http://127.0.0.1:6789/actuator/health || echo "api 未就绪"
-docker compose -f deploy/docker-compose.yml ps
-```
-
-评审结束后清理(询问用户或默认保留,视会话上下文而定):
-
-```bash
-docker compose -f deploy/docker-compose.yml down
-```
-
----
-
-## 复原工作区
-
-评审完成后切回原分支并恢复暂存:
-
-```bash
-git checkout "$ORIGINAL_BRANCH"
-git stash pop 2>/dev/null || true
-# 如需删除评审分支:git branch -D pr-review-<PR>
-```
-
----
-
-## 输出:PR 分析总结
-
-用中文输出一份结构化 Markdown 总结,包含以下部分:
-
-### 1. 概览
-- 标题、作者、状态、源分支 → 目标分支、PR 链接。
-- **PR 标题规范**:✅/❌(不规范时标出具体问题并附建议标题)。
-- 变更体量:`+additions / -deletions`,改动文件数。
-- 关联 Issue:编号、标题、链接(若有)。
-
-### 2. 需求对齐
-- 简述 PR 目标(来自正文)。
-- 若有关联 Issue,逐条比对 Issue 诉求与 PR 实现,判断是否覆盖 / 部分覆盖 / 偏离。
-
-### 3. 构建与运行结果
-用表格汇总:
-
-| 检查项 | 结果 | 说明 |
-|---|---|---|
-| 后端编译 (`mvn package`) | ✅/❌ | 关键报错摘要 |
-| 前端编译 (`npm run build`) | ✅/❌ | 关键报错摘要 |
-| docker compose 启动 | ✅/❌ | 容器状态 / 端口 6789、8888 |
-| 前端健康检查 | ✅/❌ | HTTP 状态码 |
-
-### 4. 变更分析
-- 按模块归类改动(前端页面/组件、后端 controller/service/domain、部署、文档等)。
-- 结合六边形架构(server 用 ArchUnit 约束)判断分层是否合理。
-- i18n:新增前端文案是否中英文双语(`web/src/i18n/`)。
-- 提交规范:commit message 是否符合 `[Studio] type: description` 格式。
-
-### 5. 风险与建议
-- 潜在逻辑问题、边界情况、安全风险(如凭据明文、公网暴露)。
-- 缺失的测试 / 文档 / i18n。
-- 明确的改进建议(可执行、可定位到文件)。
-
-### 6. 评审结论
-给出倾向:**Approve** / **Request Changes** / **Comment**,并用一句话说明理由。
-
----
-
-## 自由发挥补充能力
-
-- **变更规模自适应**:大 PR(改动文件多)先按目录聚合概述再抽样精读核心文件;小 PR 可逐文件过。
-- **checkstyle / lint 单独复核**:后端 `mvn checkstyle:check`,前端 `npm run
lint`,把风格问题与逻辑问题分开汇报。
-- **失败即止但汇报完整**:任一步失败不阻断总结,如实记录并继续能做的检查。
-- **可选发布评审意见**:用户明确要求时,可用 `gh pr comment <PR> --repo apache/rocketmq-dashboard
--body-file <file>` 或 `gh pr review <PR> --comment/--approve/--request-changes`
提交(默认只本地产出,不自动发布)。
-
-## 注意事项
-
-- **不修改端口**:任何环节都使用既有端口,不得改 compose / nginx / env 中的端口映射。
-- **不污染主分支**:评审在独立分支进行,结束后复原。
-- **只读默认**:默认不向 GitHub 写入评论/评审,除非用户明确要求。
-- **凭据安全**:分析 diff 时若发现 AK/SK、password、token 等明文凭据,作为高风险项在总结中显著标注。
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/util/CsvUtil.java
b/server/src/main/java/org/apache/rocketmq/studio/common/util/CsvUtil.java
new file mode 100644
index 000000000..0603dd7e8
--- /dev/null
+++ b/server/src/main/java/org/apache/rocketmq/studio/common/util/CsvUtil.java
@@ -0,0 +1,48 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.common.util;
+
+/**
+ * Shared CSV rendering helpers used by export endpoints. Cells are always
quoted and
+ * values starting with formula characters ({@code = + - @ \t \r \n}) are
prefixed with
+ * a single quote to prevent spreadsheet formula injection.
+ */
+public final class CsvUtil {
+
+ public static final String CRLF = "\r\n";
+
+ private CsvUtil() {
+ }
+
+ public static void appendRow(StringBuilder csv, Object... values) {
+ for (int i = 0; i < values.length; i++) {
+ if (i > 0) {
+ csv.append(',');
+ }
+ csv.append(toCell(values[i]));
+ }
+ csv.append(CRLF);
+ }
+
+ public static String toCell(Object value) {
+ String text = value == null ? "" : value.toString();
+ if (!text.isEmpty() && "=+-@\t\r\n".indexOf(text.charAt(0)) >= 0) {
+ text = "'" + text;
+ }
+ return '"' + text.replace("\"", "\"\"") + '"';
+ }
+}
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 751a734d3..8c965bede 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
@@ -30,6 +30,8 @@ import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.List;
@RestController
@@ -86,6 +88,23 @@ public class ConsumerGroupController {
return Result.ok(metadataService.refreshConsumerGroup(instanceId,
name));
}
+ @GetMapping("/export")
+ public Result<String> exportConsumerGroups(
+ @RequestParam(required = false) String instanceId,
+ @RequestParam(required = false) String search,
+ @RequestParam(required = false) String subscriptionMode,
+ @RequestParam(required = false) String names) {
+ return Result.ok(metadataService.exportConsumerGroups(instanceId,
search,
+ subscriptionMode, parseNames(names)));
+ }
+
+ @PostMapping("/import")
+ public Result<ImportConsumerGroupsResultVO> importConsumerGroups(
+ @Valid @RequestBody ImportConsumerGroupsDTO request) {
+ String instanceId =
instanceService.normalizeIdentifier(request.getInstanceId());
+ return Result.ok(metadataService.importConsumerGroups(instanceId,
request.getGroups()));
+ }
+
@GetMapping("/{name}/progress")
public Result<List<QueueProgressVO>> getGroupProgress(
@PathVariable String name,
@@ -127,4 +146,15 @@ public class ConsumerGroupController {
request.getTimestamp(), request.getTopic());
return Result.ok();
}
+
+ private List<String> parseNames(String names) {
+ if (names == null || names.isBlank()) {
+ return Collections.emptyList();
+ }
+ return Arrays.stream(names.split(","))
+ .map(String::trim)
+ .filter(name -> !name.isEmpty())
+ .distinct()
+ .toList();
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsDTO.java
new file mode 100644
index 000000000..b300e2e30
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsDTO.java
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.group;
+
+import jakarta.validation.Valid;
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.NotEmpty;
+import jakarta.validation.constraints.Size;
+import lombok.Data;
+
+import java.util.List;
+
+@Data
+public class ImportConsumerGroupsDTO {
+
+ @NotBlank(message = "instanceId is required")
+ private String instanceId;
+
+ @Valid
+ @NotEmpty(message = "groups is required")
+ @Size(max = 100, message = "At most 100 consumer groups are allowed per
import")
+ private List<CreateConsumerGroupDTO> groups;
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsResultVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsResultVO.java
new file mode 100644
index 000000000..6e3afe7c6
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ImportConsumerGroupsResultVO.java
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.group;
+
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+import java.util.List;
+
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class ImportConsumerGroupsResultVO {
+
+ private int imported;
+
+ private int failed;
+
+ private List<ConsumerGroupVO> groups;
+
+ private List<Failure> failures;
+
+ @Data
+ @Builder
+ @NoArgsConstructor
+ @AllArgsConstructor
+ public static class Failure {
+
+ private int index;
+
+ private String name;
+
+ private String message;
+ }
+}
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 4617fae3c..cd9a47b7b 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
@@ -24,7 +24,11 @@ import
org.apache.rocketmq.studio.provider.apache.AdminClient;
import org.apache.rocketmq.studio.provider.apache.MetadataProvider;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.util.CsvUtil;
+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.ConsumerGroupSettingsVO;
@@ -37,7 +41,10 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
+import java.util.ArrayList;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
import java.util.function.Supplier;
@Slf4j
@@ -357,6 +364,99 @@ public class MetadataService {
}
}
+ private static final int MAX_IMPORT_GROUPS = 100;
+
+ public ImportConsumerGroupsResultVO importConsumerGroups(String instanceId,
+
List<CreateConsumerGroupDTO> groups) {
+ if (!StringUtils.hasText(instanceId)) {
+ throw new BusinessException(400, "instanceId is required");
+ }
+ if (groups == null || groups.isEmpty()) {
+ throw new BusinessException(400, "groups is required");
+ }
+ if (groups.size() > MAX_IMPORT_GROUPS) {
+ throw new BusinessException(400, "At most 100 consumer groups are
allowed per import");
+ }
+
+ String normalizedInstanceId = normalizeInstanceId(instanceId);
+ List<ConsumerGroupVO> imported = new ArrayList<>();
+ List<ImportConsumerGroupsResultVO.Failure> failures = new
ArrayList<>();
+ for (int index = 0; index < groups.size(); index++) {
+ CreateConsumerGroupDTO request = groups.get(index);
+ String name = request == null ? null : request.getName();
+ try {
+ if (request == null) {
+ throw new BusinessException(400, "consumer group request
is required");
+ }
+ ConsumerGroupVO group = request.toConsumerGroupVO();
+ group.setInstanceId(normalizedInstanceId);
+ imported.add(createConsumerGroup(group));
+ } catch (Exception exception) {
+ failures.add(ImportConsumerGroupsResultVO.Failure.builder()
+ .index(index)
+ .name(name)
+ .message(StringUtils.hasText(exception.getMessage())
+ ? exception.getMessage() : "Failed to create
consumer group")
+ .build());
+ }
+ }
+
+ return ImportConsumerGroupsResultVO.builder()
+ .imported(imported.size())
+ .failed(failures.size())
+ .groups(imported)
+ .failures(failures)
+ .build();
+ }
+
+ public String exportConsumerGroups(String instanceId, String search,
String subscriptionMode,
+ List<String> names) {
+ instanceId = normalizeInstanceId(instanceId);
+ List<ConsumerGroupVO> groups = new
ArrayList<>(listConsumerGroups(instanceId, null, search));
+ Set<String> selectedNames = new HashSet<>(names == null ? List.of() :
names);
+ if (!selectedNames.isEmpty()) {
+ groups.removeIf(group -> !selectedNames.contains(group.getName()));
+ }
+ if (StringUtils.hasText(subscriptionMode) &&
!"ALL".equals(subscriptionMode)) {
+ groups.removeIf(group ->
!subscriptionMode.equals(toText(group.getSubscriptionMode())));
+ }
+
+ groups.sort(this::compareNames);
+ return buildConsumerGroupCsv(groups);
+ }
+
+ private int compareNames(ConsumerGroupVO left, ConsumerGroupVO right) {
+ String leftName = left.getName() == null ? "" : left.getName();
+ String rightName = right.getName() == null ? "" : right.getName();
+ return leftName.compareTo(rightName);
+ }
+
+ private String buildConsumerGroupCsv(List<ConsumerGroupVO> groups) {
+ StringBuilder csv = new StringBuilder();
+ CsvUtil.appendRow(csv, "Name", "Namespace", "Cluster ID",
"Subscription Mode", "Consume Type",
+ "Online Instances", "Total Lag", "Delay Seconds",
"Subscription Data Type",
+ "Delivery Order Type", "Retry Max Times", "Subscribed Topics",
"Created At", "Updated At");
+ for (ConsumerGroupVO group : groups) {
+ CsvUtil.appendRow(csv, group.getName(), group.getNamespace(),
group.getClusterId(),
+ toText(group.getSubscriptionMode()),
toText(group.getConsumeType()),
+ group.getOnlineInstances(), group.getTotalLag(),
group.getDelaySeconds(),
+ group.getSubscriptionDataType(),
group.getDeliveryOrderType(), group.getRetryMaxTimes(),
+ String.join(";", group.getSubscribedTopics() == null ?
List.of() : group.getSubscribedTopics()),
+ group.getGmtCreate(), group.getGmtModified());
+ }
+ return csv.toString();
+ }
+
+ private String toText(Object value) {
+ if (value == null) {
+ return "";
+ }
+ if (value instanceof SubscriptionMode mode) {
+ return mode.name();
+ }
+ return value.toString();
+ }
+
private String requireName(String value, String fieldName) {
if (!StringUtils.hasText(value)) {
throw new BusinessException(400, fieldName + " is required");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
index 83ebf6635..b68bd5df8 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.ops.audit;
import org.apache.rocketmq.studio.auth.AuthenticatedUserContext;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.util.CsvUtil;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@@ -69,7 +70,7 @@ public class AuditService {
}
StringBuilder csv = new StringBuilder("\uFEFF").append(CSV_HEADER);
for (AuditRecordVO record : page.getItems()) {
- appendCsvRow(csv,
+ CsvUtil.appendRow(csv,
record.getTimestamp(),
record.getOperator(),
record.getOperationType(),
@@ -142,24 +143,6 @@ public class AuditService {
start, end, result, page, pageSize);
}
- private void appendCsvRow(StringBuilder csv, Object... values) {
- for (int i = 0; i < values.length; i++) {
- if (i > 0) {
- csv.append(',');
- }
- csv.append(toCsvCell(values[i]));
- }
- csv.append("\r\n");
- }
-
- private String toCsvCell(Object value) {
- String text = value == null ? "" : value.toString();
- if (!text.isEmpty() && "=+-@\t\r\n".indexOf(text.charAt(0)) >= 0) {
- text = "'" + text;
- }
- return '"' + text.replace("\"", "\"\"") + '"';
- }
-
private LocalDateTime parseDate(String dateStr, boolean startOfDay, String
parameterName) {
if (dateStr == null || dateStr.isEmpty()) {
return null;
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/common/util/CsvUtilTest.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/CsvUtilTest.java
new file mode 100644
index 000000000..bf739869d
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/CsvUtilTest.java
@@ -0,0 +1,38 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.common.util;
+
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class CsvUtilTest {
+
+ @Test
+ void appendRowShouldQuoteAndTerminateWithCrlfTest() {
+ StringBuilder csv = new StringBuilder();
+ CsvUtil.appendRow(csv, "Name", null, 7);
+ assertThat(csv.toString()).isEqualTo("\"Name\",\"\",\"7\"\r\n");
+ }
+
+ @Test
+ void toCellShouldEscapeQuotesAndFormulaPrefixesTest() {
+ assertThat(CsvUtil.toCell("say \"hi\"")).isEqualTo("\"say
\"\"hi\"\"\"");
+ assertThat(CsvUtil.toCell("=SUM(A1)")).isEqualTo("\"'=SUM(A1)\"");
+ assertThat(CsvUtil.toCell("+cmd")).isEqualTo("\"'+cmd\"");
+ }
+}
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 a7a260f57..04a23a7f7 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
@@ -99,6 +99,23 @@ class ConsumerGroupControllerTest {
verify(metadataService).listConsumerGroupsPage("instance-a",
"cluster-a", "orders", 2, 20);
}
+ @Test
+ void exportConsumerGroupsShouldPassViewFiltersAndSelectedNames() throws
Exception {
+ when(metadataService.exportConsumerGroups("instance-a", "orders",
"Pop",
+ List.of("cg-a", "cg-b"))).thenReturn("\"Name\"\n\"cg-a\"");
+
+ mockMvc.perform(get("/api/groups/export")
+ .param("instanceId", "instance-a")
+ .param("search", "orders")
+ .param("subscriptionMode", "Pop")
+ .param("names", "cg-a, cg-b,cg-a"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data").value("\"Name\"\n\"cg-a\""));
+
+ verify(metadataService).exportConsumerGroups("instance-a", "orders",
"Pop",
+ List.of("cg-a", "cg-b"));
+ }
+
@Test
void createConsumerGroupShouldPassValidatedRequest() throws Exception {
Map<String, Object> body = Map.of(
@@ -132,6 +149,43 @@ class ConsumerGroupControllerTest {
assertThat(captor.getValue().getRetryMaxTimes()).isEqualTo(8);
}
+ @Test
+ void importConsumerGroupsShouldNormalizeInstanceAndDelegateBatch() throws
Exception {
+ Map<String, Object> body = Map.of(
+ "instanceId", "7",
+ "groups", List.of(Map.of(
+ "name", "cg-orders",
+ "subscriptionMode", "Push",
+ "consumeType", "CLUSTERING",
+ "retryMaxTimes", 8
+ ))
+ );
+ ConsumerGroupVO created = new ConsumerGroupVO();
+ created.setName("cg-orders");
+ when(instanceService.normalizeIdentifier("7")).thenReturn("rocketmq1");
+ when(metadataService.importConsumerGroups(eq("rocketmq1"), any()))
+ .thenReturn(ImportConsumerGroupsResultVO.builder()
+ .imported(1)
+ .failed(0)
+ .groups(List.of(created))
+ .failures(List.of())
+ .build());
+
+ mockMvc.perform(post("/api/groups/import")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content(objectMapper.writeValueAsString(body)))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.imported").value(1))
+
.andExpect(jsonPath("$.data.groups[0].name").value("cg-orders"));
+
+ @SuppressWarnings("unchecked")
+ ArgumentCaptor<List<CreateConsumerGroupDTO>> captor =
ArgumentCaptor.forClass(List.class);
+ verify(metadataService).importConsumerGroups(eq("rocketmq1"),
captor.capture());
+ assertThat(captor.getValue()).hasSize(1);
+ assertThat(captor.getValue().get(0).getName()).isEqualTo("cg-orders");
+ assertThat(captor.getValue().get(0).getRetryMaxTimes()).isEqualTo(8);
+ }
+
@Test
void createConsumerGroupShouldRejectMissingName() throws Exception {
Map<String, Object> body = Map.of(
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index 82612234e..1c6250963 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -19,9 +19,13 @@ package org.apache.rocketmq.studio.instance.topic;
import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.apache.rocketmq.studio.common.domain.PageResult;
+import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.CreateConsumerGroupDTO;
+import org.apache.rocketmq.studio.instance.group.ImportConsumerGroupsResultVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
import org.apache.rocketmq.studio.provider.apache.AdminClient;
@@ -30,14 +34,18 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
+import java.time.LocalDateTime;
+import java.util.ArrayList;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.verify;
@@ -458,4 +466,115 @@ class MetadataServiceTest {
"cloud-instance", "topic=orders, timestamp=1784246400000",
"SUCCESS", null);
}
+ @Test
+ void exportConsumerGroupsShouldApplyFiltersSortingAndCsvEscaping() {
+ ConsumerGroupVO stale = consumerGroup("users-cg", "=formula", 2,
SubscriptionMode.Push);
+ ConsumerGroupVO unknownLag = consumerGroup("orders-unknown", "orders",
-1, SubscriptionMode.Pop);
+ ConsumerGroupVO highLag = consumerGroup("orders-high", "orders", 100,
SubscriptionMode.Pop);
+ ConsumerGroupVO lowLag = consumerGroup("orders-low", "orders", 5,
SubscriptionMode.Pop);
+ when(apacheProvider.listConsumerGroups("instance-a", "orders"))
+ .thenReturn(List.of(stale, unknownLag, lowLag, highLag));
+
+ String csv = metadataService.exportConsumerGroups("instance-a", "
orders ", "Pop",
+ List.of("orders-low", "orders-high", "orders-unknown"));
+
+ assertThat(csv).contains("\"Name\",\"Namespace\",\"Cluster ID\"");
+ assertThat(csv).contains("\"orders-high\",\"orders\"");
+ assertThat(csv).contains("\"orders-low\",\"orders\"");
+ assertThat(csv).contains("\"orders-unknown\",\"orders\"");
+ assertThat(csv).doesNotContain("users-cg");
+
assertThat(csv.indexOf("\"orders-high\"")).isLessThan(csv.indexOf("\"orders-low\""));
+
assertThat(csv.indexOf("\"orders-low\"")).isLessThan(csv.indexOf("\"orders-unknown\""));
+ assertThat(csv).contains("\"orders-topic;payments,topic\"");
+ verify(apacheProvider).listConsumerGroups("instance-a", "orders");
+ }
+
+ @Test
+ void exportConsumerGroupsShouldEscapeFormulaCells() {
+ ConsumerGroupVO group = consumerGroup("orders-cg", "=formula", 10,
SubscriptionMode.Push);
+ when(apacheProvider.listConsumerGroups("instance-a",
null)).thenReturn(List.of(group));
+
+ String csv = metadataService.exportConsumerGroups("instance-a", null,
null, List.of());
+
+ assertThat(csv).contains("\"'=formula\"");
+ }
+
+ @Test
+ void importConsumerGroupsShouldContinueAfterRowFailure() {
+ when(apacheProvider.createConsumerGroup(eq("instance-a"),
any(ConsumerGroupVO.class)))
+ .thenAnswer(invocation -> {
+ ConsumerGroupVO group = invocation.getArgument(1);
+ if ("cg-fail".equals(group.getName())) {
+ throw new BusinessException(500, "broker rejected
group");
+ }
+ group.setClusterId("cluster-a");
+ return group;
+ });
+
+ ImportConsumerGroupsResultVO result =
metadataService.importConsumerGroups("instance-a",
+ List.of(importRequest("cg-ok", "other-instance"),
importRequest("cg-fail", "other-instance")));
+
+ assertThat(result.getImported()).isEqualTo(1);
+ assertThat(result.getFailed()).isEqualTo(1);
+
assertThat(result.getGroups()).extracting(ConsumerGroupVO::getName).containsExactly("cg-ok");
+ assertThat(result.getFailures()).hasSize(1);
+ assertThat(result.getFailures().get(0).getIndex()).isEqualTo(1);
+ assertThat(result.getFailures().get(0).getName()).isEqualTo("cg-fail");
+ assertThat(result.getFailures().get(0).getMessage()).isEqualTo("broker
rejected group");
+
+ ArgumentCaptor<ConsumerGroupVO> captor =
ArgumentCaptor.forClass(ConsumerGroupVO.class);
+ verify(apacheProvider, org.mockito.Mockito.times(2))
+ .createConsumerGroup(eq("instance-a"), captor.capture());
+
assertThat(captor.getAllValues()).extracting(ConsumerGroupVO::getInstanceId)
+ .containsExactly("instance-a", "instance-a");
+ }
+
+
+ @Test
+ void importConsumerGroupsShouldRejectEmptyAndOversizedBatchesTest() {
+ assertThatThrownBy(() ->
metadataService.importConsumerGroups("instance-a", List.of()))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("groups is required");
+
+ List<CreateConsumerGroupDTO> oversized = new ArrayList<>();
+ for (int i = 0; i < 101; i++) {
+ oversized.add(importRequest("cg-" + i, null));
+ }
+ assertThatThrownBy(() ->
metadataService.importConsumerGroups("instance-a", oversized))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("At most 100");
+
+ verifyNoInteractions(apacheProvider);
+ }
+
+ private ConsumerGroupVO consumerGroup(String name, String namespace, long
lag, SubscriptionMode mode) {
+ ConsumerGroupVO group = new ConsumerGroupVO();
+ group.setName(name);
+ group.setNamespace(namespace);
+ group.setClusterId("cluster-a");
+ group.setSubscriptionMode(mode);
+ group.setConsumeType(ConsumeType.CLUSTERING);
+ group.setOnlineInstances(1);
+ group.setTotalLag(lag);
+ group.setDelaySeconds(3);
+ group.setSubscriptionDataType("NORMAL");
+ group.setRetryMaxTimes(16);
+ group.setSubscribedTopics(List.of("orders-topic", "payments,topic"));
+ group.setGmtCreate(LocalDateTime.of(2026, 8, 27, 10, 0));
+ group.setGmtModified(LocalDateTime.of(2026, 8, 27, 11, 0));
+ return group;
+ }
+
+ private CreateConsumerGroupDTO importRequest(String name, String
instanceId) {
+ CreateConsumerGroupDTO request = new CreateConsumerGroupDTO();
+ request.setName(name);
+ request.setInstanceId(instanceId);
+ request.setSubscriptionMode(SubscriptionMode.Push);
+ request.setConsumeType(ConsumeType.CLUSTERING);
+ request.setRetryMaxTimes(16);
+ request.setSubscriptionDataType("NORMAL");
+ return request;
+
+ }
+
}
diff --git a/web/src/api/metadata.test.ts b/web/src/api/metadata.test.ts
index d2c9631d8..867995474 100644
--- a/web/src/api/metadata.test.ts
+++ b/web/src/api/metadata.test.ts
@@ -21,10 +21,12 @@ import client from './client';
import {
createTopic,
deleteTopic,
+ exportConsumerGroups,
getConsumerStack,
getTopicConsumerPage,
getTopicConsumers,
getTopicRoutes,
+ importConsumerGroups,
listTopics,
sendTopicMessage,
} from './metadata';
@@ -141,4 +143,63 @@ describe('topic metadata API', () => {
sendTopicMessage({ topic: topic.name, instanceId: 'instance-1', body:
'{"id":1}' }),
).resolves.toMatchObject({ msgId: 'msg-1' });
});
+
+ it('passes consumer group import and export contracts through API
endpoints', async () => {
+ mock.onGet('/groups/export').reply((config) => {
+ expect(config.params).toEqual({
+ instanceId: 'instance-1',
+ search: 'orders',
+ subscriptionMode: 'Pop',
+ names: 'cg-a,cg-b',
+ });
+ return [200, { code: 200, data: '"Name"\n"cg-a"' }];
+ });
+ mock.onPost('/groups/import').reply((config) => {
+ expect(JSON.parse(config.data)).toEqual({
+ instanceId: 'instance-1',
+ groups: [
+ {
+ name: 'cg-a',
+ subscriptionMode: 'Push',
+ consumeType: 'CLUSTERING',
+ retryMaxTimes: 16,
+ },
+ ],
+ });
+ return [
+ 200,
+ {
+ code: 200,
+ data: {
+ imported: 1,
+ failed: 0,
+ groups: [],
+ failures: [],
+ },
+ },
+ ];
+ });
+
+ await expect(
+ exportConsumerGroups({
+ instanceId: 'instance-1',
+ search: 'orders',
+ subscriptionMode: 'Pop',
+ names: ['cg-a', 'cg-b'],
+ }),
+ ).resolves.toBe('"Name"\n"cg-a"');
+ await expect(
+ importConsumerGroups({
+ instanceId: 'instance-1',
+ groups: [
+ {
+ name: 'cg-a',
+ subscriptionMode: 'Push',
+ consumeType: 'CLUSTERING',
+ retryMaxTimes: 16,
+ },
+ ],
+ }),
+ ).resolves.toMatchObject({ imported: 1, failed: 0 });
+ });
});
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index 118fdf81d..b919f32b3 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -153,10 +153,27 @@ export interface ConsumerGroupPageQuery extends
ConsumerGroupQuery {
pageSize?: number;
}
-export interface ResetConsumerOffsetRequest {
- name: string;
- timestamp: number;
- topic: string;
+export interface ConsumerGroupExportQuery extends ConsumerGroupQuery {
+ names?: string[];
+ subscriptionMode?: string;
+}
+
+export interface ImportConsumerGroupsRequest {
+ instanceId: string;
+ groups: Partial<ConsumerGroup>[];
+}
+
+export interface ImportConsumerGroupsFailure {
+ index: number;
+ name?: string;
+ message: string;
+}
+
+export interface ImportConsumerGroupsResult {
+ imported: number;
+ failed: number;
+ groups: ConsumerGroup[];
+ failures: ImportConsumerGroupsFailure[];
}
// ─── Topic API ──────────────────────────────────────────────────
@@ -326,13 +343,14 @@ export async function resetConsumerOffset(data:
ResetConsumerOffsetRequest) {
await client.post('/groups/reset-offset', data);
}
-export async function importConsumerGroups(data: string) {
- await client.post('/groups/import', { data });
+export async function importConsumerGroups(data: ImportConsumerGroupsRequest) {
+ const res = await client.post<{ data: ImportConsumerGroupsResult
}>('/groups/import', data);
+ return res.data.data;
}
-export async function exportConsumerGroups(names?: string[]) {
+export async function exportConsumerGroups(params?: ConsumerGroupExportQuery) {
const res = await client.get<{ data: string }>('/groups/export', {
- params: { names: names?.join(',') },
+ params: { ...params, names: params?.names?.join(',') },
});
return res.data.data;
}
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index ab1966df9..eda784c6c 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -31,11 +31,13 @@ vi.mock('../../../services/consumerService', () => ({
batchDeleteConsumerGroups: vi.fn(),
createConsumerGroup: vi.fn(),
deleteConsumerGroup: vi.fn(),
+ exportConsumerGroups: vi.fn(),
getConsumerGroup: vi.fn(),
getConsumerProgress: vi.fn(),
getConsumerStack: vi.fn(),
getConsumerSubscriptions: vi.fn(),
getConsumerGroupSettings: vi.fn(),
+ importConsumerGroups: vi.fn(),
listAllConsumerGroups: vi.fn(),
updateConsumerGroupSettings: vi.fn(),
listConsumerGroupPage: vi.fn(),
@@ -134,6 +136,7 @@ describe('Consumer page', () => {
]);
vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(groupPage([group]));
vi.mocked(consumerService.listAllConsumerGroups).mockResolvedValue([group]);
+
vi.mocked(consumerService.exportConsumerGroups).mockResolvedValue('"Name"\n"remote-cg"');
vi.mocked(consumerService.refreshConsumerGroup).mockResolvedValue({
...group,
totalLag: 42,
@@ -155,6 +158,26 @@ describe('Consumer page', () => {
gmtModified: '2026-07-24T00:00:00Z',
}) as ConsumerGroup,
);
+ vi.mocked(consumerService.importConsumerGroups).mockResolvedValue({
+ imported: 1,
+ failed: 0,
+ groups: [
+ {
+ ...group,
+ name: 'imported-cg',
+ namespace: 'default',
+ clusterId: 'server-cluster',
+ onlineInstances: 0,
+ totalLag: 0,
+ delaySeconds: 0,
+ instances: [],
+ subscribedTopics: [],
+ gmtCreate: '2026-07-24T00:00:00Z',
+ gmtModified: '2026-07-24T00:00:00Z',
+ },
+ ],
+ failures: [],
+ });
vi.mocked(consumerService.getConsumerProgress).mockResolvedValue([
{
topic: 'remote-topic',
@@ -319,10 +342,13 @@ describe('Consumer page', () => {
size: params?.pageSize ?? 20,
});
});
- vi.mocked(consumerService.listAllConsumerGroups).mockImplementation(async
(params) =>
- params?.search
- ? exportGroups.filter((item) => item.name.includes(params.search ??
''))
- : exportGroups,
+ vi.mocked(consumerService.exportConsumerGroups).mockResolvedValue(
+ [
+ '"Name","Namespace","Subscribed Topics"',
+ '"orders-cg","remote-ns","orders-topic;payments,topic"',
+ '"orders-cg-archive","\'=archive","orders-topic"',
+ '"orders-cg-formula","\'\r=formula-risk","orders-topic"',
+ ].join('\n'),
);
renderWithProviders(<ConsumerPage />);
@@ -332,11 +358,13 @@ describe('Consumer page', () => {
await user.click(screen.getByRole('button', { name: /导出/ }));
await waitFor(() =>
- expect(consumerService.listAllConsumerGroups).toHaveBeenCalledWith({
+ expect(consumerService.exportConsumerGroups).toHaveBeenCalledWith({
instanceId: 'instance-1',
search: 'orders',
+ subscriptionMode: undefined,
}),
);
+ expect(consumerService.listAllConsumerGroups).not.toHaveBeenCalled();
expect(URL.createObjectURL).toHaveBeenCalledTimes(1);
expect(clickSpy).toHaveBeenCalledTimes(1);
expect(URL.revokeObjectURL).toHaveBeenCalledWith('blob:consumer-group-export');
@@ -757,25 +785,30 @@ describe('Consumer page', () => {
it('keeps per-row state when consumer group CSV import partially fails',
async () => {
vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(groupPage([]));
- vi.mocked(consumerService.createConsumerGroup).mockImplementation(
- async (data: Partial<ConsumerGroup>) => {
- if (data.name === 'cg-fail') throw new Error('broker rejected group');
- return {
+ vi.mocked(consumerService.importConsumerGroups).mockResolvedValue({
+ imported: 1,
+ failed: 1,
+ groups: [
+ {
...group,
- ...data,
- name: data.name ?? '',
+ name: 'cg-ok',
+ subscriptionMode: 'Push',
+ consumeType: 'CLUSTERING',
+ retryMaxTimes: 16,
+ subscriptionDataType: 'NORMAL',
namespace: 'default',
clusterId: 'server-cluster',
onlineInstances: 0,
totalLag: 0,
delaySeconds: 0,
instances: [],
- subscribedTopics: data.subscribedTopics ?? [],
+ subscribedTopics: [],
gmtCreate: '2026-07-24T00:00:00Z',
gmtModified: '2026-07-24T00:00:00Z',
- } as ConsumerGroup;
- },
- );
+ },
+ ],
+ failures: [{ index: 1, name: 'cg-fail', message: 'broker rejected group'
}],
+ });
const user = userEvent.setup();
instanceServiceMocks.listInstances.mockResolvedValue([
@@ -803,19 +836,32 @@ describe('Consumer page', () => {
screen.getByTestId('consumer-group-import-file'),
new File([csv], 'groups.csv'),
);
- expect(await screen.findByText('检测到 2 个
Group,将按顺序调用创建接口')).toBeInTheDocument();
+ expect(await screen.findByText('检测到 2 个
Group,将通过后端批量导入')).toBeInTheDocument();
await user.click(screen.getByRole('button', { name: '开始导入' }));
- await waitFor(() =>
expect(consumerService.createConsumerGroup).toHaveBeenCalledTimes(2));
- expect(consumerService.createConsumerGroup).toHaveBeenNthCalledWith(1, {
- name: 'cg-ok',
- subscriptionMode: 'Push',
- consumeType: 'CLUSTERING',
- retryMaxTimes: 16,
- subscriptionDataType: 'NORMAL',
- subscribedTopics: [],
- instanceId: 'instance-proxy-1',
- });
+ await waitFor(() =>
expect(consumerService.importConsumerGroups).toHaveBeenCalledTimes(1));
+
expect(consumerService.importConsumerGroups).toHaveBeenCalledWith('instance-proxy-1',
[
+ {
+ name: 'cg-ok',
+ subscriptionMode: 'Push',
+ consumeType: 'CLUSTERING',
+ retryMaxTimes: 16,
+ subscriptionDataType: 'NORMAL',
+ subscribedTopics: [],
+ instanceId: 'instance-proxy-1',
+ },
+ {
+ name: 'cg-fail',
+ subscriptionMode: 'Pop',
+ consumeType: 'BROADCASTING',
+ retryMaxTimes: 4,
+ subscriptionDataType: 'FIFO',
+ deliveryOrderType: 'PARTITON_ORDER',
+ subscribedTopics: [],
+ instanceId: 'instance-proxy-1',
+ },
+ ]);
+ expect(consumerService.createConsumerGroup).not.toHaveBeenCalled();
expect(await screen.findByText('已导入 1 个 Group,1 个失败')).toBeInTheDocument();
expect(screen.getByText('broker rejected group')).toBeInTheDocument();
expect(screen.getAllByText('cg-ok').length).toBeGreaterThan(0);
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index 8cc71b3c7..a5c1cf993 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -78,10 +78,11 @@ import {
batchDeleteConsumerGroups,
createConsumerGroup,
deleteConsumerGroup,
+ exportConsumerGroups,
getConsumerProgress,
getConsumerStack,
getConsumerSubscriptions,
- listAllConsumerGroups,
+ importConsumerGroups,
listConsumerGroupPage,
refreshConsumerGroup,
resetConsumerOffset,
@@ -96,7 +97,7 @@ import {
validateConsumerGroupCsvImport,
type ResourceImportRow,
} from '../../utils/resourceCsvImport';
-import { buildCsv, downloadCsv, type CsvColumn } from '../../utils/download';
+import { downloadCsv } from '../../utils/download';
import { formatLag, isLagAvailable, lagSortValue } from
'../../utils/consumerLag';
import { tableScrollX } from '../../utils/table';
@@ -140,25 +141,6 @@ const formatDelay = (totalSeconds: number): string => {
return parts.length > 0 ? parts.join('') : '0秒';
};
-const GROUP_EXPORT_COLUMNS: CsvColumn<ConsumerGroup>[] = [
- { header: 'Name', value: (group) => group.name },
- { header: 'Namespace', value: (group) => group.namespace },
- { header: 'Cluster ID', value: (group) => group.clusterId },
- { header: 'Subscription Mode', value: (group) => group.subscriptionMode },
- { header: 'Consume Type', value: (group) => group.consumeType },
- { header: 'Online Instances', value: (group) => group.onlineInstances },
- { header: 'Total Lag', value: (group) => group.totalLag },
- { header: 'Delay Seconds', value: (group) => group.delaySeconds },
- { header: 'Subscription Data Type', value: (group) =>
group.subscriptionDataType },
- { header: 'Delivery Order Type', value: (group) => group.deliveryOrderType },
- { header: 'Retry Max Times', value: (group) => group.retryMaxTimes },
- { header: 'Subscribed Topics', value: (group) => (group.subscribedTopics ??
[]).join(';') },
- { header: 'Created At', value: (group) => group.gmtCreate },
- { header: 'Updated At', value: (group) => group.gmtModified },
-];
-
-const buildConsumerGroupCsv = (groups: ConsumerGroup[]) =>
buildCsv(GROUP_EXPORT_COLUMNS, groups);
-
const visibleConsumerGroups = (groups: ConsumerGroup[], modeFilter: string):
ConsumerGroup[] => {
let data = groups;
@@ -383,16 +365,13 @@ const ConsumerPageContent = ({
const handleExport = async () => {
setExporting(true);
try {
- const allGroups = await listAllConsumerGroups({
+ const csv = await exportConsumerGroups({
instanceId: selectedInstanceId || undefined,
search: search.trim() || undefined,
+ subscriptionMode: modeFilter !== 'ALL' ? modeFilter : undefined,
});
- const exportGroups = visibleConsumerGroups(allGroups, modeFilter);
- downloadCsv(
- `rocketmq-consumer-groups-${new Date().toISOString().slice(0,
10)}.csv`,
- buildConsumerGroupCsv(exportGroups),
- );
- message.success(`已导出 ${exportGroups.length} 个 Group`);
+ downloadCsv(`rocketmq-consumer-groups-${new
Date().toISOString().slice(0, 10)}.csv`, csv);
+ message.success('Group 导出完成');
} catch {
message.error('导出 Group 失败,请稍后重试');
} finally {
@@ -582,22 +561,37 @@ const ConsumerPageContent = ({
setImporting(true);
const nextRows = importRows.map((row) => ({ ...row }));
- const createdGroups: ConsumerGroup[] = [];
+ let createdGroups: ConsumerGroup[] = [];
- for (const { row, index } of targetIndexes) {
- try {
- const created = await createConsumerGroup(row.payload);
- createdGroups.push(created);
- nextRows[index] = { ...nextRows[index], status: 'success', message:
'已创建' };
- } catch (error) {
+ try {
+ const result = await importConsumerGroups(
+ selectedInstanceId,
+ targetIndexes.map(({ row }) => row.payload),
+ );
+ createdGroups = result.groups;
+ const failureByIndex = new Map(result.failures.map((failure) =>
[failure.index, failure]));
+ targetIndexes.forEach(({ index }, requestIndex) => {
+ const failure = failureByIndex.get(requestIndex);
+ nextRows[index] = failure
+ ? {
+ ...nextRows[index],
+ status: 'failed',
+ message: failure.message || '创建失败',
+ }
+ : { ...nextRows[index], status: 'success', message: '已创建' };
+ });
+ } catch (error) {
+ for (const { index } of targetIndexes) {
nextRows[index] = {
...nextRows[index],
status: 'failed',
message: error instanceof Error ? error.message : '创建失败',
};
}
- setImportRows([...nextRows]);
+ } finally {
+ setImporting(false);
}
+ setImportRows([...nextRows]);
if (createdGroups.length > 0) {
setGroups((previous) => {
@@ -619,7 +613,6 @@ const ConsumerPageContent = ({
} else {
message.error(`${failedCount} 个 Group 导入失败`);
}
- setImporting(false);
};
const consumerGroupImportColumns:
ColumnsType<ResourceImportRow<Partial<ConsumerGroup>>> = [
@@ -1897,7 +1890,7 @@ const ConsumerPageContent = ({
<Alert
type="info"
showIcon
- message={`检测到 ${importRows.length} 个 Group,将按顺序调用创建接口`}
+ message={`检测到 ${importRows.length} 个 Group,将通过后端批量导入`}
description="仅导入可创建字段;CSV 中的 Namespace、Cluster ID 和运行状态列会被忽略。"
/>
)}
diff --git a/web/src/services/consumerService.ts
b/web/src/services/consumerService.ts
index c2b854638..24d9767a2 100644
--- a/web/src/services/consumerService.ts
+++ b/web/src/services/consumerService.ts
@@ -2,21 +2,40 @@ import { isMockMode } from './dataMode';
import * as metadataApi from '../api/metadata';
import type {
ConsumerGroup,
+ ConsumerGroupExportQuery,
ConsumerGroupPageQuery,
ConsumerGroupQuery,
ConsumerGroupSettings,
ConsumerGroupDetail,
ConsumerStackTrace,
+ ImportConsumerGroupsResult,
PageResult,
QueueProgress,
ResetConsumerOffsetRequest,
SubscriptionEntry,
} from '../api/metadata';
import { mockConsumerGroups, mockQueueProgress, mockSubscriptions } from
'../mock/consumers';
+import { buildCsv, type CsvColumn } from '../utils/download';
const consumerGroupsState = mockConsumerGroups as unknown as ConsumerGroup[];
const EXPORT_PAGE_SIZE = 100;
const MAX_EXPORT_PAGES = 100;
+const GROUP_EXPORT_COLUMNS: CsvColumn<ConsumerGroup>[] = [
+ { header: 'Name', value: (group) => group.name },
+ { header: 'Namespace', value: (group) => group.namespace },
+ { header: 'Cluster ID', value: (group) => group.clusterId },
+ { header: 'Subscription Mode', value: (group) => group.subscriptionMode },
+ { header: 'Consume Type', value: (group) => group.consumeType },
+ { header: 'Online Instances', value: (group) => group.onlineInstances },
+ { header: 'Total Lag', value: (group) => group.totalLag },
+ { header: 'Delay Seconds', value: (group) => group.delaySeconds },
+ { header: 'Subscription Data Type', value: (group) =>
group.subscriptionDataType },
+ { header: 'Delivery Order Type', value: (group) => group.deliveryOrderType },
+ { header: 'Retry Max Times', value: (group) => group.retryMaxTimes },
+ { header: 'Subscribed Topics', value: (group) => (group.subscribedTopics ??
[]).join(';') },
+ { header: 'Created At', value: (group) => group.gmtCreate },
+ { header: 'Updated At', value: (group) => group.gmtModified },
+];
function copyConsumerInstance(
instance: ConsumerGroup['instances'][number],
@@ -60,6 +79,18 @@ function filterConsumerGroups(params?: ConsumerGroupQuery):
ConsumerGroup[] {
return result;
}
+function visibleConsumerGroups(groups: ConsumerGroup[], params?:
ConsumerGroupExportQuery) {
+ let result = groups;
+ if (params?.names?.length) {
+ const selectedNames = new Set(params.names);
+ result = result.filter((group) => selectedNames.has(group.name));
+ }
+ if (params?.subscriptionMode && params.subscriptionMode !== 'ALL') {
+ result = result.filter((group) => group.subscriptionMode ===
params.subscriptionMode);
+ }
+ return [...result].sort((left, right) =>
left.name.localeCompare(right.name));
+}
+
export async function listConsumerGroups(params?: ConsumerGroupQuery):
Promise<ConsumerGroup[]> {
if (isMockMode()) {
return filterConsumerGroups(params).map(copyConsumerGroup);
@@ -216,6 +247,41 @@ export async function createConsumerGroup(data:
Partial<ConsumerGroup>): Promise
return metadataApi.createConsumerGroup(data);
}
+export async function importConsumerGroups(
+ instanceId: string,
+ groups: Partial<ConsumerGroup>[],
+): Promise<ImportConsumerGroupsResult> {
+ if (isMockMode()) {
+ const imported: ConsumerGroup[] = [];
+ const failures: ImportConsumerGroupsResult['failures'] = [];
+ for (const [index, group] of groups.entries()) {
+ try {
+ imported.push(await createConsumerGroup({ ...group, instanceId }));
+ } catch (error) {
+ failures.push({
+ index,
+ name: group.name,
+ message: error instanceof Error ? error.message : '创建失败',
+ });
+ }
+ }
+ return { imported: imported.length, failed: failures.length, groups:
imported, failures };
+ }
+ const result = await metadataApi.importConsumerGroups({ instanceId, groups
});
+ return {
+ ...result,
+ groups: result.groups.map(normalizeConsumerGroup),
+ };
+}
+
+export async function exportConsumerGroups(params: ConsumerGroupExportQuery =
{}): Promise<string> {
+ if (isMockMode()) {
+ const groups = await listAllConsumerGroups(params);
+ return buildCsv(GROUP_EXPORT_COLUMNS, visibleConsumerGroups(groups,
params));
+ }
+ return metadataApi.exportConsumerGroups(params);
+}
+
export async function deleteConsumerGroup(name: string, instanceId?: string):
Promise<void> {
if (isMockMode()) {
const idx = consumerGroupsState.findIndex((group) => group.name === name);