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 04f9eacb feat: persist instances and scope topics/groups by instance
(#951)
04f9eacb is described below
commit 04f9eacb554b3e9128d7f6274705e652761b9731
Author: lizhimins <[email protected]>
AuthorDate: Tue Aug 4 15:45:17 2026 +0800
feat: persist instances and scope topics/groups by instance (#951)
Add the rmq_instance table with five demo instances (two direct, three
proxy) behind a MyBatis-Plus repository, replacing the in-memory store.
Topics and consumer groups gain an instance_id column so instance pages
can report real per-instance counts via GROUP BY, and all five instance
sub-pages (topic/group/message/acl/dlq) share a useInstanceFilter hook
with instance-scoped routes (/instance/:instanceId/<section>). Drop the
namespace concept from topic/group UI and forms, query topic consumers
through queryTopicConsumeByWho instead of scanning every subscription
group, and tidy the message query layout (instance + query mode on one
row, no tag filter in topic mode).
---
deploy/mysql/init.sql | 1 +
deploy/mysql/upgrade-demo-instance.sql | 132 +++++++++++++
.../instance/InMemoryInstanceRepository.java | 124 ------------
.../instance/MybatisPlusInstanceRepository.java | 180 ++++++++++++++++++
.../studio/instance/group/ConsumerGroupVO.java | 1 +
.../rocketmq/studio/instance/topic/TopicVO.java | 1 +
.../studio/persistence/entity/RmqGroup.java | 2 +
.../entity/{RmqGroup.java => RmqInstance.java} | 20 +-
.../studio/persistence/entity/RmqTopic.java | 2 +
.../RmqInstanceMapper.java} | 35 +---
.../studio/rocketmq/RocketMQAdminClientImpl.java | 6 +
.../studio/rocketmq/RocketMQMetadataProvider.java | 98 +++++-----
server/src/main/resources/db/schema.sql | 158 ++++++++++++----
.../instance/InMemoryInstanceRepositoryTest.java | 44 -----
.../MybatisPlusInstanceRepositoryTest.java | 161 ++++++++++++++++
web/src/App.tsx | 5 +
web/src/api/metadata.ts | 3 +
web/src/hooks/useInstanceFilter.ts | 81 ++++++++
web/src/layouts/MainLayout.tsx | 15 +-
web/src/mock/consumers.ts | 11 ++
web/src/mock/instances.ts | 86 ++++-----
web/src/mock/topics.ts | 13 ++
.../pages/instance/__tests__/ConsumerPage.test.tsx | 8 +-
web/src/pages/instance/__tests__/DLQPage.test.tsx | 8 +-
.../pages/instance/__tests__/InstancePage.test.tsx | 5 +-
.../pages/instance/__tests__/MessagePage.test.tsx | 15 +-
.../__tests__/MessagePageAsyncState.test.tsx | 12 +-
.../pages/instance/__tests__/TopicPage.test.tsx | 54 +++++-
web/src/pages/instance/consumer.tsx | 38 ++--
web/src/pages/instance/dlq.tsx | 31 ++-
web/src/pages/instance/index.tsx | 30 ++-
web/src/pages/instance/message.tsx | 63 +++++--
web/src/pages/instance/topic.tsx | 208 +++++++++++++--------
web/src/services/instanceService.test.ts | 13 +-
34 files changed, 1178 insertions(+), 486 deletions(-)
diff --git a/deploy/mysql/init.sql b/deploy/mysql/init.sql
index c50e5727..51b52ac6 100644
--- a/deploy/mysql/init.sql
+++ b/deploy/mysql/init.sql
@@ -1,3 +1,4 @@
+SET NAMES utf8mb4;
CREATE DATABASE IF NOT EXISTS rocketmq_studio DEFAULT CHARACTER SET utf8mb4
COLLATE utf8mb4_unicode_ci;
USE rocketmq_studio;
SOURCE /docker-entrypoint-initdb.d/schema.sql;
diff --git a/deploy/mysql/upgrade-demo-instance.sql
b/deploy/mysql/upgrade-demo-instance.sql
new file mode 100644
index 00000000..0df0a17b
--- /dev/null
+++ b/deploy/mysql/upgrade-demo-instance.sql
@@ -0,0 +1,132 @@
+-- deploy/mysql/upgrade-demo-instance.sql
+-- 存量 MySQL 数据卷增量迁移(2026-08-03):实例持久化 + topic/group 增加 instance_id
+-- 适用:数据卷已初始化、docker-entrypoint-initdb.d 不会再执行的存量部署。
+-- 全新数据卷由 server/src/main/resources/db/schema.sql 直接覆盖,无需本脚本。
+-- 幂等:可重复执行。
+--
+-- 用法(远程容器内执行):
+-- docker exec -i rocketmq-studio-mysql mysql -uroot -pstudio123
rocketmq_studio < upgrade-demo-instance.sql
+
+-- 固定连接编码,防止 mysql 客户端以 latin1 解释 UTF-8 字节导致中文双重编码
+SET NAMES utf8mb4;
+
+-- 1. 实例注册表
+CREATE TABLE IF NOT EXISTS rmq_instance (
+ id VARCHAR(64) PRIMARY KEY,
+ name VARCHAR(128) NOT NULL,
+ remark VARCHAR(255),
+ type VARCHAR(32) NOT NULL COMMENT 'PROXY/DIRECT',
+ endpoint VARCHAR(512) NOT NULL,
+ created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
+ updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
+
+-- 2. rmq_topic / rmq_group 增加 instance_id 列与索引(MySQL 8.0 无 ADD COLUMN IF NOT
EXISTS,用 information_schema 判断)
+SET @sql = (SELECT IF(
+ COUNT(*) > 0, 'SELECT ''rmq_topic.instance_id already exists'' AS msg',
+ 'ALTER TABLE rmq_topic ADD COLUMN instance_id VARCHAR(64) COMMENT
''归属实例,引用 rmq_instance.id'' AFTER cluster_id, ADD INDEX idx_instance
(instance_id)')
+ FROM information_schema.COLUMNS
+ WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'rmq_topic' AND COLUMN_NAME
= 'instance_id');
+PREPARE stmt FROM @sql; EXECUTE stmt; DEALLOCATE PREPARE stmt;
+
+SET @sql = (SELECT IF(
+ COUNT(*) > 0, 'SELECT ''rmq_group.instance_id already exists'' AS msg',
+ 'ALTER TABLE rmq_group ADD COLUMN instance_id VARCHAR(64) COMMENT
''归属实例,引用 rmq_instance.id'' AFTER cluster_id, ADD INDEX idx_instance
(instance_id)')
+ FROM information_schema.COLUMNS
+ WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'rmq_group' AND COLUMN_NAME
= 'instance_id');
+PREPARE stmt FROM @sql; EXECUTE stmt; DEALLOCATE PREPARE stmt;
+
+-- 3. 默认 5 个实例(幂等)
+INSERT IGNORE INTO rmq_instance (id, name, remark, type, endpoint) VALUES
+ ('instance-direct-1', 'instance-direct-1', '直连实例 1,交易核心链路(NameServer 直连)',
'DIRECT', '10.0.1.11:9876'),
+ ('instance-direct-2', 'instance-direct-2', '直连实例 2,风控与审计链路(NameServer 直连)',
'DIRECT', '10.0.1.12:9876'),
+ ('instance-proxy-1', 'instance-proxy-1', 'Proxy 实例 1,电商交易主链路', 'PROXY',
'10.0.2.21:8080'),
+ ('instance-proxy-2', 'instance-proxy-2', 'Proxy 实例 2,营销与会员链路', 'PROXY',
'10.0.2.22:8080'),
+ ('instance-proxy-3', 'instance-proxy-3', 'Proxy 实例 3,物流与大数据链路', 'PROXY',
'10.0.2.23:8080');
+
+-- 4. 旧种子数据回填 instance_id(旧部署里这 9 个 topic、8 个 group 已存在,INSERT IGNORE 不会更新它们)
+UPDATE rmq_topic SET instance_id = 'instance-proxy-1'
+ WHERE instance_id IS NULL AND name IN (
+ 'order_create_event', 'order_status_change', 'order_timeout_cancel',
+ 'payment_result_notify', 'inventory_deduct_command');
+UPDATE rmq_topic SET instance_id = 'instance-proxy-2'
+ WHERE instance_id IS NULL AND name IN ('marketing_coupon_issue');
+UPDATE rmq_topic SET instance_id = 'instance-proxy-3'
+ WHERE instance_id IS NULL AND name IN ('logistics_tracking_update',
'settlement_daily_archive');
+UPDATE rmq_topic SET instance_id = 'instance-direct-2'
+ WHERE instance_id IS NULL AND name IN ('risk_control_audit');
+
+UPDATE rmq_group SET instance_id = 'instance-proxy-1'
+ WHERE instance_id IS NULL AND name IN (
+ 'GID_fulfillment_order', 'GID_inventory_deduct', 'GID_payment_result');
+UPDATE rmq_group SET instance_id = 'instance-proxy-2'
+ WHERE instance_id IS NULL AND name IN ('GID_marketing_coupon');
+UPDATE rmq_group SET instance_id = 'instance-proxy-3'
+ WHERE instance_id IS NULL AND name IN (
+ 'GID_logistics_tracking', 'GID_settlement_archive',
'GID_bi_realtime_report');
+UPDATE rmq_group SET instance_id = 'instance-direct-2'
+ WHERE instance_id IS NULL AND name IN ('studio-trace-consumer');
+
+-- 5. 补充新种子 topic/group(与 schema.sql 一致,已存在的行被 IGNORE 跳过)
+INSERT IGNORE INTO rmq_topic
+ (cluster_id, instance_id, name, topic_type, read_queue_nums,
write_queue_nums, perm, remark, status, created_by)
+VALUES
+ ('rocketmq-studio', 'instance-proxy-1', 'refund_apply_event',
'NORMAL', 4, 4, 6,
+ '退款申请事件,客服与财务系统订阅', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'cart_sync_event',
'NORMAL', 4, 4, 6,
+ '购物车多端同步事件', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'trade_close_archive',
'NORMAL', 2, 2, 4,
+ '交易关单归档,只读供对账回溯', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'marketing_campaign_push',
'NORMAL', 8, 8, 6,
+ '大促活动 push 触达,按人群包分批投递', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_register_event',
'NORMAL', 4, 4, 6,
+ '新会员注册事件,积分与权益系统订阅', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_points_change', 'FIFO',
4, 4, 6,
+ '会员积分变动,按会员 ID 分区保序', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_level_upgrade',
'DELAY', 4, 4, 6,
+ '会员升级权益延迟发放', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'sms_send_command',
'NORMAL', 8, 8, 6,
+ '短信下发指令,网关限流后消费', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'logistics_dispatch_order', 'FIFO',
8, 8, 6,
+ '运单调度指令,同单有序', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'bi_realtime_report',
'NORMAL', 16, 16, 6,
+ '实时报表数据流,BI 大屏消费', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'user_behavior_log',
'NORMAL', 16, 16, 6,
+ '用户行为埋点日志,离线分析入湖', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'click_stream_etl',
'NORMAL', 8, 8, 6,
+ '点击流 ETL 中间结果', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'trade_core_order_flow', 'FIFO',
8, 8, 6,
+ '交易核心订单流水,直连低延迟链路', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'payment_channel_callback',
'NORMAL', 8, 8, 6,
+ '支付渠道回调通知', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'account_ledger_entry',
'TRANSACTION', 8, 8, 6,
+ '账户记账分录,与账务落库同事务', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'ledger_reconcile_task',
'DELAY', 4, 4, 6,
+ '对账任务延迟触发,T+1 凌晨执行', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'risk_event_alert',
'NORMAL', 4, 4, 6,
+ '风控命中事件告警,实时推送处置平台', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'audit_operation_log',
'NORMAL', 8, 8, 6,
+ '操作审计日志,合规留存 180 天', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'compliance_report_daily',
'DELAY', 2, 2, 6,
+ '合规日报延迟生成任务', 'ACTIVE', 'seed');
+
+INSERT IGNORE INTO rmq_group
+ (cluster_id, instance_id, name, consume_type, message_model, max_retry,
status, created_by)
+VALUES
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_refund_process', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_cart_sync', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_trade_archive', 'PULL',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_campaign_push', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_member_points', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_member_benefit', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_sms_gateway', 'PUSH',
'CLUSTERING', 5, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_logistics_dispatch', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_behavior_ingest', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_click_stream_etl', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_trade_core_flow', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_pay_channel_cb', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_ledger_entry', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_reconcile_task', 'PULL',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'GID_risk_alert', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'GID_audit_archive', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'GID_compliance_daily', 'PULL',
'CLUSTERING', 1, 'ACTIVE', 'seed');
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepository.java
deleted file mode 100644
index ec38ba59..00000000
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepository.java
+++ /dev/null
@@ -1,124 +0,0 @@
-/*
- * 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;
-
-import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
-import org.springframework.stereotype.Component;
-
-import java.time.LocalDateTime;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Locale;
-import java.util.Map;
-import java.util.Optional;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.stream.Collectors;
-
-@Component
-public class InMemoryInstanceRepository implements InstanceRepository {
-
- private final Map<String, InstanceVO> store = new ConcurrentHashMap<>();
-
- public InMemoryInstanceRepository() {
- initData();
- }
-
- private void initData() {
- addInstance("inst-1", "production-proxy", "Production Proxy
InstanceVO",
- InstanceType.PROXY, "10.0.1.100:8080", 42, 18,
- LocalDateTime.now().minusDays(90));
- addInstance("inst-2", "staging-proxy", "Staging Proxy InstanceVO",
- InstanceType.PROXY, "10.0.2.100:8080", 15, 8,
- LocalDateTime.now().minusDays(60));
- addInstance("inst-3", "dev-direct", "Development Direct InstanceVO",
- InstanceType.DIRECT, "10.0.3.100:10911", 8, 3,
- LocalDateTime.now().minusDays(30));
- }
-
- private void addInstance(String id, String name, String remark,
InstanceType type,
- String endpoint, int topicCount, int
consumerGroupCount,
- LocalDateTime createdAt) {
- InstanceVO instance = InstanceVO.builder()
- .name(name)
- .remark(remark)
- .type(type)
- .endpoint(endpoint)
- .topicCount(topicCount)
- .consumerGroupCount(consumerGroupCount)
- .build();
- instance.setId(id);
- instance.setCreatedAt(createdAt);
- instance.setUpdatedAt(createdAt);
- store.put(id, instance);
- }
-
- @Override
- public List<InstanceVO> findAll() {
- return new ArrayList<>(store.values());
- }
-
- @Override
- public List<InstanceVO> findByType(InstanceType type) {
- return store.values().stream()
- .filter(i -> i.getType() == type)
- .collect(Collectors.toList());
- }
-
- @Override
- public List<InstanceVO> search(String keyword) {
- String lower = keyword.toLowerCase(Locale.ROOT);
- return store.values().stream()
- .filter(instance -> matchesSearch(instance, lower))
- .collect(Collectors.toList());
- }
-
- @Override
- public List<InstanceVO> findByTypeAndSearch(InstanceType type, String
keyword) {
- String lower = keyword.toLowerCase(Locale.ROOT);
- return store.values().stream()
- .filter(i -> i.getType() == type)
- .filter(instance -> matchesSearch(instance, lower))
- .collect(Collectors.toList());
- }
-
- private boolean matchesSearch(InstanceVO instance, String lowerKeyword) {
- return containsIgnoreCase(instance.getName(), lowerKeyword)
- || containsIgnoreCase(instance.getEndpoint(), lowerKeyword)
- || containsIgnoreCase(instance.getRemark(), lowerKeyword);
- }
-
- private boolean containsIgnoreCase(String value, String lowerKeyword) {
- return value != null &&
value.toLowerCase(Locale.ROOT).contains(lowerKeyword);
- }
-
- @Override
- public Optional<InstanceVO> findById(String id) {
- return Optional.ofNullable(store.get(id));
- }
-
- @Override
- public InstanceVO save(InstanceVO instance) {
- store.put(instance.getId(), instance);
- return instance;
- }
-
- @Override
- public void deleteById(String id) {
- store.remove(id);
- }
-}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
new file mode 100644
index 00000000..e2f600e7
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
@@ -0,0 +1,180 @@
+/*
+ * 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;
+
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.persistence.entity.RmqGroup;
+import org.apache.rocketmq.studio.persistence.entity.RmqInstance;
+import org.apache.rocketmq.studio.persistence.entity.RmqTopic;
+import org.apache.rocketmq.studio.persistence.mapper.RmqGroupMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqInstanceMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqTopicMapper;
+import org.springframework.stereotype.Repository;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.stream.Collectors;
+
+@Repository
+public class MybatisPlusInstanceRepository implements InstanceRepository {
+
+ private final RmqInstanceMapper instanceMapper;
+ private final RmqTopicMapper topicMapper;
+ private final RmqGroupMapper groupMapper;
+
+ public MybatisPlusInstanceRepository(RmqInstanceMapper instanceMapper,
+ RmqTopicMapper topicMapper,
+ RmqGroupMapper groupMapper) {
+ this.instanceMapper = instanceMapper;
+ this.topicMapper = topicMapper;
+ this.groupMapper = groupMapper;
+ }
+
+ @Override
+ public List<InstanceVO> findAll() {
+ return withCounts(instanceMapper.selectList(
+ new QueryWrapper<RmqInstance>().orderByAsc("id")));
+ }
+
+ @Override
+ public List<InstanceVO> findByType(InstanceType type) {
+ return withCounts(instanceMapper.selectList(
+ new QueryWrapper<RmqInstance>()
+ .eq("type", type.name())
+ .orderByAsc("id")));
+ }
+
+ @Override
+ public List<InstanceVO> search(String keyword) {
+ return withCounts(instanceMapper.selectList(
+ new QueryWrapper<RmqInstance>()
+ .and(w -> w.like("name", keyword)
+ .or().like("endpoint", keyword)
+ .or().like("remark", keyword))
+ .orderByAsc("id")));
+ }
+
+ @Override
+ public List<InstanceVO> findByTypeAndSearch(InstanceType type, String
keyword) {
+ return withCounts(instanceMapper.selectList(
+ new QueryWrapper<RmqInstance>()
+ .eq("type", type.name())
+ .and(w -> w.like("name", keyword)
+ .or().like("endpoint", keyword)
+ .or().like("remark", keyword))
+ .orderByAsc("id")));
+ }
+
+ @Override
+ public Optional<InstanceVO> findById(String id) {
+ RmqInstance entity = instanceMapper.selectById(id);
+ if (entity == null) {
+ return Optional.empty();
+ }
+ Map<String, Long> topicCounts = countByInstance(topicMapper.selectMaps(
+ new QueryWrapper<RmqTopic>().select("instance_id", "COUNT(*)
AS total")
+ .isNotNull("instance_id").eq("instance_id",
id).groupBy("instance_id")));
+ Map<String, Long> groupCounts = countByInstance(groupMapper.selectMaps(
+ new QueryWrapper<RmqGroup>().select("instance_id", "COUNT(*)
AS total")
+ .isNotNull("instance_id").eq("instance_id",
id).groupBy("instance_id")));
+ return Optional.of(toVO(entity, topicCounts, groupCounts));
+ }
+
+ @Override
+ public InstanceVO save(InstanceVO instance) {
+ RmqInstance entity = toEntity(instance);
+ if (instanceMapper.selectById(entity.getId()) != null) {
+ instanceMapper.updateById(entity);
+ } else {
+ instanceMapper.insert(entity);
+ }
+ return instance;
+ }
+
+ @Override
+ public void deleteById(String id) {
+ instanceMapper.deleteById(id);
+ }
+
+ private List<InstanceVO> withCounts(List<RmqInstance> entities) {
+ if (entities.isEmpty()) {
+ return List.of();
+ }
+ Map<String, Long> topicCounts = countByInstance(topicMapper.selectMaps(
+ new QueryWrapper<RmqTopic>().select("instance_id", "COUNT(*)
AS total")
+ .isNotNull("instance_id").groupBy("instance_id")));
+ Map<String, Long> groupCounts = countByInstance(groupMapper.selectMaps(
+ new QueryWrapper<RmqGroup>().select("instance_id", "COUNT(*)
AS total")
+ .isNotNull("instance_id").groupBy("instance_id")));
+ return entities.stream()
+ .map(entity -> toVO(entity, topicCounts, groupCounts))
+ .collect(Collectors.toList());
+ }
+
+ private Map<String, Long> countByInstance(List<Map<String, Object>> rows) {
+ Map<String, Long> counts = new HashMap<>();
+ for (Map<String, Object> row : rows) {
+ Object instanceId = row.get("instance_id");
+ Object total = row.get("total");
+ if (instanceId != null && total != null) {
+ counts.put(instanceId.toString(), ((Number)
total).longValue());
+ }
+ }
+ return counts;
+ }
+
+ private InstanceVO toVO(RmqInstance entity,
+ Map<String, Long> topicCounts,
+ Map<String, Long> groupCounts) {
+ InstanceVO vo = InstanceVO.builder()
+ .name(entity.getName())
+ .remark(entity.getRemark())
+ .type(parseType(entity.getType()))
+ .endpoint(entity.getEndpoint())
+ .topicCount(topicCounts.getOrDefault(entity.getId(),
0L).intValue())
+ .consumerGroupCount(groupCounts.getOrDefault(entity.getId(),
0L).intValue())
+ .build();
+ vo.setId(entity.getId());
+ vo.setCreatedAt(entity.getCreatedAt());
+ vo.setUpdatedAt(entity.getUpdatedAt());
+ return vo;
+ }
+
+ private InstanceType parseType(String type) {
+ try {
+ return InstanceType.valueOf(type);
+ } catch (IllegalArgumentException | NullPointerException ex) {
+ return InstanceType.PROXY;
+ }
+ }
+
+ private RmqInstance toEntity(InstanceVO vo) {
+ RmqInstance entity = new RmqInstance();
+ entity.setId(vo.getId());
+ entity.setName(vo.getName());
+ entity.setRemark(vo.getRemark());
+ entity.setType(vo.getType() == null ? null : vo.getType().name());
+ entity.setEndpoint(vo.getEndpoint());
+ entity.setCreatedAt(vo.getCreatedAt());
+ entity.setUpdatedAt(vo.getUpdatedAt());
+ return entity;
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupVO.java
index eb10ec86..f5e4a9e3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupVO.java
@@ -30,6 +30,7 @@ public class ConsumerGroupVO extends BaseEntity {
private String name;
private String namespace;
private String clusterId;
+ private String instanceId;
private SubscriptionMode subscriptionMode;
private ConsumeType consumeType;
private int onlineInstances;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicVO.java
index 79a5a14f..52a27712 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicVO.java
@@ -28,6 +28,7 @@ public class TopicVO extends BaseEntity {
private String name;
private String namespace;
private String clusterId;
+ private String instanceId;
private TopicType type;
private int writeQueues;
private int readQueues;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
index 6d6cc0a1..b97e25c7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
@@ -32,6 +32,8 @@ public class RmqGroup {
private String clusterId;
+ private String instanceId;
+
private String name;
private String consumeType;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
similarity index 80%
copy from
server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
copy to
server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
index 6d6cc0a1..7c4cdbfd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
@@ -24,25 +24,19 @@ import lombok.Data;
import java.time.LocalDateTime;
@Data
-@TableName("rmq_group")
-public class RmqGroup {
+@TableName("rmq_instance")
+public class RmqInstance {
- @TableId(type = IdType.AUTO)
- private Long id;
-
- private String clusterId;
+ @TableId(type = IdType.INPUT)
+ private String id;
private String name;
- private String consumeType;
-
- private String messageModel;
-
- private Integer maxRetry;
+ private String remark;
- private String status;
+ private String type;
- private String createdBy;
+ private String endpoint;
private LocalDateTime createdAt;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTopic.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTopic.java
index d6511d17..a99f2547 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTopic.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTopic.java
@@ -32,6 +32,8 @@ public class RmqTopic {
private String clusterId;
+ private String instanceId;
+
private String name;
private String topicType;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqInstanceMapper.java
similarity index 54%
copy from
server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
copy to
server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqInstanceMapper.java
index 6d6cc0a1..abd0cbfd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqGroup.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqInstanceMapper.java
@@ -14,37 +14,10 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.persistence.entity;
+package org.apache.rocketmq.studio.persistence.mapper;
-import com.baomidou.mybatisplus.annotation.IdType;
-import com.baomidou.mybatisplus.annotation.TableId;
-import com.baomidou.mybatisplus.annotation.TableName;
-import lombok.Data;
+import com.baomidou.mybatisplus.core.mapper.BaseMapper;
+import org.apache.rocketmq.studio.persistence.entity.RmqInstance;
-import java.time.LocalDateTime;
-
-@Data
-@TableName("rmq_group")
-public class RmqGroup {
-
- @TableId(type = IdType.AUTO)
- private Long id;
-
- private String clusterId;
-
- private String name;
-
- private String consumeType;
-
- private String messageModel;
-
- private Integer maxRetry;
-
- private String status;
-
- private String createdBy;
-
- private LocalDateTime createdAt;
-
- private LocalDateTime updatedAt;
+public interface RmqInstanceMapper extends BaseMapper<RmqInstance> {
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQAdminClientImpl.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQAdminClientImpl.java
index 69a0d039..616370e6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQAdminClientImpl.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQAdminClientImpl.java
@@ -168,6 +168,9 @@ public class RocketMQAdminClientImpl implements AdminClient
{
entity.setClusterId(clusterName);
entity.setCreatedAt(LocalDateTime.now());
}
+ if (StringUtils.hasText(topic.getInstanceId())) {
+ entity.setInstanceId(topic.getInstanceId());
+ }
entity.setTopicType(topic.getType() != null ?
topic.getType().name() : "NORMAL");
entity.setReadQueueNums(readQueues);
entity.setWriteQueueNums(writeQueues);
@@ -380,6 +383,9 @@ public class RocketMQAdminClientImpl implements AdminClient
{
entity.setClusterId(groupClusterName);
entity.setCreatedAt(LocalDateTime.now());
}
+ if (StringUtils.hasText(group.getInstanceId())) {
+ entity.setInstanceId(group.getInstanceId());
+ }
entity.setConsumeType(group.getConsumeType() != null ?
group.getConsumeType().name() : "CLUSTERING");
entity.setMessageModel(group.getSubscriptionMode() != null ?
group.getSubscriptionMode().name() : "Push");
entity.setMaxRetry(config.getRetryMaxTimes());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
index f8bcbf84..e017d746 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
@@ -22,12 +22,11 @@ import
org.apache.rocketmq.remoting.protocol.admin.ConsumeStats;
import org.apache.rocketmq.remoting.protocol.admin.OffsetWrapper;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection;
-import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
+import org.apache.rocketmq.remoting.protocol.body.GroupList;
import org.apache.rocketmq.remoting.protocol.heartbeat.SubscriptionData;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.remoting.protocol.route.QueueData;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
-import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
@@ -118,6 +117,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
vo.setId(entity.getName());
vo.setName(entity.getName());
vo.setClusterId(entity.getClusterId());
+ vo.setInstanceId(entity.getInstanceId());
vo.setType(parseTopicType(entity.getTopicType()));
vo.setReadQueues(entity.getReadQueueNums() == null ? 0 :
entity.getReadQueueNums());
vo.setWriteQueues(entity.getWriteQueueNums() == null ? 0 :
entity.getWriteQueueNums());
@@ -164,6 +164,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
vo.setId(entity.getName());
vo.setName(entity.getName());
vo.setClusterId(entity.getClusterId());
+ vo.setInstanceId(entity.getInstanceId());
vo.setConsumeType(parseConsumeType(entity.getMessageModel()));
vo.setRetryMaxTimes(entity.getMaxRetry() == null ? 0 :
entity.getMaxRetry());
vo.setCreatedAt(entity.getCreatedAt());
@@ -240,20 +241,30 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
try {
- Set<String> allGroups = collectAllConsumerGroups();
- List<TopicConsumerVO> consumers = new ArrayList<>();
+ // Ask the broker who consumes this topic instead of scanning
every subscription
+ // group, which floods the result with system groups.
+ GroupList groupList = adminExt.queryTopicConsumeByWho(name);
+ Set<String> subscribingGroups = new HashSet<>();
+ if (groupList != null && groupList.getGroupList() != null) {
+ for (String group : groupList.getGroupList()) {
+ if (!isSystemConsumerGroup(group)) {
+ subscribingGroups.add(group);
+ }
+ }
+ }
- for (String group : allGroups) {
+ List<TopicConsumerVO> consumers = new ArrayList<>();
+ for (String group : subscribingGroups) {
try {
ConsumeStats stats = adminExt.examineConsumeStats(group,
name);
- if (stats == null || stats.getOffsetTable() == null ||
stats.getOffsetTable().isEmpty()) {
- continue;
- }
-
long diffTotal = 0;
- for (Map.Entry<MessageQueue, OffsetWrapper> entry :
stats.getOffsetTable().entrySet()) {
- OffsetWrapper ow = entry.getValue();
- diffTotal += Math.max(0, ow.getBrokerOffset() -
ow.getConsumerOffset());
+ double consumeTps = 0;
+ if (stats != null && stats.getOffsetTable() != null) {
+ for (Map.Entry<MessageQueue, OffsetWrapper> entry :
stats.getOffsetTable().entrySet()) {
+ OffsetWrapper ow = entry.getValue();
+ diffTotal += Math.max(0, ow.getBrokerOffset() -
ow.getConsumerOffset());
+ }
+ consumeTps = stats.getConsumeTps();
}
ConsumeType consumeType = ConsumeType.CLUSTERING;
@@ -274,13 +285,21 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
.group(group)
.consumeType(consumeType)
.messageModel(messageModel)
- .consumeTps(stats.getConsumeTps())
+ .consumeTps(consumeTps)
.diffTotal(diffTotal)
.build());
} catch (Exception ignored) {
- // This group does not subscribe to this topic
+ // stats unavailable for this group, still list it below
without numbers
+ consumers.add(TopicConsumerVO.builder()
+ .group(group)
+ .consumeType(ConsumeType.CLUSTERING)
+ .messageModel("CLUSTERING")
+ .consumeTps(0)
+ .diffTotal(0)
+ .build());
}
}
+ consumers.sort((a, b) ->
a.getGroup().compareToIgnoreCase(b.getGroup()));
return consumers;
} catch (Exception e) {
log.warn("Failed to get consumers for topic {}: {}", name,
e.getMessage());
@@ -288,6 +307,18 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
}
+ private boolean isSystemConsumerGroup(String group) {
+ if (group == null || group.isEmpty()) {
+ return true;
+ }
+ return group.startsWith("%RETRY%")
+ || group.startsWith("%DLQ%")
+ || group.startsWith("CID_RMQ_SYS_")
+ || group.startsWith("CID_ONS_")
+ || group.startsWith("TOOLS_CONSUMER")
+ || group.startsWith("FILTERSRV_CONSUMER");
+ }
+
@Override
public List<QueueProgressVO> getGroupProgress(String name) {
if (adminExt == null) {
@@ -386,45 +417,6 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
}
- private Set<String> collectAllConsumerGroups() {
- Set<String> allGroups = new HashSet<>();
- try {
- ClusterInfo clusterInfo = adminExt.examineBrokerClusterInfo();
- if (clusterInfo == null || clusterInfo.getBrokerAddrTable() ==
null) {
- return allGroups;
- }
-
- Set<String> processedAddrs = new HashSet<>();
- for (BrokerData brokerData :
clusterInfo.getBrokerAddrTable().values()) {
- if (brokerData.getBrokerAddrs() == null) {
- continue;
- }
- // Use master address (brokerId = 0) preferentially
- String masterAddr = brokerData.getBrokerAddrs().get(0L);
- if (masterAddr == null &&
!brokerData.getBrokerAddrs().isEmpty()) {
- masterAddr =
brokerData.getBrokerAddrs().values().iterator().next();
- }
- if (masterAddr == null || processedAddrs.contains(masterAddr))
{
- continue;
- }
- processedAddrs.add(masterAddr);
-
- try {
- SubscriptionGroupWrapper wrapper =
adminExt.getAllSubscriptionGroup(masterAddr, 5000);
- if (wrapper != null && wrapper.getSubscriptionGroupTable()
!= null) {
- java.util.concurrent.ConcurrentMap<String,
SubscriptionGroupConfig> table = wrapper.getSubscriptionGroupTable();
- allGroups.addAll(table.keySet());
- }
- } catch (Exception e) {
- log.debug("Failed to get subscription groups from broker
{}: {}", masterAddr, e.getMessage());
- }
- }
- } catch (Exception e) {
- log.warn("Failed to collect consumer groups: {}", e.getMessage());
- }
- return allGroups;
- }
-
private void enrichGroupWithConnectionInfo(ConsumerGroupVO vo, String
groupName) {
try {
ConsumerConnection conn =
adminExt.examineConsumerConnectionInfo(groupName);
diff --git a/server/src/main/resources/db/schema.sql
b/server/src/main/resources/db/schema.sql
index f335b632..a918d5f9 100644
--- a/server/src/main/resources/db/schema.sql
+++ b/server/src/main/resources/db/schema.sql
@@ -3,6 +3,9 @@
-- 此文件为唯一权威 DDL 来源,MyBatis-Plus Entity 与此保持同步
-- 注意:MyBatis-Plus 不自动建表,需通过此 SQL 初始化(docker-compose 挂载执行)
+-- 固定连接编码,防止 mysql 客户端以 latin1 解释 UTF-8 字节导致中文双重编码
+SET NAMES utf8mb4;
+
-- 1. NameServer / 集群地址注册表
CREATE TABLE IF NOT EXISTS rmq_nameserver (
id VARCHAR(64) PRIMARY KEY,
@@ -15,10 +18,22 @@ CREATE TABLE IF NOT EXISTS rmq_nameserver (
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 2. Topic 管理记录(通过 Studio 创建/管理的 Topic 元数据)
+-- 2. 实例注册表(实例管理页的数据源,topic/group 按 instance_id 归属统计)
+CREATE TABLE IF NOT EXISTS rmq_instance (
+ id VARCHAR(64) PRIMARY KEY,
+ name VARCHAR(128) NOT NULL,
+ remark VARCHAR(255),
+ type VARCHAR(32) NOT NULL COMMENT 'PROXY/DIRECT',
+ endpoint VARCHAR(512) NOT NULL,
+ created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
+ updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
+
+-- 3. Topic 管理记录(通过 Studio 创建/管理的 Topic 元数据)
CREATE TABLE IF NOT EXISTS rmq_topic (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
cluster_id VARCHAR(64) NOT NULL,
+ instance_id VARCHAR(64) COMMENT '归属实例,引用 rmq_instance.id',
name VARCHAR(255) NOT NULL,
topic_type VARCHAR(32) DEFAULT 'NORMAL',
read_queue_nums INT DEFAULT 8,
@@ -29,13 +44,15 @@ CREATE TABLE IF NOT EXISTS rmq_topic (
created_by VARCHAR(64),
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
- UNIQUE KEY uk_cluster_topic (cluster_id, name)
+ UNIQUE KEY uk_cluster_topic (cluster_id, name),
+ INDEX idx_instance (instance_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 3. Consumer Group 管理记录
+-- 4. Consumer Group 管理记录
CREATE TABLE IF NOT EXISTS rmq_group (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
cluster_id VARCHAR(64) NOT NULL,
+ instance_id VARCHAR(64) COMMENT '归属实例,引用 rmq_instance.id',
name VARCHAR(255) NOT NULL,
consume_type VARCHAR(32) DEFAULT 'CONCURRENTLY',
message_model VARCHAR(32) DEFAULT 'CLUSTERING',
@@ -44,10 +61,11 @@ CREATE TABLE IF NOT EXISTS rmq_group (
created_by VARCHAR(64),
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
- UNIQUE KEY uk_cluster_group (cluster_id, name)
+ UNIQUE KEY uk_cluster_group (cluster_id, name),
+ INDEX idx_instance (instance_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 4. K8s 证书管理
+-- 5. K8s 证书管理
CREATE TABLE IF NOT EXISTS rmq_k8s_certificate (
id VARCHAR(64) PRIMARY KEY,
name VARCHAR(128) NOT NULL,
@@ -64,7 +82,7 @@ CREATE TABLE IF NOT EXISTS rmq_k8s_certificate (
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 5. 消息查询记录
+-- 6. 消息查询记录
CREATE TABLE IF NOT EXISTS rmq_message_query (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
query_type VARCHAR(32) NOT NULL COMMENT 'TOPIC/KEY/MSG_ID',
@@ -81,7 +99,7 @@ CREATE TABLE IF NOT EXISTS rmq_message_query (
INDEX idx_topic (topic)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 6. 消息轨迹查询记录
+-- 7. 消息轨迹查询记录
CREATE TABLE IF NOT EXISTS rmq_trace_query (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
msg_id VARCHAR(128) NOT NULL,
@@ -94,7 +112,7 @@ CREATE TABLE IF NOT EXISTS rmq_trace_query (
INDEX idx_queried_at (queried_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 7. 操作审计日志(所有写操作)
+-- 8. 操作审计日志(所有写操作)
CREATE TABLE IF NOT EXISTS rmq_operation_audit (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
operation VARCHAR(64) NOT NULL COMMENT
'CREATE_TOPIC/DELETE_TOPIC/CREATE_GROUP/RESET_OFFSET/SEND_MESSAGE/UPDATE_CONFIG/...',
@@ -111,14 +129,14 @@ CREATE TABLE IF NOT EXISTS rmq_operation_audit (
INDEX idx_operation (operation)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 8. 通用设置(单行)
+-- 9. 通用设置(单行)
CREATE TABLE IF NOT EXISTS rmq_settings (
id VARCHAR(16) PRIMARY KEY DEFAULT 'singleton',
json TEXT NOT NULL COMMENT 'GeneralSettingsVO JSON',
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--- 9. 数据源配置
+-- 10. 数据源配置
CREATE TABLE IF NOT EXISTS rmq_data_source (
ds_key VARCHAR(64) PRIMARY KEY,
json TEXT NOT NULL COMMENT 'DataSourceVO JSON',
@@ -127,41 +145,115 @@ CREATE TABLE IF NOT EXISTS rmq_data_source (
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- ============================================================
--- 样例数据(幂等):topic / group 列表以本库为准,创建时写库、读取时读库。
--- 以电商交易链路为例:下单 -> 支付 -> 库存 -> 履约 -> 物流 -> 结算。
+-- 样例数据(幂等):instance / topic / group 列表以本库为准,创建时写库、读取时读库。
+-- 实例管理页默认 5 个实例:2 个 DIRECT(instance-direct-1/2)+ 3 个
PROXY(instance-proxy-1/2/3)。
+-- topic/group 通过 instance_id 归属实例,实例页的 topic/group 数量按 instance_id group by
实时统计。
-- cluster_id 需与 NameServer 上报的集群名一致,否则页面按集群过滤时查不到。
-- ============================================================
+INSERT IGNORE INTO rmq_instance (id, name, remark, type, endpoint) VALUES
+ ('instance-direct-1', 'instance-direct-1', '直连实例 1,交易核心链路(NameServer 直连)',
'DIRECT', '10.0.1.11:9876'),
+ ('instance-direct-2', 'instance-direct-2', '直连实例 2,风控与审计链路(NameServer 直连)',
'DIRECT', '10.0.1.12:9876'),
+ ('instance-proxy-1', 'instance-proxy-1', 'Proxy 实例 1,电商交易主链路', 'PROXY',
'10.0.2.21:8080'),
+ ('instance-proxy-2', 'instance-proxy-2', 'Proxy 实例 2,营销与会员链路', 'PROXY',
'10.0.2.22:8080'),
+ ('instance-proxy-3', 'instance-proxy-3', 'Proxy 实例 3,物流与大数据链路', 'PROXY',
'10.0.2.23:8080');
+
INSERT IGNORE INTO rmq_topic
- (cluster_id, name, topic_type, read_queue_nums, write_queue_nums, perm,
remark, status, created_by)
+ (cluster_id, instance_id, name, topic_type, read_queue_nums,
write_queue_nums, perm, remark, status, created_by)
VALUES
- ('rocketmq-studio', 'order_create_event', 'NORMAL', 8, 8, 6,
+ -- instance-proxy-1:电商交易主链路
+ ('rocketmq-studio', 'instance-proxy-1', 'order_create_event',
'NORMAL', 8, 8, 6,
'下单成功事件,履约、营销、风控多方订阅', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'order_status_change', 'FIFO', 4, 4, 6,
+ ('rocketmq-studio', 'instance-proxy-1', 'order_status_change', 'FIFO',
4, 4, 6,
'订单状态流转,按订单号分区保证同单有序', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'order_timeout_cancel', 'DELAY', 8, 8, 6,
+ ('rocketmq-studio', 'instance-proxy-1', 'order_timeout_cancel',
'DELAY', 8, 8, 6,
'未支付订单超时关单,延迟 30 分钟投递', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'payment_result_notify', 'TRANSACTION', 8, 8, 6,
+ ('rocketmq-studio', 'instance-proxy-1', 'payment_result_notify',
'TRANSACTION', 8, 8, 6,
'支付结果通知,与支付流水落库同事务', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'inventory_deduct_command', 'NORMAL', 16, 16, 6,
+ ('rocketmq-studio', 'instance-proxy-1', 'inventory_deduct_command',
'NORMAL', 16, 16, 6,
'库存扣减指令,大促期间扩容至 16 队列', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'logistics_tracking_update', 'NORMAL', 8, 8, 6,
- '物流轨迹更新,承运商回调后投递', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'marketing_coupon_issue', 'NORMAL', 4, 4, 6,
+ ('rocketmq-studio', 'instance-proxy-1', 'refund_apply_event',
'NORMAL', 4, 4, 6,
+ '退款申请事件,客服与财务系统订阅', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'cart_sync_event',
'NORMAL', 4, 4, 6,
+ '购物车多端同步事件', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'trade_close_archive',
'NORMAL', 2, 2, 4,
+ '交易关单归档,只读供对账回溯', 'ACTIVE', 'seed'),
+ -- instance-proxy-2:营销与会员链路
+ ('rocketmq-studio', 'instance-proxy-2', 'marketing_coupon_issue',
'NORMAL', 4, 4, 6,
'营销发券,活动期间异步发放优惠券', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'settlement_daily_archive', 'NORMAL', 2, 2, 4,
+ ('rocketmq-studio', 'instance-proxy-2', 'marketing_campaign_push',
'NORMAL', 8, 8, 6,
+ '大促活动 push 触达,按人群包分批投递', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_register_event',
'NORMAL', 4, 4, 6,
+ '新会员注册事件,积分与权益系统订阅', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_points_change', 'FIFO',
4, 4, 6,
+ '会员积分变动,按会员 ID 分区保序', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'member_level_upgrade',
'DELAY', 4, 4, 6,
+ '会员升级权益延迟发放', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'sms_send_command',
'NORMAL', 8, 8, 6,
+ '短信下发指令,网关限流后消费', 'ACTIVE', 'seed'),
+ -- instance-proxy-3:物流与大数据链路
+ ('rocketmq-studio', 'instance-proxy-3', 'logistics_tracking_update',
'NORMAL', 8, 8, 6,
+ '物流轨迹更新,承运商回调后投递', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'logistics_dispatch_order', 'FIFO',
8, 8, 6,
+ '运单调度指令,同单有序', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'settlement_daily_archive',
'NORMAL', 2, 2, 4,
'日结账单归档,已停止写入仅供回溯消费', 'ACTIVE', 'seed'),
- ('rocketmq-studio', 'risk_control_audit', 'NORMAL', 4, 4, 2,
- '风控审计流水,仅生产侧写入,消费方待接入', 'ACTIVE', 'seed');
+ ('rocketmq-studio', 'instance-proxy-3', 'bi_realtime_report',
'NORMAL', 16, 16, 6,
+ '实时报表数据流,BI 大屏消费', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'user_behavior_log',
'NORMAL', 16, 16, 6,
+ '用户行为埋点日志,离线分析入湖', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'click_stream_etl',
'NORMAL', 8, 8, 6,
+ '点击流 ETL 中间结果', 'ACTIVE', 'seed'),
+ -- instance-direct-1:交易核心直连链路
+ ('rocketmq-studio', 'instance-direct-1', 'trade_core_order_flow', 'FIFO',
8, 8, 6,
+ '交易核心订单流水,直连低延迟链路', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'payment_channel_callback',
'NORMAL', 8, 8, 6,
+ '支付渠道回调通知', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'account_ledger_entry',
'TRANSACTION', 8, 8, 6,
+ '账户记账分录,与账务落库同事务', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'ledger_reconcile_task',
'DELAY', 4, 4, 6,
+ '对账任务延迟触发,T+1 凌晨执行', 'ACTIVE', 'seed'),
+ -- instance-direct-2:风控与审计链路
+ ('rocketmq-studio', 'instance-direct-2', 'risk_control_audit',
'NORMAL', 4, 4, 2,
+ '风控审计流水,仅生产侧写入,消费方待接入', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'risk_event_alert',
'NORMAL', 4, 4, 6,
+ '风控命中事件告警,实时推送处置平台', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'audit_operation_log',
'NORMAL', 8, 8, 6,
+ '操作审计日志,合规留存 180 天', 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'compliance_report_daily',
'DELAY', 2, 2, 6,
+ '合规日报延迟生成任务', 'ACTIVE', 'seed');
INSERT IGNORE INTO rmq_group
- (cluster_id, name, consume_type, message_model, max_retry, status,
created_by)
+ (cluster_id, instance_id, name, consume_type, message_model, max_retry,
status, created_by)
VALUES
- ('rocketmq-studio', 'GID_fulfillment_order', 'PUSH', 'CLUSTERING', 16,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_inventory_deduct', 'PUSH', 'CLUSTERING', 16,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_payment_result', 'PUSH', 'CLUSTERING', 16,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_logistics_tracking', 'PUSH', 'CLUSTERING', 5,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_marketing_coupon', 'PUSH', 'CLUSTERING', 3,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_settlement_archive', 'PULL', 'CLUSTERING', 3,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'GID_bi_realtime_report', 'PUSH', 'BROADCASTING', 1,
'ACTIVE', 'seed'),
- ('rocketmq-studio', 'studio-trace-consumer', 'PUSH', 'CLUSTERING', 16,
'ACTIVE', 'seed');
+ -- instance-proxy-1
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_fulfillment_order', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_inventory_deduct', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_payment_result', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_refund_process', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_cart_sync', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-1', 'GID_trade_archive', 'PULL',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ -- instance-proxy-2
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_marketing_coupon', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_campaign_push', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_member_points', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_member_benefit', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-2', 'GID_sms_gateway', 'PUSH',
'CLUSTERING', 5, 'ACTIVE', 'seed'),
+ -- instance-proxy-3
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_logistics_tracking', 'PUSH',
'CLUSTERING', 5, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_logistics_dispatch', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_settlement_archive', 'PULL',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_bi_realtime_report', 'PUSH',
'BROADCASTING', 1, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_behavior_ingest', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-proxy-3', 'GID_click_stream_etl', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ -- instance-direct-1
+ ('rocketmq-studio', 'instance-direct-1', 'GID_trade_core_flow', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_pay_channel_cb', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_ledger_entry', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-1', 'GID_reconcile_task', 'PULL',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ -- instance-direct-2
+ ('rocketmq-studio', 'instance-direct-2', 'GID_risk_alert', 'PUSH',
'CLUSTERING', 8, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'GID_audit_archive', 'PUSH',
'CLUSTERING', 3, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'GID_compliance_daily', 'PULL',
'CLUSTERING', 1, 'ACTIVE', 'seed'),
+ ('rocketmq-studio', 'instance-direct-2', 'studio-trace-consumer', 'PUSH',
'CLUSTERING', 16, 'ACTIVE', 'seed');
+
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepositoryTest.java
deleted file mode 100644
index 06a7bddb..00000000
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InMemoryInstanceRepositoryTest.java
+++ /dev/null
@@ -1,44 +0,0 @@
-/*
- * 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;
-
-import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
-import org.junit.jupiter.api.Test;
-
-import static org.assertj.core.api.Assertions.assertThat;
-
-class InMemoryInstanceRepositoryTest {
-
- private final InMemoryInstanceRepository repository = new
InMemoryInstanceRepository();
-
- @Test
- void searchShouldMatchInstanceEndpoints() {
- assertThat(repository.search("10.0.2.100:8080"))
- .extracting(InstanceVO::getId)
- .containsExactly("inst-2");
- }
-
- @Test
- void typeAndSearchShouldMatchEndpointsWithinTheSelectedType() {
- assertThat(repository.findByTypeAndSearch(InstanceType.DIRECT,
"10.0.3.100"))
- .extracting(InstanceVO::getId)
- .containsExactly("inst-3");
-
- assertThat(repository.findByTypeAndSearch(InstanceType.PROXY,
"10.0.3.100"))
- .isEmpty();
- }
-}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
new file mode 100644
index 00000000..6f71d802
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
@@ -0,0 +1,161 @@
+/*
+ * 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;
+
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.persistence.entity.RmqInstance;
+import org.apache.rocketmq.studio.persistence.mapper.RmqGroupMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqInstanceMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqTopicMapper;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.time.LocalDateTime;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+class MybatisPlusInstanceRepositoryTest {
+
+ @Mock
+ private RmqInstanceMapper instanceMapper;
+
+ @Mock
+ private RmqTopicMapper topicMapper;
+
+ @Mock
+ private RmqGroupMapper groupMapper;
+
+ @InjectMocks
+ private MybatisPlusInstanceRepository repository;
+
+ @Test
+ void findAllShouldPopulateCountsGroupedByInstance() {
+ when(instanceMapper.selectList(any(QueryWrapper.class)))
+ .thenReturn(List.of(entity("instance-direct-1",
InstanceType.DIRECT),
+ entity("instance-proxy-1", InstanceType.PROXY)));
+ when(topicMapper.selectMaps(any(QueryWrapper.class)))
+ .thenReturn(List.of(Map.of("instance_id", "instance-proxy-1",
"total", 3L)));
+ when(groupMapper.selectMaps(any(QueryWrapper.class)))
+ .thenReturn(List.of(Map.of("instance_id", "instance-proxy-1",
"total", 2L),
+ Map.of("instance_id", "instance-direct-1", "total",
1L)));
+
+ List<InstanceVO> result = repository.findAll();
+
+ assertThat(result).hasSize(2);
+ InstanceVO direct = result.stream()
+ .filter(i ->
"instance-direct-1".equals(i.getId())).findFirst().orElseThrow();
+ InstanceVO proxy = result.stream()
+ .filter(i ->
"instance-proxy-1".equals(i.getId())).findFirst().orElseThrow();
+ assertThat(direct.getTopicCount()).isZero();
+ assertThat(direct.getConsumerGroupCount()).isEqualTo(1);
+ assertThat(proxy.getTopicCount()).isEqualTo(3);
+ assertThat(proxy.getConsumerGroupCount()).isEqualTo(2);
+ }
+
+ @Test
+ void findAllShouldReturnEmptyWithoutCountQueriesWhenNoInstances() {
+
when(instanceMapper.selectList(any(QueryWrapper.class))).thenReturn(List.of());
+
+ assertThat(repository.findAll()).isEmpty();
+ verify(topicMapper, never()).selectMaps(any(QueryWrapper.class));
+ verify(groupMapper, never()).selectMaps(any(QueryWrapper.class));
+ }
+
+ @Test
+ void findByIdShouldReturnEmptyWhenMissing() {
+ when(instanceMapper.selectById("missing")).thenReturn(null);
+
+ assertThat(repository.findById("missing")).isEmpty();
+ }
+
+ @Test
+ void findByIdShouldPopulateCountsForSingleInstance() {
+
when(instanceMapper.selectById("instance-proxy-1")).thenReturn(entity("instance-proxy-1",
InstanceType.PROXY));
+ when(topicMapper.selectMaps(any(QueryWrapper.class)))
+ .thenReturn(List.of(Map.of("instance_id", "instance-proxy-1",
"total", 5L)));
+
when(groupMapper.selectMaps(any(QueryWrapper.class))).thenReturn(List.of());
+
+ Optional<InstanceVO> result = repository.findById("instance-proxy-1");
+
+ assertThat(result).isPresent();
+ assertThat(result.get().getTopicCount()).isEqualTo(5);
+ assertThat(result.get().getConsumerGroupCount()).isZero();
+ }
+
+ @Test
+ void saveShouldInsertWhenInstanceAbsent() {
+ InstanceVO vo = vo("instance-proxy-2", InstanceType.PROXY);
+ when(instanceMapper.selectById("instance-proxy-2")).thenReturn(null);
+
+ repository.save(vo);
+
+ verify(instanceMapper).insert(any(RmqInstance.class));
+ verify(instanceMapper, never()).updateById(any(RmqInstance.class));
+ }
+
+ @Test
+ void saveShouldUpdateWhenInstanceExists() {
+ InstanceVO vo = vo("instance-proxy-2", InstanceType.PROXY);
+
when(instanceMapper.selectById("instance-proxy-2")).thenReturn(entity("instance-proxy-2",
InstanceType.PROXY));
+
+ repository.save(vo);
+
+ verify(instanceMapper).updateById(any(RmqInstance.class));
+ verify(instanceMapper, never()).insert(any(RmqInstance.class));
+ }
+
+ @Test
+ void deleteByIdShouldDelegateToMapper() {
+ repository.deleteById("instance-direct-1");
+
+ verify(instanceMapper).deleteById("instance-direct-1");
+ }
+
+ private RmqInstance entity(String id, InstanceType type) {
+ RmqInstance entity = new RmqInstance();
+ entity.setId(id);
+ entity.setName(id);
+ entity.setType(type.name());
+ entity.setEndpoint("10.0.0.1:9876");
+ entity.setCreatedAt(LocalDateTime.of(2026, 8, 3, 0, 0));
+ entity.setUpdatedAt(LocalDateTime.of(2026, 8, 3, 0, 0));
+ return entity;
+ }
+
+ private InstanceVO vo(String id, InstanceType type) {
+ InstanceVO vo = InstanceVO.builder()
+ .name(id)
+ .type(type)
+ .endpoint("10.0.0.1:9876")
+ .build();
+ vo.setId(id);
+ return vo;
+ }
+}
diff --git a/web/src/App.tsx b/web/src/App.tsx
index 4e98d064..86035f3d 100644
--- a/web/src/App.tsx
+++ b/web/src/App.tsx
@@ -145,10 +145,15 @@ function App() {
<Route index element={<HomePage />} />
<Route path="instance" element={<InstancePage />} />
<Route path="instance/topic" element={<TopicPage />} />
+ <Route path="instance/:instanceId/topic" element={<TopicPage />} />
<Route path="instance/consumer" element={<ConsumerPage />} />
+ <Route path="instance/:instanceId/consumer" element={<ConsumerPage
/>} />
<Route path="instance/message" element={<MessagePage />} />
+ <Route path="instance/:instanceId/message" element={<MessagePage
/>} />
<Route path="instance/acl" element={<AclPage />} />
+ <Route path="instance/:instanceId/acl" element={<AclPage />} />
<Route path="instance/dlq" element={<DlqPage />} />
+ <Route path="instance/:instanceId/dlq" element={<DlqPage />} />
<Route path="cluster" element={<ClusterPage />} />
<Route path="cluster/certs" element={<K8sCertsPage />} />
<Route path="cluster/clients" element={<ClientsPage />} />
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index bb85ec7b..f3897b18 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -11,6 +11,7 @@ export interface Topic {
namespace: string;
type: string;
clusterId: string;
+ instanceId?: string;
writeQueues: number;
readQueues: number;
perm: string;
@@ -24,6 +25,7 @@ export interface Topic {
export interface TopicQuery {
clusterId?: string;
+ instanceId?: string;
type?: string;
search?: string;
}
@@ -49,6 +51,7 @@ export interface ConsumerGroup {
name: string;
namespace: string;
clusterId: string;
+ instanceId?: string;
subscriptionMode: string;
consumeType: string;
onlineInstances: number;
diff --git a/web/src/hooks/useInstanceFilter.ts
b/web/src/hooks/useInstanceFilter.ts
new file mode 100644
index 00000000..31e7a33a
--- /dev/null
+++ b/web/src/hooks/useInstanceFilter.ts
@@ -0,0 +1,81 @@
+/*
+ * 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.
+ */
+
+import { useEffect, useState } from 'react';
+import { useLocation, useNavigate } from 'react-router-dom';
+import { listInstances } from '../services/instanceService';
+import type { Instance } from '../api/instance';
+
+const INSTANCE_SCOPED_PATH =
/^\/instance\/([^/]+)\/(topic|consumer|message|acl|dlq)$/;
+const STATIC_SECTION_PATH = /^\/instance\/(topic|consumer|message|acl|dlq)$/;
+
+/**
+ * 实例维度页面的公共筛选逻辑:从 /instance/:instanceId/<section> 路由解析当前实例,
+ * 无实例参数时重定向到第一个实例;实例列表加载失败时降级为不过滤。
+ */
+export function useInstanceFilter() {
+ const navigate = useNavigate();
+ const { pathname } = useLocation();
+
+ const scopedMatch = pathname.match(INSTANCE_SCOPED_PATH);
+ const staticMatch = pathname.match(STATIC_SECTION_PATH);
+ const routeInstanceId = scopedMatch?.[1];
+ const section = scopedMatch?.[2] ?? staticMatch?.[1] ?? 'topic';
+
+ const [instances, setInstances] = useState<Instance[]>([]);
+
+ useEffect(() => {
+ let cancelled = false;
+ void listInstances()
+ .then((nextInstances) => {
+ if (cancelled) return;
+ setInstances(nextInstances);
+ if (!routeInstanceId && nextInstances.length > 0) {
+ navigate(`/instance/${nextInstances[0].id}/${section}`, { replace:
true });
+ }
+ })
+ .catch(() => {
+ // 实例列表加载失败时不做实例过滤,保持页面数据可用
+ });
+ return () => {
+ cancelled = true;
+ };
+ }, [navigate, routeInstanceId, section]);
+
+ const selectedInstanceId =
+ routeInstanceId && instances.some((instance) => instance.id ===
routeInstanceId)
+ ? routeInstanceId
+ : (instances[0]?.id ?? '');
+ const selectedInstance = instances.find((instance) => instance.id ===
selectedInstanceId);
+
+ const selectInstance = (id: string) => {
+ navigate(`/instance/${id}/${section}`);
+ };
+
+ const instanceOptions = instances.map((instance) => ({
+ value: instance.id,
+ label: instance.name,
+ }));
+
+ return {
+ instances,
+ selectedInstanceId,
+ selectedInstance,
+ selectInstance,
+ instanceOptions,
+ };
+}
diff --git a/web/src/layouts/MainLayout.tsx b/web/src/layouts/MainLayout.tsx
index cb8e0f2d..045c9d6f 100644
--- a/web/src/layouts/MainLayout.tsx
+++ b/web/src/layouts/MainLayout.tsx
@@ -97,6 +97,13 @@ const MainLayout = () => {
return () => window.removeEventListener('keydown', openSearchWithShortcut);
}, []);
+ const instanceScopedMatch = location.pathname.match(
+ /^\/instance\/[^/]+\/(topic|consumer|message|acl|dlq)$/,
+ );
+ const selectedMenuKey = instanceScopedMatch
+ ? `/instance/${instanceScopedMatch[1]}`
+ : location.pathname;
+
const menuItems = [
{ key: '/', icon: <House size={iconSize} weight="duotone" />, label:
t('nav.home') },
{
@@ -166,8 +173,12 @@ const MainLayout = () => {
},
...pathSnippets.map((_, index) => {
const path = '/' + pathSnippets.slice(0, index + 1).join('/');
+ const isSectionLeaf = instanceScopedMatch && index ===
pathSnippets.length - 1;
+ const leafTitle = isSectionLeaf
+ ? breadcrumbMap[`/instance/${instanceScopedMatch[1]}`]
+ : undefined;
return {
- title: breadcrumbMap[path] || path,
+ title: breadcrumbMap[path] || leafTitle || path,
key: path,
};
}),
@@ -252,7 +263,7 @@ const MainLayout = () => {
<Menu
theme={darkMode ? 'dark' : 'light'}
mode="inline"
- selectedKeys={[location.pathname]}
+ selectedKeys={[selectedMenuKey]}
defaultOpenKeys={['instance-group', 'cluster-ops-group']}
items={menuItems}
onClick={({ key }) => navigate(key)}
diff --git a/web/src/mock/consumers.ts b/web/src/mock/consumers.ts
index 0d805085..9cf87540 100644
--- a/web/src/mock/consumers.ts
+++ b/web/src/mock/consumers.ts
@@ -29,6 +29,7 @@ export interface ConsumerInstance {
export interface ConsumerGroup {
name: string;
namespace: string;
+ instanceId: string;
clusterId: string;
subscriptionMode: 'Push' | 'Pop';
consumeType: 'CLUSTERING' | 'BROADCASTING';
@@ -65,6 +66,7 @@ export interface SubscriptionEntry {
export const mockConsumerGroups: ConsumerGroup[] = [
{
name: 'cg-order-notify',
+ instanceId: 'instance-proxy-1',
namespace: 'trade',
clusterId: 'hz-prod',
subscriptionMode: 'Push',
@@ -130,6 +132,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-payment-callback',
+ instanceId: 'instance-proxy-1',
namespace: 'trade',
clusterId: 'hz-prod',
subscriptionMode: 'Push',
@@ -179,6 +182,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-user-activity',
+ instanceId: 'instance-proxy-3',
namespace: 'user',
clusterId: 'hz-prod',
subscriptionMode: 'Push',
@@ -220,6 +224,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-inventory-sync',
+ instanceId: 'instance-proxy-1',
namespace: 'supply',
clusterId: 'sh-prod',
subscriptionMode: 'Push',
@@ -254,6 +259,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-log-collector',
+ instanceId: 'instance-direct-2',
namespace: 'infra',
clusterId: 'hz-prod',
subscriptionMode: 'Pop',
@@ -335,6 +341,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-notification-push',
+ instanceId: 'instance-proxy-2',
namespace: 'message',
clusterId: 'hz-prod',
subscriptionMode: 'Push',
@@ -376,6 +383,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-ai-task-worker',
+ instanceId: 'instance-proxy-3',
namespace: 'ai',
clusterId: 'hz-prod',
subscriptionMode: 'Pop',
@@ -425,6 +433,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-metrics-aggregator',
+ instanceId: 'instance-direct-1',
namespace: 'infra',
clusterId: 'sh-prod',
subscriptionMode: 'Push',
@@ -458,6 +467,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-risk-control',
+ instanceId: 'instance-direct-2',
namespace: 'risk',
clusterId: 'hz-prod',
subscriptionMode: 'Push',
@@ -515,6 +525,7 @@ export const mockConsumerGroups: ConsumerGroup[] = [
},
{
name: 'cg-data-sync',
+ instanceId: 'instance-proxy-2',
namespace: 'data',
clusterId: 'sh-prod',
subscriptionMode: 'Push',
diff --git a/web/src/mock/instances.ts b/web/src/mock/instances.ts
index bdf45f93..f9c994b1 100644
--- a/web/src/mock/instances.ts
+++ b/web/src/mock/instances.ts
@@ -19,58 +19,58 @@ import type { Instance } from '../api/instance';
export const mockInstances: Instance[] = [
{
- id: '1',
- name: 'rocketmq-trade',
- remark: '核心交易链路,承载订单、支付等主要业务',
- type: 'PROXY',
- endpoint: 'proxy-hz.rocketmq.internal:8080',
- topicCount: 128,
- consumerGroupCount: 56,
- createdAt: '2024-03-15 08:30:00',
- updatedAt: '2026-06-20 14:15:00',
+ id: 'instance-direct-1',
+ name: 'instance-direct-1',
+ remark: '直连实例 1,交易核心链路(NameServer 直连)',
+ type: 'DIRECT',
+ endpoint: '10.0.1.11:9876',
+ topicCount: 4,
+ consumerGroupCount: 4,
+ createdAt: '2026-08-03 10:00:00',
+ updatedAt: '2026-08-03 10:00:00',
},
{
- id: '2',
- name: 'rocketmq-dr',
- remark: '灾备集群,与 trade 集群互为双活',
- type: 'PROXY',
- endpoint: 'proxy-sh.rocketmq.internal:8080',
- topicCount: 96,
- consumerGroupCount: 42,
- createdAt: '2024-05-10 10:00:00',
- updatedAt: '2026-06-18 09:30:00',
+ id: 'instance-direct-2',
+ name: 'instance-direct-2',
+ remark: '直连实例 2,风控与审计链路(NameServer 直连)',
+ type: 'DIRECT',
+ endpoint: '10.0.1.12:9876',
+ topicCount: 4,
+ consumerGroupCount: 4,
+ createdAt: '2026-08-03 10:00:00',
+ updatedAt: '2026-08-03 10:00:00',
},
{
- id: '3',
- name: 'rocketmq-debug',
- remark: '开发测试环境,仅供内部调试使用',
+ id: 'instance-proxy-1',
+ name: 'instance-proxy-1',
+ remark: 'Proxy 实例 1,电商交易主链路',
type: 'PROXY',
- endpoint: 'localhost:8081',
- topicCount: 15,
- consumerGroupCount: 8,
- createdAt: '2025-01-20 14:00:00',
- updatedAt: '2026-07-01 11:45:00',
+ endpoint: '10.0.2.21:8080',
+ topicCount: 8,
+ consumerGroupCount: 6,
+ createdAt: '2026-08-03 10:00:00',
+ updatedAt: '2026-08-03 10:00:00',
},
{
- id: '4',
- name: 'rocketmq-legacy',
- remark: '旧版集群,计划 Q3 完成迁移后下线',
- type: 'DIRECT',
- endpoint: 'namesrv-legacy:9876',
- topicCount: 64,
- consumerGroupCount: 30,
- createdAt: '2022-08-01 09:00:00',
- updatedAt: '2025-12-10 16:20:00',
+ id: 'instance-proxy-2',
+ name: 'instance-proxy-2',
+ remark: 'Proxy 实例 2,营销与会员链路',
+ type: 'PROXY',
+ endpoint: '10.0.2.22:8080',
+ topicCount: 6,
+ consumerGroupCount: 5,
+ createdAt: '2026-08-03 10:00:00',
+ updatedAt: '2026-08-03 10:00:00',
},
{
- id: '5',
- name: 'rocketmq-staging',
- remark: '预发布验证环境,与生产配置一致',
+ id: 'instance-proxy-3',
+ name: 'instance-proxy-3',
+ remark: 'Proxy 实例 3,物流与大数据链路',
type: 'PROXY',
- endpoint: 'proxy-staging:8080',
- topicCount: 32,
- consumerGroupCount: 18,
- createdAt: '2024-11-05 11:30:00',
- updatedAt: '2026-06-25 08:00:00',
+ endpoint: '10.0.2.23:8080',
+ topicCount: 6,
+ consumerGroupCount: 6,
+ createdAt: '2026-08-03 10:00:00',
+ updatedAt: '2026-08-03 10:00:00',
},
];
diff --git a/web/src/mock/topics.ts b/web/src/mock/topics.ts
index 949097ce..69aa7741 100644
--- a/web/src/mock/topics.ts
+++ b/web/src/mock/topics.ts
@@ -18,6 +18,7 @@
export interface Topic {
name: string;
namespace: string;
+ instanceId: string;
type: 'NORMAL' | 'FIFO' | 'DELAY' | 'TRANSACTION' | 'LITE';
clusterId: string;
writeQueues: number;
@@ -35,6 +36,7 @@ export const topics: Topic[] = [
// NORMAL topics (4)
{
name: 'order-create',
+ instanceId: 'instance-proxy-1',
namespace: 'trade',
type: 'NORMAL',
clusterId: 'rmq-cn-v5-prod-01',
@@ -50,6 +52,7 @@ export const topics: Topic[] = [
},
{
name: 'user-activity-log',
+ instanceId: 'instance-proxy-3',
namespace: 'user',
type: 'NORMAL',
clusterId: 'rmq-cn-v5-prod-01',
@@ -65,6 +68,7 @@ export const topics: Topic[] = [
},
{
name: 'system-log',
+ instanceId: 'instance-direct-2',
namespace: 'message',
type: 'NORMAL',
clusterId: 'rmq-cn-v4-prod-02',
@@ -80,6 +84,7 @@ export const topics: Topic[] = [
},
{
name: 'notification-email',
+ instanceId: 'instance-proxy-2',
namespace: 'message',
type: 'NORMAL',
clusterId: 'rmq-cn-v5-prod-01',
@@ -97,6 +102,7 @@ export const topics: Topic[] = [
// FIFO topics (2)
{
name: 'inventory-sync',
+ instanceId: 'instance-proxy-1',
namespace: 'supply',
type: 'FIFO',
clusterId: 'rmq-cn-v5-prod-01',
@@ -112,6 +118,7 @@ export const topics: Topic[] = [
},
{
name: 'payment-sequence',
+ instanceId: 'instance-proxy-1',
namespace: 'trade',
type: 'FIFO',
clusterId: 'rmq-cn-v5-prod-01',
@@ -129,6 +136,7 @@ export const topics: Topic[] = [
// DELAY topics (2)
{
name: 'notification-push',
+ instanceId: 'instance-proxy-2',
namespace: 'message',
type: 'DELAY',
clusterId: 'rmq-cn-v5-prod-01',
@@ -144,6 +152,7 @@ export const topics: Topic[] = [
},
{
name: 'scheduled-task',
+ instanceId: 'instance-direct-1',
namespace: 'supply',
type: 'DELAY',
clusterId: 'rmq-cn-v4-prod-02',
@@ -161,6 +170,7 @@ export const topics: Topic[] = [
// TRANSACTION topics (2)
{
name: 'payment-callback',
+ instanceId: 'instance-direct-1',
namespace: 'trade',
type: 'TRANSACTION',
clusterId: 'rmq-cn-v5-prod-01',
@@ -176,6 +186,7 @@ export const topics: Topic[] = [
},
{
name: 'order-confirm',
+ instanceId: 'instance-proxy-1',
namespace: 'trade',
type: 'TRANSACTION',
clusterId: 'rmq-cn-v5-prod-01',
@@ -193,6 +204,7 @@ export const topics: Topic[] = [
// LITE topics (2)
{
name: 'chat-session',
+ instanceId: 'instance-proxy-2',
namespace: 'ai',
type: 'LITE',
clusterId: 'rmq-cn-v5-prod-01',
@@ -208,6 +220,7 @@ export const topics: Topic[] = [
},
{
name: 'ai-task-dispatch',
+ instanceId: 'instance-proxy-3',
namespace: 'ai',
type: 'LITE',
clusterId: 'rmq-cn-v5-prod-01',
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 8f5b6d7c..a7ebdb28 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -19,6 +19,7 @@ import { App } from 'antd';
import { render, screen, waitFor } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import type React from 'react';
+import { MemoryRouter } from 'react-router-dom';
import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import type { ConsumerGroup } from '../../../api/metadata';
import { LangProvider } from '../../../i18n/LangContext';
@@ -35,6 +36,9 @@ vi.mock('../../../services/consumerService', () => ({
listConsumerGroups: vi.fn(),
resetConsumerOffset: vi.fn(),
}));
+vi.mock('../../../services/instanceService', () => ({
+ listInstances: vi.fn().mockResolvedValue([]),
+}));
beforeAll(() => {
Object.defineProperty(window, 'matchMedia', {
@@ -72,7 +76,9 @@ const group: ConsumerGroup = {
const renderWithProviders = (ui: React.ReactElement) =>
render(
<App>
- <LangProvider>{ui}</LangProvider>
+ <LangProvider>
+ <MemoryRouter>{ui}</MemoryRouter>
+ </LangProvider>
</App>,
);
diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx
b/web/src/pages/instance/__tests__/DLQPage.test.tsx
index 3d64bb17..bd5d2b82 100644
--- a/web/src/pages/instance/__tests__/DLQPage.test.tsx
+++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx
@@ -19,6 +19,7 @@ import { App } from 'antd';
import { render, screen, waitFor, within } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import type React from 'react';
+import { MemoryRouter } from 'react-router-dom';
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from
'vitest';
import type { DLQGroup } from '../../../api/message';
import { LangProvider } from '../../../i18n/LangContext';
@@ -29,6 +30,9 @@ vi.mock('../../../services/messageService', () => ({
listDLQGroups: vi.fn(),
resendDLQ: vi.fn(),
}));
+vi.mock('../../../services/instanceService', () => ({
+ listInstances: vi.fn().mockResolvedValue([]),
+}));
const dlqGroup: DLQGroup = {
groupName: 'cg-order',
@@ -51,7 +55,9 @@ const secondDlqGroup: DLQGroup = {
const renderWithProviders = (ui: React.ReactElement) =>
render(
<App>
- <LangProvider>{ui}</LangProvider>
+ <LangProvider>
+ <MemoryRouter>{ui}</MemoryRouter>
+ </LangProvider>
</App>,
);
diff --git a/web/src/pages/instance/__tests__/InstancePage.test.tsx
b/web/src/pages/instance/__tests__/InstancePage.test.tsx
index 269679e9..48d81511 100644
--- a/web/src/pages/instance/__tests__/InstancePage.test.tsx
+++ b/web/src/pages/instance/__tests__/InstancePage.test.tsx
@@ -18,6 +18,7 @@
import { App } from 'antd';
import { act, fireEvent, render, screen, waitFor, within } from
'@testing-library/react';
import userEvent from '@testing-library/user-event';
+import { MemoryRouter } from 'react-router-dom';
import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import type { Instance } from '../../../api/instance';
import { LangProvider } from '../../../i18n/LangContext';
@@ -63,7 +64,9 @@ const renderPage = () =>
render(
<App>
<LangProvider>
- <InstancePage />
+ <MemoryRouter>
+ <InstancePage />
+ </MemoryRouter>
</LangProvider>
</App>,
);
diff --git a/web/src/pages/instance/__tests__/MessagePage.test.tsx
b/web/src/pages/instance/__tests__/MessagePage.test.tsx
index b1257174..630ae17b 100644
--- a/web/src/pages/instance/__tests__/MessagePage.test.tsx
+++ b/web/src/pages/instance/__tests__/MessagePage.test.tsx
@@ -19,6 +19,7 @@ import { App } from 'antd';
import { render, screen, waitFor } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import type React from 'react';
+import { MemoryRouter } from 'react-router-dom';
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from
'vitest';
import { LangProvider } from '../../../i18n/LangContext';
@@ -31,6 +32,13 @@ const QUERY_HISTORY_STORAGE_KEY =
'rocketmq-studio-message-query-history';
vi.mock('../../../services/messageService', () => messageServiceMocks);
+vi.mock('../../../services/instanceService', () => ({
+ listInstances: vi.fn().mockResolvedValue([]),
+}));
+vi.mock('../../../services/topicService', () => ({
+ listTopics: vi.fn().mockResolvedValue([]),
+}));
+
import MessagePage from '../message';
beforeAll(() => {
@@ -52,7 +60,9 @@ beforeAll(() => {
const renderWithProviders = (ui: React.ReactElement) =>
render(
<App>
- <LangProvider>{ui}</LangProvider>
+ <LangProvider>
+ <MemoryRouter>{ui}</MemoryRouter>
+ </LangProvider>
</App>,
);
@@ -182,10 +192,9 @@ describe('Message page query history', () => {
renderWithProviders(<MessagePage />);
await user.click(screen.getByRole('button', { name: /最近查询/ }));
- await user.click(await screen.findByText('Topic: order-create · Tag:
vip'));
+ await user.click(await screen.findByText('Topic: order-create'));
await waitFor(() => {
expect(messageServiceMocks.queryMessages).toHaveBeenLastCalledWith(topicParams);
- expect(screen.getByPlaceholderText('输入 Tag(可选)')).toHaveValue('vip');
});
await user.click(screen.getByRole('button', { name: /最近查询/ }));
diff --git a/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
b/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
index fe34ff8a..2f938b4b 100644
--- a/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
+++ b/web/src/pages/instance/__tests__/MessagePageAsyncState.test.tsx
@@ -16,6 +16,7 @@
*/
import { App, ConfigProvider, message } from 'antd';
+import { MemoryRouter } from 'react-router-dom';
import { act, render, screen, waitFor, within } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from
'vitest';
@@ -30,6 +31,13 @@ const serviceMocks = vi.hoisted(() => ({
vi.mock('../../../services/messageService', () => serviceMocks);
+vi.mock('../../../services/instanceService', () => ({
+ listInstances: vi.fn().mockResolvedValue([]),
+}));
+vi.mock('../../../services/topicService', () => ({
+ listTopics: vi.fn().mockResolvedValue([]),
+}));
+
beforeAll(() => {
Object.defineProperty(window, 'matchMedia', {
writable: true,
@@ -87,7 +95,9 @@ const renderPage = () =>
<ConfigProvider theme={{ token: { motion: false } }}>
<App>
<LangProvider>
- <MessagePage />
+ <MemoryRouter>
+ <MessagePage />
+ </MemoryRouter>
</LangProvider>
</App>
</ConfigProvider>,
diff --git a/web/src/pages/instance/__tests__/TopicPage.test.tsx
b/web/src/pages/instance/__tests__/TopicPage.test.tsx
index f50a4db1..e0fc385b 100644
--- a/web/src/pages/instance/__tests__/TopicPage.test.tsx
+++ b/web/src/pages/instance/__tests__/TopicPage.test.tsx
@@ -18,6 +18,7 @@
import { describe, it, expect, vi, beforeAll, beforeEach, afterEach } from
'vitest';
import { render, screen, waitFor, within } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
+import { MemoryRouter, Route, Routes } from 'react-router-dom';
import { App } from 'antd';
import { LangProvider } from '../../../i18n/LangContext';
import type { Topic } from '../../../api/metadata';
@@ -33,7 +34,12 @@ const topicServiceMocks = vi.hoisted(() => ({
sendTopicMessage: vi.fn(),
}));
+const instanceServiceMocks = vi.hoisted(() => ({
+ listInstances: vi.fn(),
+}));
+
vi.mock('../../../services/topicService', () => topicServiceMocks);
+vi.mock('../../../services/instanceService', () => instanceServiceMocks);
beforeAll(() => {
Object.defineProperty(window, 'matchMedia', {
@@ -71,11 +77,16 @@ const buildTopics = (count: number): Topic[] =>
};
});
-const renderWithProviders = () =>
+const renderWithProviders = (initialEntry = '/instance/topic') =>
render(
<App>
<LangProvider>
- <TopicPage />
+ <MemoryRouter initialEntries={[initialEntry]}>
+ <Routes>
+ <Route path="/instance/topic" element={<TopicPage />} />
+ <Route path="/instance/:instanceId/topic" element={<TopicPage />}
/>
+ </Routes>
+ </MemoryRouter>
</LangProvider>
</App>,
);
@@ -92,6 +103,7 @@ describe('TopicPage', () => {
topicServiceMocks.batchDeleteTopics.mockResolvedValue({ deleted: [],
failed: [] });
topicServiceMocks.getTopicRoutes.mockResolvedValue([]);
topicServiceMocks.getTopicConsumers.mockResolvedValue([]);
+ instanceServiceMocks.listInstances.mockResolvedValue([]);
});
afterEach(() => {
@@ -144,4 +156,42 @@ describe('TopicPage', () => {
expect(screen.getByRole('button', { name: /删除 \(1\)$/
})).toBeInTheDocument();
expect(screen.getByText('已删除 2 个 Topic,1 个删除失败')).toBeInTheDocument();
});
+
+ it('filters topics by the instance from the route and shows its endpoint',
async () => {
+ const base = buildTopics(1)[0];
+ topicServiceMocks.listTopics.mockResolvedValue([
+ { ...base, name: 'topic-a', instanceId: 'instance-proxy-1' },
+ { ...base, name: 'topic-b', instanceId: 'instance-proxy-2' },
+ ]);
+ instanceServiceMocks.listInstances.mockResolvedValue([
+ {
+ id: 'instance-proxy-1',
+ name: 'instance-proxy-1',
+ remark: '',
+ type: 'PROXY',
+ endpoint: '10.0.2.21:8080',
+ topicCount: 1,
+ consumerGroupCount: 0,
+ createdAt: '2026-01-01T00:00:00Z',
+ updatedAt: '2026-01-01T00:00:00Z',
+ },
+ {
+ id: 'instance-proxy-2',
+ name: 'instance-proxy-2',
+ remark: '',
+ type: 'PROXY',
+ endpoint: '10.0.2.22:8080',
+ topicCount: 1,
+ consumerGroupCount: 0,
+ createdAt: '2026-01-01T00:00:00Z',
+ updatedAt: '2026-01-01T00:00:00Z',
+ },
+ ]);
+
+ renderWithProviders('/instance/instance-proxy-1/topic');
+
+ expect(await screen.findByText('topic-a')).toBeInTheDocument();
+ expect(screen.queryByText('topic-b')).not.toBeInTheDocument();
+ expect(screen.getByText('10.0.2.21:8080')).toBeInTheDocument();
+ });
});
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index 9b6491c1..72f6d936 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -77,6 +77,7 @@ import {
listConsumerGroups,
resetConsumerOffset,
} from '../../services/consumerService';
+import { useInstanceFilter } from '../../hooks/useInstanceFilter';
const { Text } = Typography;
@@ -131,6 +132,7 @@ const isInconsistentSubscription = (subscription:
SubscriptionEntry): boolean =>
═══════════════════════════════════════════ */
const ConsumerPage = () => {
const { t } = useLang();
+ const { selectedInstanceId, selectInstance, instanceOptions } =
useInstanceFilter();
const [groups, setGroups] = useState<ConsumerGroup[]>([]);
const [loading, setLoading] = useState(true);
const [submitting, setSubmitting] = useState(false);
@@ -213,12 +215,13 @@ const ConsumerPage = () => {
/* ─── Filtered & sorted data ─── */
const filtered = useMemo(() => {
let data = groups.filter(
- (g) =>
- g.name.includes(search) ||
- g.namespace.includes(search) ||
- g.subscribedTopics.some((t) => t.includes(search)),
+ (g) => g.name.includes(search) || g.subscribedTopics.some((t) =>
t.includes(search)),
);
+ if (selectedInstanceId) {
+ data = data.filter((g) => g.instanceId === selectedInstanceId);
+ }
+
if (modeFilter !== 'ALL') {
data = data.filter((g) => g.subscriptionMode === modeFilter);
}
@@ -230,7 +233,7 @@ const ConsumerPage = () => {
}
return data;
- }, [groups, search, modeFilter, sortKey]);
+ }, [groups, search, modeFilter, sortKey, selectedInstanceId]);
/* ─── Open detail modal ─── */
const openModal = (group: ConsumerGroup) => {
@@ -589,8 +592,16 @@ const ConsumerPage = () => {
{/* ─── Filter Bar ─── */}
<Flex justify="space-between" align="center" style={{ marginBottom: 16
}}>
<Space size={12} wrap>
+ <Select
+ placeholder="选择实例"
+ value={selectedInstanceId || undefined}
+ onChange={selectInstance}
+ options={instanceOptions}
+ style={{ width: 220 }}
+ notFoundContent="暂无实例"
+ />
<Input.Search
- placeholder="搜索 Group 名称、命名空间或 Topic"
+ placeholder="搜索 Group 名称或 Topic"
allowClear
value={search}
onChange={(e) => setSearch(e.target.value)}
@@ -708,9 +719,6 @@ const ConsumerPage = () => {
<Space>
<Cube size={18} weight="fill" color="#1677ff" />
<span style={{ fontWeight: 600 }}>{selectedGroup.name}</span>
- <Tag color="default" style={{ fontSize: 11, borderRadius: 4 }}>
- {selectedGroup.namespace}
- </Tag>
</Space>
) : (
'Group 详情'
@@ -807,9 +815,6 @@ const ConsumerPage = () => {
<Descriptions.Item label="Group 名称">
<Text strong>{selectedGroup.name}</Text>
</Descriptions.Item>
- <Descriptions.Item label="命名空间">
- <Tag color="default">{selectedGroup.namespace}</Tag>
- </Descriptions.Item>
<Descriptions.Item label="所属集群">
{selectedGroup.clusterId}
</Descriptions.Item>
@@ -1030,7 +1035,7 @@ const ConsumerPage = () => {
.then((values) => {
Modal.confirm({
title: '确认创建',
- content: `将创建消费组 "${values.name}",命名空间: ${values.namespace ||
'default'}`,
+ content: `将创建消费组 "${values.name}"`,
okText: '确认创建',
cancelText: '取消',
onOk: async () => {
@@ -1038,13 +1043,13 @@ const ConsumerPage = () => {
try {
const created = await createConsumerGroup({
name: values.name,
- namespace: values.namespace || 'default',
subscriptionMode: values.subscriptionMode,
consumeType: values.consumeType,
retryMaxTimes: values.retryMaxTimes,
subscriptionDataType: values.dataType || 'NORMAL',
deliveryOrderType: values.deliveryOrderType,
subscribedTopics: [],
+ ...(selectedInstanceId ? { instanceId:
selectedInstanceId } : {}),
});
setGroups((prev) => [
created,
@@ -1076,7 +1081,6 @@ const ConsumerPage = () => {
layout="vertical"
style={{ marginTop: 16 }}
initialValues={{
- namespace: 'default',
subscriptionMode: 'Push',
consumeType: 'CLUSTERING',
retryMaxTimes: 16,
@@ -1096,10 +1100,6 @@ const ConsumerPage = () => {
<Input placeholder="例:cg-order-notify" />
</Form.Item>
- <Form.Item label="命名空间" name="namespace">
- <Input placeholder="例:trade" />
- </Form.Item>
-
<Form.Item label="订阅模式" name="subscriptionMode">
<Radio.Group>
<Radio.Button value="Push">Push</Radio.Button>
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index 5ee6d892..5cc88bb5 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -26,6 +26,7 @@ import {
Modal,
DatePicker,
Typography,
+ Select,
message,
} from 'antd';
import { MagnifyingGlass, Eye, ArrowsCounterClockwise, Download } from
'@phosphor-icons/react';
@@ -36,6 +37,7 @@ import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import type { DLQGroup } from '../../api/message';
import { listDLQGroups, resendDLQ } from '../../services/messageService';
+import { useInstanceFilter } from '../../hooks/useInstanceFilter';
const { Text } = Typography;
const { RangePicker } = DatePicker;
@@ -80,6 +82,7 @@ const exportDLQGroups = (groups: DLQGroup[], filename:
string) => {
═══════════════════════════════════════════ */
const DLQPage = () => {
const { t } = useLang();
+ const { selectedInstanceId, selectInstance, instanceOptions } =
useInstanceFilter();
const [groups, setGroups] = useState<DLQGroup[]>([]);
const [loading, setLoading] = useState(true);
const [refreshKey, setRefreshKey] = useState(0);
@@ -280,15 +283,25 @@ const DLQPage = () => {
{/* ── Filter Bar ── */}
<Flex justify="space-between" align="center" style={{ marginBottom: 16
}}>
- <Input.Search
- placeholder="搜索 Group 名称或 DLQ Topic"
- allowClear
- value={search}
- onChange={(e) => setSearch(e.target.value)}
- onSearch={setSearch}
- style={{ width: 320 }}
- prefix={<MagnifyingGlass size={14} color="#9CA3AF" />}
- />
+ <Space size={12} wrap>
+ <Select
+ placeholder="选择实例"
+ value={selectedInstanceId || undefined}
+ onChange={selectInstance}
+ options={instanceOptions}
+ style={{ width: 220 }}
+ notFoundContent="暂无实例"
+ />
+ <Input.Search
+ placeholder="搜索 Group 名称或 DLQ Topic"
+ allowClear
+ value={search}
+ onChange={(e) => setSearch(e.target.value)}
+ onSearch={setSearch}
+ style={{ width: 320 }}
+ prefix={<MagnifyingGlass size={14} color="#9CA3AF" />}
+ />
+ </Space>
<Button
icon={<Download size={16} />}
disabled={selectedGroups.length === 0}
diff --git a/web/src/pages/instance/index.tsx b/web/src/pages/instance/index.tsx
index 5ef21347..aec233b2 100644
--- a/web/src/pages/instance/index.tsx
+++ b/web/src/pages/instance/index.tsx
@@ -16,6 +16,7 @@
*/
import { useCallback, useEffect, useRef, useState } from 'react';
+import { useNavigate } from 'react-router-dom';
import {
Table,
Card,
@@ -28,6 +29,7 @@ import {
Form,
Flex,
Typography,
+ Alert,
message,
} from 'antd';
import { useLang } from '../../i18n/LangContext';
@@ -58,6 +60,7 @@ type InstanceTypeFilter = 'ALL' | Instance['type'];
═══════════════════════════════════════════ */
const InstancePage = () => {
const { t } = useLang();
+ const navigate = useNavigate();
const [instances, setInstances] = useState<Instance[]>([]);
const [loading, setLoading] = useState(true);
const [search, setSearch] = useState('');
@@ -65,6 +68,7 @@ const InstancePage = () => {
const [typeFilter, setTypeFilter] = useState<InstanceTypeFilter>('ALL');
const [addModalOpen, setAddModalOpen] = useState(false);
const [addForm] = Form.useForm();
+ const addInstanceType = Form.useWatch<'PROXY' | 'DIRECT' |
undefined>('type', addForm);
const [editModalOpen, setEditModalOpen] = useState(false);
const [editingInstance, setEditingInstance] = useState<Instance |
null>(null);
const [editForm] = Form.useForm();
@@ -235,7 +239,7 @@ const InstancePage = () => {
key: 'actions',
width: 160,
render: (_: unknown, record: Instance) => (
- <Flex gap={6}>
+ <Flex gap={6} onClick={(e) => e.stopPropagation()}>
<Button
size="small"
icon={<EditOutlined />}
@@ -327,7 +331,7 @@ const InstancePage = () => {
size="small"
onRow={(record) => ({
style: { cursor: 'pointer' },
- onClick: () => message.info(`进入 ${record.name}`),
+ onClick: () => navigate(`/instance/${record.id}/topic`),
})}
/>
</Card>
@@ -371,9 +375,29 @@ const InstancePage = () => {
label="接入地址"
name="endpoint"
rules={[{ required: true, message: '请输入接入地址' }]}
+ extra={
+ addInstanceType === 'DIRECT'
+ ? 'Direct 模式请填写 NameServer SLB 地址(K8s 场景下一般为 NameServer
Service 地址,如 namesrv.mq.svc:9876)'
+ : addInstanceType === 'PROXY'
+ ? 'Proxy 模式请填写 Proxy SLB 内网地址(如 proxy.mq.svc:8080)'
+ : '请先选择接入方式'
+ }
>
- <Input placeholder="例:proxy.example.com:8080" />
+ <Input
+ placeholder={
+ addInstanceType === 'DIRECT'
+ ? '例:namesrv.mq.svc.cluster.local:9876'
+ : '例:proxy.mq.svc.cluster.local:8080'
+ }
+ />
</Form.Item>
+ <Alert
+ type="info"
+ showIcon
+ style={{ marginBottom: 16 }}
+ message="接入地址为客户端访问入口"
+ description="接入地址会展示在 Topic 等页面供客户端配置使用。若客户端环境无法解析该地址(如 K8s 内部
Service 域名),可自行配置 DNS 解析或在客户端 hosts 中映射。"
+ />
<Form.Item label="备注" name="remark">
<Input.TextArea rows={2} placeholder="可选,描述实例用途" />
</Form.Item>
diff --git a/web/src/pages/instance/message.tsx
b/web/src/pages/instance/message.tsx
index cc583feb..7497f0ce 100644
--- a/web/src/pages/instance/message.tsx
+++ b/web/src/pages/instance/message.tsx
@@ -54,6 +54,8 @@ import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import type { MessageQuery, MessageRecord, TraceRecord } from
'../../api/message';
import { getMessageTrace, queryMessages } from '../../services/messageService';
+import { listTopics } from '../../services/topicService';
+import { useInstanceFilter } from '../../hooks/useInstanceFilter';
const { Paragraph, Text } = Typography;
const { RangePicker } = DatePicker;
@@ -172,7 +174,7 @@ const queryLabel = ({ mode, params }: RecentQuery): string
=> {
if (mode === 'key') {
return `Key: ${params.key || '全部'}${params.topic ? ` · Topic:
${params.topic}` : ''}`;
}
- return `Topic: ${params.topic || '全部'}${params.tag ? ` · Tag: ${params.tag}`
: ''}`;
+ return `Topic: ${params.topic || '全部'}`;
};
/* ═══════════════════════════════════════════
@@ -180,10 +182,31 @@ const queryLabel = ({ mode, params }: RecentQuery):
string => {
═══════════════════════════════════════════ */
const MessagePage = () => {
const { t } = useLang();
+ const { selectedInstanceId, selectInstance, instanceOptions } =
useInstanceFilter();
+ const [topicOptions, setTopicOptions] = useState<string[]>(TOPIC_OPTIONS);
+
+ useEffect(() => {
+ let cancelled = false;
+ void listTopics()
+ .then((nextTopics) => {
+ if (cancelled) return;
+ const scoped = selectedInstanceId
+ ? nextTopics.filter((topic) => topic.instanceId ===
selectedInstanceId)
+ : nextTopics;
+ if (scoped.length > 0) {
+ setTopicOptions(scoped.map((topic) => topic.name));
+ }
+ })
+ .catch(() => {
+ // 加载失败保持静态选项可用
+ });
+ return () => {
+ cancelled = true;
+ };
+ }, [selectedInstanceId]);
const [queryMode, setQueryMode] = useState<QueryMode>('topic');
const [selectedTopic, setSelectedTopic] = useState<string | undefined>();
const [dateRange, setDateRange] = useState<[Dayjs, Dayjs]>(getDefaultRange);
- const [tagInput, setTagInput] = useState('');
const [keyInput, setKeyInput] = useState('');
const [msgIdInput, setMsgIdInput] = useState('');
const [messages, setMessages] = useState<MessageRecord[]>([]);
@@ -209,7 +232,6 @@ const MessagePage = () => {
const handleReset = () => {
queryGenerationRef.current += 1;
setSelectedTopic(undefined);
- setTagInput('');
setKeyInput('');
setMsgIdInput('');
setDateRange(getDefaultRange());
@@ -260,7 +282,6 @@ const MessagePage = () => {
queryMode === 'topic'
? {
topic: selectedTopic,
- tag: tagInput || undefined,
startTime: dateRange[0].valueOf(),
endTime: dateRange[1].valueOf(),
}
@@ -275,7 +296,6 @@ const MessagePage = () => {
const { mode, params } = recentQuery;
setQueryMode(mode);
setSelectedTopic(params.topic);
- setTagInput(params.tag || '');
setKeyInput(params.key || '');
setMsgIdInput(params.msgId || '');
if (mode === 'topic' && params.startTime !== undefined && params.endTime
!== undefined) {
@@ -626,16 +646,26 @@ const MessagePage = () => {
═══════════════════════════════════════════ */
return (
<div style={{ padding: 24 }}>
- <PageHeader title={t('message.title')} subtitle="按 Topic、Tag、Key 或
Message ID 检索消息" />
+ <PageHeader title={t('message.title')} subtitle="按 Topic、Key 或 Message
ID 检索消息" />
{/* ── Query Form ── */}
<Card style={{ marginBottom: 16 }}>
<Space direction="vertical" size={16} style={{ width: '100%' }}>
- <Segmented
- options={QUERY_OPTIONS}
- value={queryMode}
- onChange={(v) => setQueryMode(v as QueryMode)}
- />
+ <Space size={12}>
+ <Select
+ placeholder="选择实例"
+ value={selectedInstanceId || undefined}
+ onChange={selectInstance}
+ options={instanceOptions}
+ style={{ width: 220 }}
+ notFoundContent="暂无实例"
+ />
+ <Segmented
+ options={QUERY_OPTIONS}
+ value={queryMode}
+ onChange={(v) => setQueryMode(v as QueryMode)}
+ />
+ </Space>
<Space wrap size={12}>
{queryMode === 'topic' && (
@@ -647,7 +677,7 @@ const MessagePage = () => {
onChange={setSelectedTopic}
allowClear
showSearch
- options={TOPIC_OPTIONS.map((t) => ({
+ options={topicOptions.map((t) => ({
value: t,
label: t,
}))}
@@ -662,13 +692,6 @@ const MessagePage = () => {
}
}}
/>
- <Input
- placeholder="输入 Tag(可选)"
- style={{ width: 180 }}
- value={tagInput}
- onChange={(e) => setTagInput(e.target.value)}
- allowClear
- />
</>
)}
@@ -681,7 +704,7 @@ const MessagePage = () => {
onChange={setSelectedTopic}
allowClear
showSearch
- options={TOPIC_OPTIONS.map((t) => ({
+ options={topicOptions.map((t) => ({
value: t,
label: t,
}))}
diff --git a/web/src/pages/instance/topic.tsx b/web/src/pages/instance/topic.tsx
index b8ae793a..cb71d4af 100644
--- a/web/src/pages/instance/topic.tsx
+++ b/web/src/pages/instance/topic.tsx
@@ -65,6 +65,7 @@ import {
listTopics,
sendTopicMessage,
} from '../../services/topicService';
+import { useInstanceFilter } from '../../hooks/useInstanceFilter';
const { Text } = Typography;
@@ -74,16 +75,6 @@ const CLUSTER_NAME_MAP: Record<string, { name: string; type:
string }> = {
'rmq-cn-v4-prod-02': { name: 'rmq-cn-v4-prod-02', type: 'V4_DIRECT' },
};
-// ─── Namespace options ────────────────────────────────────────────
-const NAMESPACE_OPTIONS = [
- { label: '全部', value: '' },
- { label: 'trade', value: 'trade' },
- { label: 'user', value: 'user' },
- { label: 'message', value: 'message' },
- { label: 'supply', value: 'supply' },
- { label: 'ai', value: 'ai' },
-];
-
const TYPE_OPTIONS = [
{ label: '全部', value: '' },
{ label: '普通', value: 'NORMAL' },
@@ -205,6 +196,20 @@ const RANDOM_BODY_GENERATORS = [
// ─── Format helpers ───────────────────────────────────────────────
const formatNumber = (n: number) => n.toLocaleString('zh-CN');
+
+// 解析批量粘贴的用户属性串:key=value 按换行或逗号分隔,等号只取第一个
+const parsePropsText = (text: string): Record<string, string> => {
+ const props: Record<string, string> = {};
+ for (const line of text.split(/[\n,]+/)) {
+ const trimmed = line.trim();
+ if (!trimmed) continue;
+ const eqIndex = trimmed.indexOf('=');
+ if (eqIndex <= 0) continue;
+ const key = trimmed.slice(0, eqIndex).trim();
+ if (key) props[key] = trimmed.slice(eqIndex + 1).trim();
+ }
+ return props;
+};
const formatDateTime = (iso: string): string => {
const d = new Date(iso);
const pad = (n: number) => String(n).padStart(2, '0');
@@ -214,6 +219,8 @@ const formatDateTime = (iso: string): string => {
// ═══════════════════════════════════════════════════════════════════
const TopicPage = () => {
const { t } = useLang();
+ const { selectedInstanceId, selectedInstance, selectInstance,
instanceOptions } =
+ useInstanceFilter();
// ─── State ─────────────────────────────────────────────────────
const [topics, setTopics] = useState<Topic[]>([]);
@@ -223,7 +230,6 @@ const TopicPage = () => {
const [selectedRowKeys, setSelectedRowKeys] = useState<React.Key[]>([]);
const [searchText, setSearchText] = useState('');
const [typeFilter, setTypeFilter] = useState('');
- const [nsFilter, setNsFilter] = useState('');
const [tablePage, setTablePage] = useState(1);
const [tablePageSize, setTablePageSize] = useState(20);
const [viewMode, setViewMode] = useState<string>('列表');
@@ -237,6 +243,7 @@ const TopicPage = () => {
const [sendTopic, setSendTopic] = useState<Topic | null>(null);
const [sending, setSending] = useState(false);
const [sendForm] = Form.useForm();
+ const [propsMode, setPropsMode] = useState<'form' | 'text'>('form');
const { modal } = App.useApp();
useEffect(() => {
@@ -262,13 +269,13 @@ const TopicPage = () => {
() =>
topics
.filter((t) => {
+ if (selectedInstanceId && t.instanceId !== selectedInstanceId)
return false;
if (searchText &&
!t.name.toLowerCase().includes(searchText.toLowerCase())) return false;
if (typeFilter && t.type !== typeFilter) return false;
- if (nsFilter && t.namespace !== nsFilter) return false;
return true;
})
.sort((a, b) => a.name.localeCompare(b.name)),
- [topics, searchText, typeFilter, nsFilter],
+ [topics, selectedInstanceId, searchText, typeFilter],
);
const maxTablePage = Math.max(1, Math.ceil(filteredTopics.length /
tablePageSize));
@@ -328,6 +335,7 @@ const TopicPage = () => {
void openDetail(topic);
} else if (key === 'send') {
setSendTopic(topic);
+ setPropsMode('form');
sendForm.setFieldsValue({ topic: topic.name, tag: '', key: '', body: '',
properties: [] });
setSendModalOpen(true);
} else if (key === 'delete') {
@@ -501,9 +509,6 @@ const TopicPage = () => {
{typeInfo?.labelKey ? t(typeInfo.labelKey) : topic.type}
</Tag>
</Descriptions.Item>
- <Descriptions.Item label="命名空间">
- <Tag>{topic.namespace}</Tag>
- </Descriptions.Item>
<Descriptions.Item label="集群" span={2}>
<Space>
<Text>{topic.clusterId}</Text>
@@ -552,9 +557,8 @@ const TopicPage = () => {
</Tag>
</Flex>
- {/* Namespace + cluster tags */}
+ {/* Cluster tags */}
<Space size={4} style={{ marginBottom: 16 }}>
- <Tag style={{ fontSize: 11 }}>{topic.namespace}</Tag>
{clusterType && (
<Tag color={clusterType.color} style={{ fontSize: 11 }}>
{t(clusterType.labelKey)}
@@ -600,7 +604,10 @@ const TopicPage = () => {
const handleCreate = async () => {
try {
const values = await form.validateFields();
- const created = await createTopic(values);
+ const created = await createTopic({
+ ...values,
+ ...(selectedInstanceId ? { instanceId: selectedInstanceId } : {}),
+ });
setTopics((previous) => [created, ...previous]);
message.success(`Topic「${created.name}」创建成功`);
setModalOpen(false);
@@ -612,12 +619,20 @@ const TopicPage = () => {
// ─── Send message modal submit ────────────────────────────────
const handleSend = async () => {
+ let values;
try {
- const values = await sendForm.validateFields();
- setSending(true);
- // Build properties from the list
- const props: Record<string, string> = {};
- if (values.properties && Array.isArray(values.properties)) {
+ values = await sendForm.validateFields();
+ } catch {
+ // validation error, keep the modal open
+ return;
+ }
+ setSending(true);
+ try {
+ // Build properties: batch-paste text mode or key-value form rows
+ let props: Record<string, string> = {};
+ if (propsMode === 'text') {
+ props = parsePropsText(values.propsText || '');
+ } else if (values.properties && Array.isArray(values.properties)) {
values.properties.forEach((p: { key?: string; value?: string }) => {
if (p.key) props[p.key] = p.value || '';
});
@@ -629,11 +644,10 @@ const TopicPage = () => {
body: values.body,
properties: props,
});
+ // Keep the modal open for consecutive sends
message.success(`消息发送成功!MsgId: ${result.msgId}`);
- setSendModalOpen(false);
- sendForm.resetFields();
} catch {
- // validation error, do nothing
+ message.error('消息发送失败,请稍后重试');
} finally {
setSending(false);
}
@@ -647,6 +661,26 @@ const TopicPage = () => {
{/* ── Header ────────────────────────────────────────────── */}
<PageHeader title={t('topic.title')} subtitle={`共
${filteredTopics.length} 个 Topic`} />
+ {/* ── Endpoint hint ─────────────────────────────────────── */}
+ {selectedInstance && (
+ <div style={{ marginBottom: 16, fontSize: 13, lineHeight: 1.8 }}>
+ <Space wrap size={8}>
+ <Text strong>
+ 当前实例:{selectedInstance.name}(
+ {selectedInstance.type === 'DIRECT' ? 'Direct 模式' : 'Proxy 模式'})
+ </Text>
+ <Text code copyable>
+ {selectedInstance.endpoint}
+ </Text>
+ </Space>
+ <div style={{ color: '#8c8c8c' }}>
+ {selectedInstance.type === 'DIRECT'
+ ? '接入点为 NameServer SLB 地址(K8s 场景下一般为 NameServer Service
地址),Direct 模式客户端通过该地址发现 Broker。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。'
+ : '接入点为 Proxy SLB 内网地址,gRPC/Remoting
客户端直接连接该地址收发消息。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。'}
+ </div>
+ </div>
+ )}
+
{/* ── Filter bar ────────────────────────────────────────── */}
<Flex
gap={12}
@@ -656,6 +690,17 @@ const TopicPage = () => {
justify="space-between"
>
<Space size={12} wrap>
+ <Select
+ placeholder="选择实例"
+ value={selectedInstanceId || undefined}
+ onChange={(value) => {
+ resetTablePage();
+ selectInstance(value);
+ }}
+ options={instanceOptions}
+ style={{ width: 220 }}
+ notFoundContent="暂无实例"
+ />
<Input.Search
placeholder="搜索 Topic 名称"
allowClear
@@ -681,16 +726,6 @@ const TopicPage = () => {
options={TYPE_OPTIONS}
style={{ width: 140 }}
/>
- <Select
- placeholder="命名空间"
- value={nsFilter}
- onChange={(value) => {
- setNsFilter(value);
- resetTablePage();
- }}
- options={NAMESPACE_OPTIONS}
- style={{ width: 140 }}
- />
<Segmented
value={viewMode}
onChange={(v) => setViewMode(v as string)}
@@ -881,7 +916,6 @@ const TopicPage = () => {
readQueues: 8,
perm: 'RW',
type: 'NORMAL',
- namespace: 'default',
}}
style={{ marginTop: 16 }}
>
@@ -899,18 +933,9 @@ const TopicPage = () => {
<Input placeholder="请输入 Topic 名称" />
</Form.Item>
- <Row gutter={16}>
- <Col span={12}>
- <Form.Item label="命名空间" name="namespace" rules={[{ required:
true }]}>
- <Select disabled options={[{ label: 'default', value:
'default' }]} />
- </Form.Item>
- </Col>
- <Col span={12}>
- <Form.Item label="类型" name="type" rules={[{ required: true }]}>
- <Select options={TYPE_OPTIONS.filter((o) => o.value)} />
- </Form.Item>
- </Col>
- </Row>
+ <Form.Item label="类型" name="type" rules={[{ required: true }]}>
+ <Select options={TYPE_OPTIONS.filter((o) => o.value)} />
+ </Form.Item>
<Row gutter={16}>
<Col span={12}>
@@ -1026,35 +1051,62 @@ const TopicPage = () => {
自定义属性(可选)
</Divider>
- <Form.List name="properties">
- {(fields, { add, remove }) => (
- <>
- {fields.map(({ key, name, ...rest }) => (
- <Row gutter={8} key={key} align="middle" style={{
marginBottom: 8 }}>
- <Col span={10}>
- <Form.Item {...rest} name={[name, 'key']} style={{
marginBottom: 0 }}>
- <Input placeholder="属性名" />
- </Form.Item>
- </Col>
- <Col span={10}>
- <Form.Item {...rest} name={[name, 'value']} style={{
marginBottom: 0 }}>
- <Input placeholder="属性值" />
- </Form.Item>
- </Col>
- <Col span={4}>
- <MinusCircleOutlined
- style={{ color: '#ff4d4f', fontSize: 18, cursor:
'pointer' }}
- onClick={() => remove(name)}
- />
- </Col>
- </Row>
- ))}
- <Button type="dashed" onClick={() => add()} block
icon={<PlusCircleOutlined />}>
- 添加属性
- </Button>
- </>
+ <Flex justify="space-between" align="center" style={{ marginBottom:
12 }}>
+ <Segmented
+ size="small"
+ value={propsMode}
+ onChange={(value) => setPropsMode(value as 'form' | 'text')}
+ options={[
+ { label: '逐条录入', value: 'form' },
+ { label: '批量粘贴', value: 'text' },
+ ]}
+ />
+ {propsMode === 'text' && (
+ <Text type="secondary" style={{ fontSize: 12 }}>
+ 支持 key=value,多个属性用换行或逗号分隔
+ </Text>
)}
- </Form.List>
+ </Flex>
+
+ {propsMode === 'text' ? (
+ <Form.Item name="propsText" style={{ marginBottom: 0 }}>
+ <Input.TextArea
+ rows={5}
+ placeholder={'TAGS=tagA\nKEY1=value1, KEY2=value2'}
+ style={{ fontFamily: 'monospace', fontSize: 13 }}
+ />
+ </Form.Item>
+ ) : (
+ <Form.List name="properties">
+ {(fields, { add, remove }) => (
+ <>
+ {fields.map(({ key, name, ...rest }) => (
+ <Row gutter={8} key={key} align="middle" style={{
marginBottom: 8 }}>
+ <Col span={10}>
+ <Form.Item {...rest} name={[name, 'key']} style={{
marginBottom: 0 }}>
+ <Input placeholder="属性名" />
+ </Form.Item>
+ </Col>
+ <Col span={10}>
+ <Form.Item {...rest} name={[name, 'value']} style={{
marginBottom: 0 }}>
+ <Input placeholder="属性值" />
+ </Form.Item>
+ </Col>
+ <Col span={4}>
+ <MinusCircleOutlined
+ style={{ color: '#ff4d4f', fontSize: 18, cursor:
'pointer' }}
+ onClick={() => remove(name)}
+ />
+ </Col>
+ </Row>
+ ))}
+ <Button type="dashed" onClick={() => add()} block
icon={<PlusCircleOutlined />}>
+ 添加属性
+ </Button>
+ </>
+ )}
+ </Form.List>
+ )}
</Form>
</Modal>
</div>
diff --git a/web/src/services/instanceService.test.ts
b/web/src/services/instanceService.test.ts
index b003f31d..5d285f05 100644
--- a/web/src/services/instanceService.test.ts
+++ b/web/src/services/instanceService.test.ts
@@ -42,13 +42,16 @@ describe('instanceService mock instances', () => {
it('filters mock instances with the same type and search semantics as the
API', async () => {
const byType = await listInstances({ type: 'DIRECT' });
- expect(byType.map((instance) => instance.id)).toEqual(['4']);
+ expect(byType.map((instance) => instance.id)).toEqual([
+ 'instance-direct-1',
+ 'instance-direct-2',
+ ]);
- const byEndpoint = await listInstances({ search: ' PROXY-HZ ' });
- expect(byEndpoint.map((instance) => instance.id)).toEqual(['1']);
+ const byEndpoint = await listInstances({ search: ' 10.0.2.21 ' });
+ expect(byEndpoint.map((instance) =>
instance.id)).toEqual(['instance-proxy-1']);
- const combined = await listInstances({ type: 'DIRECT', search:
'namesrv-legacy' });
- expect(combined.map((instance) => instance.id)).toEqual(['4']);
+ const combined = await listInstances({ type: 'DIRECT', search:
'instance-direct-2' });
+ expect(combined.map((instance) =>
instance.id)).toEqual(['instance-direct-2']);
});
it('does not expose created or updated store records by reference', async ()
=> {