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 3fcb3426 feat: support semantic metric queries (#766)
3fcb3426 is described below
commit 3fcb34266f4cf88375b34216cb090352499bc353
Author: yx9o <[email protected]>
AuthorDate: Mon Aug 3 11:39:10 2026 +0800
feat: support semantic metric queries (#766)
---
.../cluster/metrics/MetricProfileService.java | 18 +++++++
.../studio/cluster/metrics/MetricQueryDTO.java | 43 ++++++++++++++--
.../studio/cluster/metrics/MetricsService.java | 41 +++++++++++++---
.../cluster/metrics/MetricProfileServiceTest.java | 30 ++++++++++++
.../cluster/metrics/MetricsControllerTest.java | 42 ++++++++++++++++
.../studio/cluster/metrics/MetricsServiceTest.java | 57 ++++++++++++++++++++++
6 files changed, 221 insertions(+), 10 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
index ba7f3146..d41e574e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.cluster.metrics;
+import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Service;
import java.util.List;
@@ -30,6 +31,19 @@ public class MetricProfileService {
);
}
+ public String resolvePromql(String profileId, String semanticMetric) {
+ MetricProfileVO profile = listProfiles().stream()
+ .filter(candidate -> candidate.getId().equals(profileId))
+ .findFirst()
+ .orElseThrow(() -> badRequest("Unknown metric profile: " +
profileId));
+ return profile.getMetrics().stream()
+ .filter(metric ->
metric.getSemanticMetric().equals(semanticMetric))
+ .map(MetricProfileVO.MetricMappingVO::getPromql)
+ .findFirst()
+ .orElseThrow(() -> badRequest("Unknown semantic metric '" +
semanticMetric
+ + "' for profile '" + profileId + "'"));
+ }
+
private MetricProfileVO profile(MetricProfile profile,
List<MetricProfileVO.MetricMappingVO>
metrics) {
return MetricProfileVO.builder()
@@ -103,4 +117,8 @@ public class MetricProfileService {
.labels(List.of(labels))
.build();
}
+
+ private PrometheusException badRequest(String message) {
+ return new PrometheusException(HttpStatus.BAD_REQUEST.value(),
message);
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
index 7af771d8..b3dc389c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
@@ -16,7 +16,9 @@
*/
package org.apache.rocketmq.studio.cluster.metrics;
+import com.fasterxml.jackson.annotation.JsonIgnore;
import io.swagger.v3.oas.annotations.media.Schema;
+import jakarta.validation.constraints.AssertTrue;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Positive;
import jakarta.validation.constraints.Size;
@@ -31,13 +33,21 @@ import lombok.NoArgsConstructor;
@AllArgsConstructor
@Schema(description = "Prometheus range query")
public class MetricQueryDTO {
- @Schema(description = "PromQL expression evaluated by Prometheus",
- example = "sum(rate(rocketmq_messages_in_total[1m])) by
(node_id)", minLength = 1,
- requiredMode = Schema.RequiredMode.REQUIRED)
- @NotBlank(message = "Metric query is required")
+ @Schema(description = "Raw PromQL expression evaluated by Prometheus.
Mutually exclusive with profileId and "
+ + "semanticMetric", example =
"sum(rate(rocketmq_messages_in_total[1m])) by (node_id)", minLength = 1)
@Size(max = 4096, message = "Metric query must not exceed 4096 characters")
private String metric;
+ @Schema(description = "Metric profile used to resolve a semantic metric",
example = "rocketmq5-native",
+ minLength = 1)
+ @Size(max = 128, message = "Metric profile ID must not exceed 128
characters")
+ private String profileId;
+
+ @Schema(description = "Semantic metric key resolved through the selected
profile",
+ example = "consumer_lag_messages", minLength = 1)
+ @Size(max = 128, message = "Semantic metric must not exceed 128
characters")
+ private String semanticMetric;
+
@Schema(description = "Range start as a Unix timestamp in seconds",
example = "1784112606",
requiredMode = Schema.RequiredMode.REQUIRED)
@Positive(message = "Metric query start must be positive")
@@ -53,4 +63,29 @@ public class MetricQueryDTO {
@NotBlank(message = "Metric query step is required")
@Size(max = 32, message = "Metric query step must not exceed 32
characters")
private String step;
+
+ @JsonIgnore
+ @Schema(hidden = true)
+ @AssertTrue(message = "Metric query is required")
+ public boolean isMetricSelectionPresent() {
+ return hasText(metric) || hasText(profileId) ||
hasText(semanticMetric);
+ }
+
+ @JsonIgnore
+ @Schema(hidden = true)
+ @AssertTrue(message = "Metric profile and semantic metric are required
together")
+ public boolean isSemanticMetricSelectionComplete() {
+ return hasText(metric) || hasText(profileId) ==
hasText(semanticMetric);
+ }
+
+ @JsonIgnore
+ @Schema(hidden = true)
+ @AssertTrue(message = "Metric query cannot be combined with a semantic
metric selection")
+ public boolean isMetricSelectionExclusive() {
+ return !hasText(metric) || !hasText(profileId) &&
!hasText(semanticMetric);
+ }
+
+ private boolean hasText(String value) {
+ return value != null && !value.isBlank();
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
index a68eba17..f4fcb219 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
@@ -46,18 +46,20 @@ public class MetricsService {
);
private final MetricsSource metricsSource;
+ private final MetricProfileService metricProfileService;
public MetricDataVO query(MetricQueryDTO query) {
- validateQueryWindow(query);
+ if (query == null) {
+ throw badRequest("Metric query is required");
+ }
+ MetricQueryDTO resolvedQuery = resolveMetricQuery(query);
+ validateQueryWindow(resolvedQuery);
log.debug("Querying metrics: start={}, end={}, step={}",
- query.getStart(), query.getEnd(), query.getStep());
- return metricsSource.query(query);
+ resolvedQuery.getStart(), resolvedQuery.getEnd(),
resolvedQuery.getStep());
+ return metricsSource.query(resolvedQuery);
}
private void validateQueryWindow(MetricQueryDTO query) {
- if (query == null) {
- throw badRequest("Metric query is required");
- }
long rangeSeconds = query.getEnd() - query.getStart();
if (rangeSeconds <= 0) {
throw badRequest("Metric query end must be later than start");
@@ -78,6 +80,33 @@ public class MetricsService {
}
}
+ private MetricQueryDTO resolveMetricQuery(MetricQueryDTO query) {
+ boolean hasMetric = StringUtils.hasText(query.getMetric());
+ boolean hasProfile = StringUtils.hasText(query.getProfileId());
+ boolean hasSemanticMetric =
StringUtils.hasText(query.getSemanticMetric());
+ if (hasMetric) {
+ if (hasProfile || hasSemanticMetric) {
+ throw badRequest("Metric query cannot be combined with a
semantic metric selection");
+ }
+ return query;
+ }
+ if (hasProfile != hasSemanticMetric) {
+ throw badRequest("Metric profile and semantic metric are required
together");
+ }
+ if (!hasProfile) {
+ throw badRequest("Metric query is required");
+ }
+
+ String promql = metricProfileService.resolvePromql(
+ query.getProfileId().strip(),
query.getSemanticMetric().strip());
+ return MetricQueryDTO.builder()
+ .metric(promql)
+ .start(query.getStart())
+ .end(query.getEnd())
+ .step(query.getStep())
+ .build();
+ }
+
private BigDecimal parseStepMillis(String step) {
if (!StringUtils.hasText(step)) {
throw badRequest("Metric query step is required");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
index 3cd93129..e310cdf6 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
@@ -24,6 +24,7 @@ import java.util.function.Function;
import java.util.stream.Collectors;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
class MetricProfileServiceTest {
@@ -82,6 +83,26 @@ class MetricProfileServiceTest {
assertThat(messageOut.getName()).isEqualTo("Message Out TPS");
}
+ @Test
+ void resolvePromqlShouldReturnVersionSpecificMapping() {
+ assertThat(service.resolvePromql("rocketmq4-exporter",
"consumer_lag_messages"))
+ .isEqualTo("sum(rocketmq_message_accumulation) by (cluster,
group, topic)");
+ assertThat(service.resolvePromql("rocketmq5-native",
"consumer_lag_messages"))
+ .isEqualTo("sum(rocketmq_consumer_lag_messages) by (cluster,
topic, consumer_group)");
+ }
+
+ @Test
+ void resolvePromqlShouldRejectUnknownProfile() {
+ assertBadRequest(() -> service.resolvePromql("rocketmq6-native",
"message_in_tps"),
+ "Unknown metric profile: rocketmq6-native");
+ }
+
+ @Test
+ void resolvePromqlShouldRejectUnknownSemanticMetric() {
+ assertBadRequest(() -> service.resolvePromql("rocketmq5-native",
"queue_depth"),
+ "Unknown semantic metric 'queue_depth' for profile
'rocketmq5-native'");
+ }
+
private MetricProfileVO findProfile(String id) {
return service.listProfiles().stream()
.filter(profile -> profile.getId().equals(id))
@@ -107,4 +128,13 @@ class MetricProfileServiceTest {
.map(SemanticMetric::getKey)
.toList();
}
+
+ private void assertBadRequest(Runnable action, String message) {
+ assertThatExceptionOfType(PrometheusException.class)
+ .isThrownBy(action::run)
+ .satisfies(exception -> {
+ assertThat(exception.getStatusCode()).isEqualTo(400);
+ assertThat(exception.getMessage()).isEqualTo(message);
+ });
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
index 3b1af182..33bead3b 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
@@ -26,7 +26,10 @@ import org.springframework.test.web.servlet.MockMvc;
import java.util.List;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get;
import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.post;
@@ -83,6 +86,45 @@ class MetricsControllerTest {
verifyNoInteractions(metricsService);
}
+ @Test
+ void queryShouldAcceptSemanticMetricSelection() throws Exception {
+
when(metricsService.query(any(MetricQueryDTO.class))).thenReturn(MetricDataVO.builder()
+ .resultType("matrix")
+ .series(List.of())
+ .warnings(List.of())
+ .build());
+
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+
{"profileId":"rocketmq5-native","semanticMetric":"consumer_lag_messages",
+
"start":1784107658,"end":1784108558,"step":"30s"}
+ """))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200));
+
+ verify(metricsService).query(argThat(query ->
+ "rocketmq5-native".equals(query.getProfileId())
+ &&
"consumer_lag_messages".equals(query.getSemanticMetric())
+ && query.getMetric() == null));
+ }
+
+ @Test
+ void queryShouldRejectIncompleteSemanticMetricSelection() throws Exception
{
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+ {"profileId":"rocketmq5-native",
+
"start":1784107658,"end":1784108558,"step":"30s"}
+ """))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message")
+ .value("Metric profile and semantic metric are
required together"));
+
+ verifyNoInteractions(metricsService);
+ }
+
@Test
void queryShouldReturnBadRequestWhenStartIsNotPositive() throws Exception {
mockMvc.perform(post("/api/metrics/query")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
index b67e4adb..a0d790cd 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.cluster.metrics;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
@@ -39,6 +40,9 @@ class MetricsServiceTest {
@Mock
private MetricsSource metricsSource;
+ @Mock
+ private MetricProfileService metricProfileService;
+
@InjectMocks
private MetricsService metricsService;
@@ -100,6 +104,59 @@ class MetricsServiceTest {
verify(metricsSource).query(query);
}
+ @Test
+ void queryShouldResolveSemanticMetricBeforeCallingSource() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .profileId("rocketmq5-native")
+ .semanticMetric("consumer_lag_messages")
+ .start(1700000000L)
+ .end(1700003600L)
+ .step("1m")
+ .build();
+ String promql = "sum(rocketmq_consumer_lag_messages) by (cluster,
topic, consumer_group)";
+ when(metricProfileService.resolvePromql("rocketmq5-native",
"consumer_lag_messages"))
+ .thenReturn(promql);
+
when(metricsSource.query(any(MetricQueryDTO.class))).thenReturn(emptyMetricData());
+
+ metricsService.query(query);
+
+ ArgumentCaptor<MetricQueryDTO> captor =
ArgumentCaptor.forClass(MetricQueryDTO.class);
+ verify(metricsSource).query(captor.capture());
+ assertThat(captor.getValue())
+ .extracting(MetricQueryDTO::getMetric,
MetricQueryDTO::getStart,
+ MetricQueryDTO::getEnd, MetricQueryDTO::getStep)
+ .containsExactly(promql, 1700000000L, 1700003600L, "1m");
+ assertThat(query.getMetric()).isNull();
+ }
+
+ @Test
+ void queryShouldRejectIncompleteSemanticMetricSelection() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .profileId("rocketmq5-native")
+ .start(1700000000L)
+ .end(1700003600L)
+ .step("1m")
+ .build();
+
+ assertBadRequest(query, "Metric profile and semantic metric are
required together");
+ verifyNoInteractions(metricsSource, metricProfileService);
+ }
+
+ @Test
+ void queryShouldRejectMixedRawAndSemanticMetricSelection() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("up")
+ .profileId("rocketmq5-native")
+ .semanticMetric("broker_health")
+ .start(1700000000L)
+ .end(1700003600L)
+ .step("1m")
+ .build();
+
+ assertBadRequest(query, "Metric query cannot be combined with a
semantic metric selection");
+ verifyNoInteractions(metricsSource, metricProfileService);
+ }
+
@Test
void queryShouldHandleVariousStepSizes() {
MetricQueryDTO query15s =
MetricQueryDTO.builder().metric("cpu").start(1L).end(2L).step("15s").build();