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 836944e2f fix(group): report cloud queue offsets as unknown instead of
zero (#4902)
836944e2f is described below
commit 836944e2f280cb48d9a8176ef31768bee637efe3
Author: 烤化の初雪 <[email protected]>
AuthorDate: Thu Sep 24 18:21:31 2026 +0800
fix(group): report cloud queue offsets as unknown instead of zero (#4902)
fix(group): report cloud queue offsets as unknown instead of zero
The Tencent and Aliyun lag APIs report the lag per topic and expose no
per-queue offsets, but both providers filled brokerOffset and consumerOffset
with 0 for every row they build (TencentInstanceProvider#getGroupProgress,
AliyunConverters#toQueueProgressRows). The streamed offsets therefore read
as
real measurements in the 消费进度 table: a row showed "Broker Offset 0" and
"Consumer Offset 0" next to a real 堆积量, and the reset-offset preview took
the same zero as the current consumer offset.
Report the unknown sentinel instead, the value this codebase already uses
for
"cannot be determined" (ConsumerLagResolver.UNKNOWN, the -1 the broker sends
for a gRPC lag, onlineInstances, the reset preview's own offsets), and
render a
negative offset as unavailable in the progress table, which is what
formatOffsetValue already does for the reset preview.
QueueProgressVO gains UNKNOWN_OFFSET so the providers and the console share
one
name for the sentinel; the AI group-detail contract keeps receiving
numbers, so
its required schema is unaffected.
Regression tests: two Aliyun converter cases and one Tencent provider case
fail
with "expected: -1L but was: 0L" before the change, and a ConsumerPage case
pins that such an offset renders as unavailable instead of a number.
---
.../studio/instance/group/QueueProgressVO.java | 9 ++++++
.../studio/provider/alibaba/AliyunConverters.java | 10 ++++---
.../provider/tencent/TencentInstanceProvider.java | 7 +++--
.../provider/alibaba/AliyunConvertersLagTest.java | 30 ++++++++++++++++++++
.../tencent/TencentInstanceProviderTest.java | 19 +++++++++++++
.../pages/instance/__tests__/ConsumerPage.test.tsx | 33 ++++++++++++++++++++++
web/src/pages/instance/consumer.tsx | 7 +++--
7 files changed, 107 insertions(+), 8 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
index d7149b6df..9fc0599f5 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
@@ -26,6 +26,15 @@ import lombok.NoArgsConstructor;
@NoArgsConstructor
@AllArgsConstructor
public class QueueProgressVO {
+
+ /**
+ * Sentinel for an offset the provider cannot report, e.g. the per-topic
lag rows of the cloud
+ * providers, which carry no per-queue offsets at all. It matches the
{@code -1} the broker uses
+ * for an undeterminable lag, and the console renders a negative offset as
unavailable instead
+ * of a number that would read like a measurement.
+ */
+ public static final long UNKNOWN_OFFSET = -1L;
+
private String topic;
private String broker;
private int queueId;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
index 8209b18fc..0ccf5e2c7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
@@ -195,12 +195,14 @@ final class AliyunConverters {
for (Map.Entry<String, DataTopicLagMapValue> entry :
topicLagMap.entrySet()) {
long ready = entry.getValue() == null ||
entry.getValue().getReadyCount() == null
? 0L : entry.getValue().getReadyCount();
+ // The Aliyun API reports the lag per topic, so the row
carries no queue offsets;
+ // report the unknown sentinel instead of a zero that reads
like a measurement.
rows.add(QueueProgressVO.builder()
.topic(entry.getKey())
.broker("topic:" + entry.getKey())
.queueId(0)
- .brokerOffset(0L)
- .consumerOffset(0L)
+ .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+ .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
.diffTotal(ready)
.build());
}
@@ -213,8 +215,8 @@ final class AliyunConverters {
rows.add(QueueProgressVO.builder()
.broker("total")
.queueId(0)
- .brokerOffset(0L)
- .consumerOffset(0L)
+ .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+ .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
.diffTotal(totalLag.getReadyCount())
.build());
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
index 95f2990aa..93a6f0752 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
@@ -535,8 +535,11 @@ public class TencentInstanceProvider implements
InstanceProvider {
.topic(subscription.getTopic())
.broker("topic:" + subscription.getTopic())
.queueId(0)
- .brokerOffset(0L)
- .consumerOffset(0L)
+ // The Tencent API reports the lag per topic, so this row
carries no queue
+ // offsets; report the unknown sentinel instead of a zero
that the console
+ // would render as a real measurement next to the real lag.
+ .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+ .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
.diffTotal(subscription.getConsumerLag() == null ? 0L :
subscription.getConsumerLag())
.build());
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
index 9fd14f418..c9d980cfd 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
@@ -57,4 +57,34 @@ class AliyunConvertersLagTest {
assertThat(row.getDiffTotal()).isEqualTo(100L);
});
}
+
+ @Test
+ void queueOffsetsShouldBeUnknownWhenOnlyTheTopicLagIsKnownTest() {
+ GetConsumerGroupLagResponseBody.Data data =
GetConsumerGroupLagResponseBody.Data.builder()
+ .topicLagMap(Map.of("orders",
DataTopicLagMapValue.builder().readyCount(40L).build()))
+ .build();
+
+ // The Aliyun API reports the lag per topic and no per-queue offsets,
so the row must not
+ // claim offsets of zero next to the real lag.
+ assertThat(AliyunConverters.toQueueProgressRows(data)).singleElement()
+ .satisfies(row -> {
+
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+ assertThat(row.getDiffTotal()).isEqualTo(40L);
+ });
+ }
+
+ @Test
+ void aggregateRowOffsetsShouldBeUnknownTooTest() {
+ GetConsumerGroupLagResponseBody.Data data =
GetConsumerGroupLagResponseBody.Data.builder()
+
.totalLag(GetConsumerGroupLagResponseBody.TotalLag.builder().readyCount(100L).build())
+ .build();
+
+ assertThat(AliyunConverters.toQueueProgressRows(data)).singleElement()
+ .satisfies(row -> {
+ assertThat(row.getBroker()).isEqualTo("total");
+
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+ });
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
index 67371481a..91393f953 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
@@ -730,6 +730,25 @@ class TencentInstanceProviderTest {
assertThat(subscriptions.get(1).getConsistency()).isNull();
}
+ @Test
+ void getGroupProgressShouldReportUnknownQueueOffsetsTest() throws
Exception {
+ SubscriptionData subscription = new SubscriptionData();
+ subscription.setTopic("orders");
+ subscription.setConsumerLag(42L);
+ DescribeTopicListByGroupResponse response = new
DescribeTopicListByGroupResponse();
+ response.setData(new SubscriptionData[]{subscription});
+ when(client.DescribeTopicListByGroup(any())).thenReturn(response);
+
+ // The Tencent API exposes the lag per topic and no per-queue offsets,
so the row must not
+ // claim offsets of zero next to the real lag.
+ assertThat(provider.getGroupProgress(STUDIO_INSTANCE_ID,
"GID_test")).singleElement()
+ .satisfies(row -> {
+
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+ assertThat(row.getDiffTotal()).isEqualTo(42L);
+ });
+ }
+
@Test
void previewResetOffsetShouldAllowLimitedCloudPreviewTest() throws
Exception {
SubscriptionData subscription = new SubscriptionData();
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 4e371a1a9..f5dd217c0 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -703,6 +703,39 @@ describe('Consumer page', () => {
await waitFor(() =>
expect(within(panel).queryByText(/消费进度加载失败/)).toBeInTheDocument());
});
+ it('renders an offset the provider cannot report as unavailable', async ()
=> {
+ vi.mocked(consumerService.getConsumerProgress).mockResolvedValue([
+ {
+ topic: 'remote-topic',
+ // A cloud provider reports the lag per topic and no per-queue
offsets, so it sends the
+ // negative sentinel the backend uses for a value it cannot determine.
+ broker: 'topic:remote-topic',
+ queueId: 0,
+ brokerOffset: -1,
+ consumerOffset: -1,
+ diffTotal: 42,
+ },
+ ]);
+ const user = userEvent.setup({ pointerEventsCheck: 0 });
+ renderWithProviders(<ConsumerPage />);
+
+ await user.click(await screen.findByRole('button', { name: /详情/ }));
+ await user.click(await screen.findByRole('tab', { name: /消费进度/ }));
+ const progressPanel = await screen.findByRole('tabpanel', { name: /消费进度/
});
+ await waitFor(() =>
+
expect(within(progressPanel).getByText('remote-topic')).toBeInTheDocument(),
+ );
+
+ const row = within(progressPanel)
+ .getAllByRole('row')
+ .find((candidate) => within(candidate).queryByText('remote-topic'));
+ expect(row).toBeDefined();
+ const cells = within(row!)
+ .getAllByRole('cell')
+ .map((cell) => cell.textContent?.trim());
+ expect(cells.filter((cell) => cell === '-')).toHaveLength(2);
+ });
+
it('shows group health diagnostics from subscriptions, progress and
clients', async () => {
const riskyGroup: ConsumerGroup = {
...group,
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index 4e2036d91..cd0ce23b5 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -1268,8 +1268,11 @@ const ConsumerPageContent = ({
key: 'brokerOffset',
width: 140,
align: 'right',
+ // Cloud providers report the lag per topic and cannot supply per-queue
offsets, so they
+ // send the negative sentinel: show it as unavailable instead of a
number that would read
+ // like a measurement next to the real lag.
render: (offset: number) => (
- <Text style={{ fontFamily: 'monospace'
}}>{offset.toLocaleString()}</Text>
+ <Text style={{ fontFamily: 'monospace'
}}>{formatOffsetValue(offset)}</Text>
),
},
{
@@ -1279,7 +1282,7 @@ const ConsumerPageContent = ({
width: 150,
align: 'right',
render: (offset: number) => (
- <Text style={{ fontFamily: 'monospace'
}}>{offset.toLocaleString()}</Text>
+ <Text style={{ fontFamily: 'monospace'
}}>{formatOffsetValue(offset)}</Text>
),
},
{