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 5f4e01b5 fix(consumer): paginate topic consumer diagnostics (#1677)
5f4e01b5 is described below
commit 5f4e01b50dd66b85dfc26b7dc0f9dd7668df4e44
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 11 20:48:33 2026 +0800
fix(consumer): paginate topic consumer diagnostics (#1677)
---
.../studio/instance/topic/MetadataService.java | 11 +++++++
.../studio/instance/topic/TopicConsumerPageVO.java | 31 +++++++++++++++++++
.../studio/instance/topic/TopicController.java | 9 ++++++
.../rocketmq/studio/provider/InstanceProvider.java | 14 +++++++++
.../provider/apache/ApacheInstanceProvider.java | 6 ++++
.../studio/provider/apache/MetadataProvider.java | 14 +++++++++
.../provider/apache/RocketMQMetadataProvider.java | 30 ++++++++++++++-----
.../studio/instance/topic/MetadataServiceTest.java | 11 +++++++
.../studio/instance/topic/TopicControllerTest.java | 16 ++++++++++
.../apache/RocketMQMetadataProviderTest.java | 28 +++++++++++++++--
web/src/api/metadata.test.ts | 15 ++++++++++
web/src/api/metadata.ts | 20 +++++++++++++
.../pages/instance/__tests__/TopicPage.test.tsx | 13 ++++++++
web/src/pages/instance/topic.tsx | 35 +++++++++++++++++-----
web/src/services/topicService.test.ts | 15 +++++++++-
web/src/services/topicService.ts | 22 ++++++++++++++
16 files changed, 272 insertions(+), 18 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index 31b70cb7..eaf9fd81 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -103,6 +103,17 @@ public class MetadataService {
return resolve(instanceId).getTopicConsumers(instanceId, topicName);
}
+ public TopicConsumerPageVO getTopicConsumersPage(String instanceId, String
name, int page, int pageSize) {
+ String topicName = requireName(name, "topic name");
+ if (page < 1) {
+ throw new BusinessException(400, "page must be greater than zero");
+ }
+ if (pageSize < 1 || pageSize > 100) {
+ throw new BusinessException(400, "pageSize must be between 1 and
100");
+ }
+ return resolve(instanceId).getTopicConsumersPage(instanceId,
topicName, page, pageSize);
+ }
+
public SendMessageVO sendMessage(SendMessageDTO request) {
requireSendMessageRequest(request);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicConsumerPageVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicConsumerPageVO.java
new file mode 100644
index 00000000..7afe57a9
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicConsumerPageVO.java
@@ -0,0 +1,31 @@
+/*
+ * 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.topic;
+
+import lombok.Builder;
+import lombok.Value;
+
+import java.util.List;
+
+@Value
+@Builder
+public class TopicConsumerPageVO {
+ List<TopicConsumerVO> items;
+ int total;
+ int page;
+ int pageSize;
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
index 6caeda88..262633a3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
@@ -79,6 +79,15 @@ public class TopicController {
return Result.ok(metadataService.getTopicConsumers(instanceId, name));
}
+ @GetMapping("/{name}/consumers/page")
+ public Result<TopicConsumerPageVO> getTopicConsumersPage(
+ @PathVariable String name,
+ @RequestParam(required = false) String instanceId,
+ @RequestParam(defaultValue = "1") int page,
+ @RequestParam(defaultValue = "20") int pageSize) {
+ return Result.ok(metadataService.getTopicConsumersPage(instanceId,
name, page, pageSize));
+ }
+
@PostMapping("/send")
public Result<SendMessageVO> sendMessage(@Valid @RequestBody(required =
false) SendMessageDTO request) {
requireSendMessageRequest(request);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
index e3f2fa19..4582a4a6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
@@ -23,6 +23,7 @@ import
org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import java.util.List;
@@ -50,6 +51,19 @@ public interface InstanceProvider {
List<TopicConsumerVO> getTopicConsumers(String instanceId, String
topicName);
+ default TopicConsumerPageVO getTopicConsumersPage(String instanceId,
String topicName, int page, int pageSize) {
+ List<TopicConsumerVO> consumers = getTopicConsumers(instanceId,
topicName);
+ int total = consumers.size();
+ int from = Math.min((page - 1) * pageSize, total);
+ int to = Math.min(from + pageSize, total);
+ return TopicConsumerPageVO.builder()
+ .items(consumers.subList(from, to))
+ .total(total)
+ .page(page)
+ .pageSize(pageSize)
+ .build();
+ }
+
List<ConsumerGroupVO> listConsumerGroups(String instanceId, String search);
ConsumerGroupVO createConsumerGroup(String instanceId, ConsumerGroupVO
group);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index 045f1564..3bf81839 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -25,6 +25,7 @@ import
org.apache.rocketmq.studio.instance.message.MessageProvider;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
import org.springframework.stereotype.Component;
@@ -88,6 +89,11 @@ public class ApacheInstanceProvider implements
InstanceProvider {
return metadataProvider.getTopicConsumers(instanceId, topicName);
}
+ @Override
+ public TopicConsumerPageVO getTopicConsumersPage(String instanceId, String
topicName, int page, int pageSize) {
+ return metadataProvider.getTopicConsumersPage(instanceId, topicName,
page, pageSize);
+ }
+
@Override
public List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
search) {
return metadataProvider.listConsumerGroups(instanceId, null, search);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
index 7d243478..f66b7fb1 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.provider.apache;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.BrokerRouteVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
@@ -35,6 +36,19 @@ public interface MetadataProvider {
List<BrokerRouteVO> getTopicRoutes(String instanceId, String name);
List<TopicConsumerVO> getTopicConsumers(String instanceId, String name);
+
+ default TopicConsumerPageVO getTopicConsumersPage(String instanceId,
String name, int page, int pageSize) {
+ List<TopicConsumerVO> consumers = getTopicConsumers(instanceId, name);
+ int total = consumers.size();
+ int from = Math.min((page - 1) * pageSize, total);
+ int to = Math.min(from + pageSize, total);
+ return TopicConsumerPageVO.builder()
+ .items(consumers.subList(from, to))
+ .total(total)
+ .page(page)
+ .pageSize(pageSize)
+ .build();
+ }
List<QueueProgressVO> getGroupProgress(String instanceId, String name);
List<SubscriptionEntryVO> getGroupSubscriptions(String instanceId, String
name);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 5ce9b517..9db74b72 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -41,6 +41,7 @@ import
org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.topic.BrokerRouteVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.persistence.entity.RmqGroup;
import org.apache.rocketmq.studio.persistence.entity.RmqTopic;
@@ -265,16 +266,22 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
@Override
public List<TopicConsumerVO> getTopicConsumers(String instanceId, String
name) {
+ return getTopicConsumersPage(instanceId, name, 1,
Integer.MAX_VALUE).getItems();
+ }
+
+ @Override
+ public TopicConsumerPageVO getTopicConsumersPage(String instanceId, String
name, int page, int pageSize) {
if (StringUtils.hasText(instanceId)) {
- return runtimeAdminClientResolver.execute(instanceId, admin ->
getTopicConsumers(admin, name));
+ return runtimeAdminClientResolver.execute(instanceId,
+ admin -> getTopicConsumersPage(admin, name, page,
pageSize));
}
if (!hasAdmin()) {
- return Collections.emptyList();
+ return
TopicConsumerPageVO.builder().items(List.of()).total(0).page(page).pageSize(pageSize).build();
}
- return adminExecute(admin -> getTopicConsumers(admin, name));
+ return adminExecute(admin -> getTopicConsumersPage(admin, name, page,
pageSize));
}
- private List<TopicConsumerVO> getTopicConsumers(MQAdminExt admin, String
name) {
+ private TopicConsumerPageVO getTopicConsumersPage(MQAdminExt admin, String
name, int page, int pageSize) {
try {
// Ask the broker who consumes this topic instead of scanning
every subscription
// group, which floods the result with system groups.
@@ -288,8 +295,13 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
}
+ List<String> sortedGroups = new ArrayList<>(subscribingGroups);
+ sortedGroups.sort(String::compareToIgnoreCase);
+ int total = sortedGroups.size();
+ int from = Math.min((page - 1) * pageSize, total);
+ int to = Math.min(from + pageSize, total);
List<TopicConsumerVO> consumers = new ArrayList<>();
- for (String group : subscribingGroups) {
+ for (String group : sortedGroups.subList(from, to)) {
try {
ConsumeStats stats = admin.examineConsumeStats(group,
name);
long diffTotal = 0;
@@ -333,8 +345,12 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
.build());
}
}
- consumers.sort((a, b) ->
a.getGroup().compareToIgnoreCase(b.getGroup()));
- return consumers;
+ return TopicConsumerPageVO.builder()
+ .items(consumers)
+ .total(total)
+ .page(page)
+ .pageSize(pageSize)
+ .build();
} catch (Exception e) {
log.warn("Failed to get consumers for topic {}: {}", name,
e.getMessage());
throw new BusinessException(502, "Failed to get consumers for
topic " + name + ": " + e.getMessage());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index 918a7fad..d3da87f4 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -171,6 +171,17 @@ class MetadataServiceTest {
verify(apacheProvider).getTopicConsumers("instance-a", "orders");
}
+ @Test
+ void topicConsumerPageShouldDelegateWithSelectedInstance() {
+ TopicConsumerPageVO page = TopicConsumerPageVO.builder()
+ .items(List.of()).total(3).page(1).pageSize(20).build();
+ when(apacheProvider.getTopicConsumersPage("instance-a", "orders", 1,
20)).thenReturn(page);
+
+ assertThat(metadataService.getTopicConsumersPage("instance-a",
"orders", 1, 20)).isSameAs(page);
+
+ verify(apacheProvider).getTopicConsumersPage("instance-a", "orders",
1, 20);
+ }
+
@Test
void
runtimeDiagnosticsShouldRejectBlankTopicAndGroupNamesBeforeProviderResolution()
{
assertThatThrownBy(() -> metadataService.getTopicRoutes("instance-a",
" "))
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
index 555658a3..a404edab 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
@@ -103,6 +103,22 @@ class TopicControllerTest {
verify(metadataService).getTopicConsumers("instance-a", "orders");
}
+ @Test
+ void topicConsumerPageShouldPassSelectedInstanceAndPaging() throws
Exception {
+ TopicConsumerPageVO page = TopicConsumerPageVO.builder()
+ .items(List.of()).total(3).page(2).pageSize(20).build();
+ when(metadataService.getTopicConsumersPage("instance-a", "orders", 2,
20)).thenReturn(page);
+
+ mockMvc.perform(get("/api/topics/orders/consumers/page")
+ .param("instanceId", "instance-a")
+ .param("page", "2")
+ .param("pageSize", "20"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.total").value(3));
+
+ verify(metadataService).getTopicConsumersPage("instance-a", "orders",
2, 20);
+ }
+
@Test
void createTopicShouldReturnCreatedTopic() throws Exception {
TopicVO input = new TopicVO();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
index ef521f19..d1870695 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
@@ -21,11 +21,13 @@ import
org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
import org.apache.rocketmq.tools.admin.MQAdminExt;
+import org.apache.rocketmq.remoting.protocol.body.GroupList;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.topic.BrokerRouteVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.persistence.entity.RmqGroup;
import org.apache.rocketmq.studio.persistence.mapper.RmqGroupMapper;
@@ -38,6 +40,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.util.HashSet;
import java.util.List;
+import java.util.HashSet;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -47,6 +50,7 @@ import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -116,14 +120,34 @@ class RocketMQMetadataProviderTest {
@Test
void getTopicConsumersShouldUseSelectedInstanceRuntimeClient() {
- List<TopicConsumerVO> consumers =
List.of(TopicConsumerVO.builder().group("cg-orders").build());
+ TopicConsumerPageVO consumers = TopicConsumerPageVO.builder()
+
.items(List.of(TopicConsumerVO.builder().group("cg-orders").build()))
+ .total(1).page(1).pageSize(Integer.MAX_VALUE).build();
when(runtimeAdminClientResolver.execute(eq("instance-a"),
any())).thenReturn(consumers);
RocketMQMetadataProvider provider = newProvider();
- assertThat(provider.getTopicConsumers("instance-a",
"orders")).containsExactlyElementsOf(consumers);
+ assertThat(provider.getTopicConsumers("instance-a", "orders"))
+
.extracting(TopicConsumerVO::getGroup).containsExactly("cg-orders");
verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
}
+ @Test
+ void
getTopicConsumersPageShouldOnlyFetchDiagnosticsForTheRequestedGroups() throws
Exception {
+ DefaultMQAdminExt admin = mock(DefaultMQAdminExt.class);
+ GroupList groups = new GroupList();
+ groups.setGroupList(new HashSet<>(List.of("group-a", "group-b",
"group-c")));
+ when(admin.queryTopicConsumeByWho("TopicA")).thenReturn(groups);
+
+ TopicConsumerPageVO result =
newLiveProvider(admin).getTopicConsumersPage(null, "TopicA", 2, 2);
+
+ assertThat(result.getTotal()).isEqualTo(3);
+
assertThat(result.getItems()).extracting(TopicConsumerVO::getGroup).containsExactly("group-c");
+ verify(admin).examineConsumeStats("group-c", "TopicA");
+ verify(admin).examineConsumerConnectionInfo("group-c");
+ verify(admin, never()).examineConsumeStats("group-a", "TopicA");
+ verify(admin, never()).examineConsumeStats("group-b", "TopicA");
+ }
+
@Test
void groupRuntimeDiagnosticsShouldUseSelectedInstanceRuntimeClient() {
List<QueueProgressVO> progress =
List.of(QueueProgressVO.builder().broker("broker-a").build());
diff --git a/web/src/api/metadata.test.ts b/web/src/api/metadata.test.ts
index dc77b50b..cbe0df1d 100644
--- a/web/src/api/metadata.test.ts
+++ b/web/src/api/metadata.test.ts
@@ -22,6 +22,7 @@ import {
createTopic,
deleteTopic,
getConsumerStack,
+ getTopicConsumerPage,
getTopicConsumers,
getTopicRoutes,
listTopics,
@@ -59,9 +60,23 @@ describe('topic metadata API', () => {
mock
.onGet('/topics/%25DLQ%25cg-order/consumers', { params: { instanceId:
'instance-a' } })
.reply(200, { code: 200, data: [] });
+ mock
+ .onGet('/topics/%25DLQ%25cg-order/consumers/page', {
+ params: { instanceId: 'instance-a', page: 2, pageSize: 20 },
+ })
+ .reply(200, {
+ code: 200,
+ data: { items: [], total: 21, page: 2, pageSize: 20 },
+ });
await expect(getTopicRoutes(topicName, 'instance-a')).resolves.toEqual([]);
await expect(getTopicConsumers(topicName,
'instance-a')).resolves.toEqual([]);
+ await expect(getTopicConsumerPage(topicName, 'instance-a', 2,
20)).resolves.toEqual({
+ items: [],
+ total: 21,
+ page: 2,
+ pageSize: 20,
+ });
});
it('encodes consumer stack route parameters and passes instanceId', async ()
=> {
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index a780a695..d614f6f0 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -47,6 +47,13 @@ export interface ConsumerGroupInfo {
metricsAvailable?: boolean;
}
+export interface TopicConsumerPage {
+ items: ConsumerGroupInfo[];
+ total: number;
+ page: number;
+ pageSize: number;
+}
+
// ─── Consumer Group (matches mock/consumers.ts) ─────────────────
export interface ConsumerGroup {
name: string;
@@ -161,6 +168,19 @@ export async function getTopicConsumers(name: string,
instanceId?: string) {
return res.data.data;
}
+export async function getTopicConsumerPage(
+ name: string,
+ instanceId: string | undefined,
+ page: number,
+ pageSize: number,
+) {
+ const res = await client.get<{ data: TopicConsumerPage }>(
+ `/topics/${encodeURIComponent(name)}/consumers/page`,
+ { params: { ...(instanceId ? { instanceId } : {}), page, pageSize } },
+ );
+ return res.data.data;
+}
+
export interface SendTopicMessageRequest {
topic: string;
instanceId?: string;
diff --git a/web/src/pages/instance/__tests__/TopicPage.test.tsx
b/web/src/pages/instance/__tests__/TopicPage.test.tsx
index 48dbc90c..6b2f37e9 100644
--- a/web/src/pages/instance/__tests__/TopicPage.test.tsx
+++ b/web/src/pages/instance/__tests__/TopicPage.test.tsx
@@ -29,6 +29,7 @@ const topicServiceMocks = vi.hoisted(() => ({
createTopic: vi.fn(),
deleteTopic: vi.fn(),
getTopicConsumers: vi.fn(),
+ getTopicConsumerPage: vi.fn(),
getTopicRoutes: vi.fn(),
listTopics: vi.fn(),
sendTopicMessage: vi.fn(),
@@ -135,6 +136,12 @@ describe('TopicPage', () => {
}));
topicServiceMocks.getTopicRoutes.mockResolvedValue([]);
topicServiceMocks.getTopicConsumers.mockResolvedValue([]);
+ topicServiceMocks.getTopicConsumerPage.mockResolvedValue({
+ items: [],
+ total: 0,
+ page: 1,
+ pageSize: 20,
+ });
instanceServiceMocks.listInstances.mockResolvedValue([
{
id: 'instance-proxy-1',
@@ -248,6 +255,12 @@ describe('TopicPage', () => {
await waitFor(() =>
expect(topicServiceMocks.getTopicRoutes).toHaveBeenCalledWith('topic-21',
'instance-a'),
);
+ expect(topicServiceMocks.getTopicConsumerPage).toHaveBeenCalledWith(
+ 'topic-21',
+ 'instance-a',
+ 1,
+ 20,
+ );
const closeButton = document.querySelector('.ant-modal-close');
expect(closeButton).not.toBeNull();
diff --git a/web/src/pages/instance/topic.tsx b/web/src/pages/instance/topic.tsx
index b7a23eeb..d5540c42 100644
--- a/web/src/pages/instance/topic.tsx
+++ b/web/src/pages/instance/topic.tsx
@@ -55,12 +55,12 @@ import {
import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import { TOPIC_TYPE_MAP, CLUSTER_TYPE_MAP } from '../../constants/theme';
-import type { Topic, BrokerRoute, ConsumerGroupInfo } from
'../../api/metadata';
+import type { Topic, BrokerRoute, ConsumerGroupInfo, TopicConsumerPage } from
'../../api/metadata';
import {
batchDeleteTopics,
createTopic,
deleteTopic,
- getTopicConsumers,
+ getTopicConsumerPage,
getTopicRoutes,
listTopics,
sendTopicMessage,
@@ -279,7 +279,7 @@ const TopicPage = () => {
const [topics, setTopics] = useState<Topic[]>([]);
const [loading, setLoading] = useState(false);
const [routesByTopic, setRoutesByTopic] = useState<Record<string,
BrokerRoute[]>>({});
- const [consumersByTopic, setConsumersByTopic] = useState<Record<string,
ConsumerGroupInfo[]>>({});
+ const [consumersByTopic, setConsumersByTopic] = useState<Record<string,
TopicConsumerPage>>({});
const [selectedRowKeys, setSelectedRowKeys] = useState<React.Key[]>([]);
const [searchText, setSearchText] = useState('');
const [typeFilter, setTypeFilter] = useState('');
@@ -363,18 +363,27 @@ const TopicPage = () => {
setTablePage(1);
};
+ const loadTopicConsumers = async (topic: Topic, page = 1, pageSize = 20) => {
+ const consumers = await getTopicConsumerPage(
+ topic.name,
+ selectedInstanceId || undefined,
+ page,
+ pageSize,
+ );
+ setConsumersByTopic((previous) => ({ ...previous, [topic.name]: consumers
}));
+ };
+
// ─── Open detail modal ────────────────────────────────────────
const openDetail = async (topic: Topic) => {
setSelectedTopic(topic);
setDetailModalOpen(true);
setDetailLoading(true);
try {
- const consumers = await getTopicConsumers(topic.name, selectedInstanceId
|| undefined);
+ await loadTopicConsumers(topic);
if (!isCloudInstance) {
const routes = await getTopicRoutes(topic.name, selectedInstanceId ||
undefined);
setRoutesByTopic((previous) => ({ ...previous, [topic.name]: routes
}));
}
- setConsumersByTopic((previous) => ({ ...previous, [topic.name]:
consumers }));
} catch {
message.error('Topic 详情加载失败,请稍后重试');
} finally {
@@ -406,7 +415,8 @@ const TopicPage = () => {
// ─── Route / consumer helpers ─────────────────────────────────
const getRoutes = (name: string): BrokerRoute[] => routesByTopic[name] ?? [];
- const getConsumers = (name: string): ConsumerGroupInfo[] =>
consumersByTopic[name] ?? [];
+ const getConsumerPage = (name: string): TopicConsumerPage =>
+ consumersByTopic[name] ?? { items: [], total: 0, page: 1, pageSize: 20 };
const handleAction = (key: string, topic: Topic) => {
if (key === 'detail') {
@@ -1115,9 +1125,18 @@ const TopicPage = () => {
</Text>
<Table<ConsumerGroupInfo>
columns={consumerColumns}
- dataSource={getConsumers(selectedTopic.name)}
+ dataSource={getConsumerPage(selectedTopic.name).items}
rowKey="group"
- pagination={false}
+ pagination={{
+ current: getConsumerPage(selectedTopic.name).page,
+ pageSize: getConsumerPage(selectedTopic.name).pageSize,
+ total: getConsumerPage(selectedTopic.name).total,
+ showSizeChanger: true,
+ pageSizeOptions: [10, 20, 50, 100],
+ onChange: (page, pageSize) => {
+ void loadTopicConsumers(selectedTopic, page, pageSize);
+ },
+ }}
size="small"
/>
</>
diff --git a/web/src/services/topicService.test.ts
b/web/src/services/topicService.test.ts
index 26415163..7aa9cbb4 100644
--- a/web/src/services/topicService.test.ts
+++ b/web/src/services/topicService.test.ts
@@ -15,7 +15,13 @@
* limitations under the License.
*/
-import { createTopic, getTopicConsumers, getTopicRoutes, listTopics } from
'./topicService';
+import {
+ createTopic,
+ getTopicConsumerPage,
+ getTopicConsumers,
+ getTopicRoutes,
+ listTopics,
+} from './topicService';
vi.mock('./dataMode', () => ({ isMockMode: () => true }));
vi.mock('../config', () => ({
@@ -56,6 +62,13 @@ describe('topic service mock data', () => {
expect(second[0]).not.toBe(first[0]);
});
+ it('paginates copied topic consumer rows', async () => {
+ const page = await getTopicConsumerPage('order-create', undefined, 1, 1);
+
+ expect(page).toMatchObject({ total: 4, page: 1, pageSize: 1 });
+ expect(page.items[0].group).toBe('GID_order_service');
+ });
+
it('trims search text before filtering topic names', async () => {
const topics = await listTopics({ search: ' ORDER-CREATE ' });
diff --git a/web/src/services/topicService.ts b/web/src/services/topicService.ts
index fc0f7e3f..5d552f7c 100644
--- a/web/src/services/topicService.ts
+++ b/web/src/services/topicService.ts
@@ -5,6 +5,7 @@ import type {
TopicQuery,
BrokerRoute,
ConsumerGroupInfo,
+ TopicConsumerPage,
SendTopicMessageRequest,
SendTopicMessageResult,
} from '../api/metadata';
@@ -106,6 +107,27 @@ export async function getTopicConsumers(
return metadataApi.getTopicConsumers(name, instanceId);
}
+export async function getTopicConsumerPage(
+ name: string,
+ instanceId: string | undefined,
+ page: number,
+ pageSize: number,
+): Promise<TopicConsumerPage> {
+ if (isMockMode()) {
+ const consumers = cloneConsumers(
+ (topicConsumers[name] as unknown as ConsumerGroupInfo[]) ?? [],
+ );
+ const from = Math.min((page - 1) * pageSize, consumers.length);
+ return {
+ items: consumers.slice(from, from + pageSize),
+ total: consumers.length,
+ page,
+ pageSize,
+ };
+ }
+ return metadataApi.getTopicConsumerPage(name, instanceId, page, pageSize);
+}
+
export async function sendTopicMessage(
data: SendTopicMessageRequest,
): Promise<SendTopicMessageResult> {