From 26f84e2a86b99181e59e15ce59b7b462dfc2b8b8 Mon Sep 17 00:00:00 2001 From: limbo Date: Tue, 25 Aug 2026 14:31:35 +0800 Subject: [PATCH] feat: add monitorId parameter to metric history query API and update VictoriaMetrics storage to support monitor-based filtering --- .../controller/MetricsDataController.java | 8 +- .../warehouse/service/MetricsDataService.java | 17 +++ .../service/impl/MetricsDataServiceImpl.java | 19 +++- .../store/history/tsdb/HistoryDataReader.java | 32 ++++++ .../vm/VictoriaMetricsClusterDataStorage.java | 20 +++- .../tsdb/vm/VictoriaMetricsDataStorage.java | 22 +++- .../controller/MetricsDataControllerTest.java | 5 +- .../impl/MetricsDataServiceImplTest.java | 14 +++ ...VictoriaMetricsClusterDataStorageTest.java | 105 ++++++++++++++++++ .../vm/VictoriaMetricsDataStorageTest.java | 60 +++++++++- .../monitor-data-chart.component.ts | 1 + .../src/app/service/monitor.service.spec.ts | 15 +++ web-app/src/app/service/monitor.service.ts | 2 + 13 files changed, 305 insertions(+), 15 deletions(-) create mode 100644 hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorageTest.java diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/controller/MetricsDataController.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/controller/MetricsDataController.java index 832059cd2d1..be46486678c 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/controller/MetricsDataController.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/controller/MetricsDataController.java @@ -85,7 +85,9 @@ public ResponseEntity> getMetricHistoryData( @Parameter(description = "query historical time period, default 6h-6 hours: s-seconds, M-minutes, h-hours, d-days, w-weeks", example = "6h") @RequestParam(required = false) String history, @Parameter(description = "aggregate data calc. off by default; 4-hour window, query limit >1 week", example = "false") - @RequestParam(required = false) Boolean interval + @RequestParam(required = false) Boolean interval, + @Parameter(description = "monitor id", example = "343254354") + @RequestParam(required = false) Long monitorId ) { if (!metricsDataService.getWarehouseStorageServerStatus()) { return ResponseEntity.ok(Message.fail(FAIL_CODE, "time series database not available")); @@ -97,7 +99,9 @@ public ResponseEntity> getMetricHistoryData( String app = names[0]; String metrics = names[1]; String metric = names[2]; - MetricsHistoryData historyData = metricsDataService.getMetricHistoryData(instance, app, metrics, metric, history, interval); + MetricsHistoryData historyData = monitorId == null + ? metricsDataService.getMetricHistoryData(instance, app, metrics, metric, history, interval) + : metricsDataService.getMetricHistoryData(monitorId, instance, app, metrics, metric, history, interval); return ResponseEntity.ok(Message.success(historyData)); } } diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/MetricsDataService.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/MetricsDataService.java index 6c39b78d585..56cb39cd0f6 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/MetricsDataService.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/MetricsDataService.java @@ -51,4 +51,21 @@ public interface MetricsDataService { * @return metrics history data */ MetricsHistoryData getMetricHistoryData(String instance, String app, String metrics, String metric, String history, Boolean interval); + + /** + * Queries historical data for a specified monitor metric + * + * @param monitorId Monitor Id + * @param instance Instance e.g. ip:port or ip or domain + * @param app Monitor Type + * @param metrics Metrics Name + * @param metric Metrics Field Name + * @param history Query Historical Time Period + * @param interval aggregate data calc + * @return metrics history data + */ + default MetricsHistoryData getMetricHistoryData(Long monitorId, String instance, String app, String metrics, + String metric, String history, Boolean interval) { + return getMetricHistoryData(instance, app, metrics, metric, history, interval); + } } diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImpl.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImpl.java index df336e5e3ae..7819c59929a 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImpl.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImpl.java @@ -141,6 +141,17 @@ public MetricsData getMetricsData(Long monitorId, String metrics) { @Override public MetricsHistoryData getMetricHistoryData(String instance, String app, String metrics, String metric, String history, Boolean interval) { + return queryMetricHistoryData(null, instance, app, metrics, metric, history, interval); + } + + @Override + public MetricsHistoryData getMetricHistoryData(Long monitorId, String instance, String app, String metrics, + String metric, String history, Boolean interval) { + return queryMetricHistoryData(monitorId, instance, app, metrics, metric, history, interval); + } + + private MetricsHistoryData queryMetricHistoryData(Long monitorId, String instance, String app, String metrics, + String metric, String history, Boolean interval) { if (history == null) { history = "6h"; } @@ -151,9 +162,13 @@ public MetricsHistoryData getMetricHistoryData(String instance, String app, Stri validateInstance(instance); Map> instanceValuesMap; if (interval == null || !interval) { - instanceValuesMap = historyDataReader.get().getHistoryMetricData(instance, app, metrics, metric, history); + instanceValuesMap = monitorId != null + ? historyDataReader.get().getHistoryMetricData(monitorId, instance, app, metrics, metric, history) + : historyDataReader.get().getHistoryMetricData(instance, app, metrics, metric, history); } else { - instanceValuesMap = historyDataReader.get().getHistoryIntervalMetricData(instance, app, metrics, metric, history); + instanceValuesMap = monitorId != null + ? historyDataReader.get().getHistoryIntervalMetricData(monitorId, instance, app, metrics, metric, history) + : historyDataReader.get().getHistoryIntervalMetricData(instance, app, metrics, metric, history); } if (instanceValuesMap.containsKey("{}")) { instanceValuesMap.put("", instanceValuesMap.get("{}")); diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/HistoryDataReader.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/HistoryDataReader.java index 08a42bb9190..80e55b9e095 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/HistoryDataReader.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/HistoryDataReader.java @@ -52,6 +52,22 @@ default boolean supportsLogQuery() { */ Map> getHistoryMetricData(String instance, String app, String metrics, String metric, String history); + /** + * query history range metrics data from tsdb for a monitor + * + * @param monitorId monitor id + * @param instance instance e.g. ip:port or ip or domain + * @param app monitor type + * @param metrics metrics + * @param metric metric + * @param history range + * @return metrics data + */ + default Map> getHistoryMetricData(Long monitorId, String instance, String app, + String metrics, String metric, String history) { + return getHistoryMetricData(instance, app, metrics, metric, history); + } + /** * query history range interval metrics data from tsdb * max min mean metrics value @@ -65,6 +81,22 @@ default boolean supportsLogQuery() { */ Map> getHistoryIntervalMetricData(String instance, String app, String metrics, String metric, String history); + /** + * query history range interval metrics data from tsdb for a monitor + * + * @param monitorId monitor id + * @param instance instance e.g. ip:port or ip or domain + * @param app monitor type + * @param metrics metrics + * @param metric metric + * @param history history range + * @return metrics data + */ + default Map> getHistoryIntervalMetricData(Long monitorId, String instance, String app, + String metrics, String metric, String history) { + return getHistoryIntervalMetricData(instance, app, metrics, metric, history); + } + /** * Query logs with multiple filter conditions * @param startTime start time in milliseconds diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorage.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorage.java index ae93a1e7e2c..63d33e87c34 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorage.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorage.java @@ -286,13 +286,21 @@ public void destroy() { @Override public Map> getHistoryMetricData(String instance, String app, String metrics, String metric, String history) { + return getHistoryMetricData(null, instance, app, metrics, metric, history); + } + + @Override + public Map> getHistoryMetricData(Long monitorId, String instance, String app, String metrics, + String metric, String history) { String labelName = metrics + SPILT + metric; if (app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX)) { labelName = metrics; } String timeSeriesSelector = Stream.of( LABEL_KEY_NAME + "=\"" + labelName + "\"", - LABEL_KEY_INSTANCE + "=\"" + instance + "\"", + monitorId == null + ? LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + : LABEL_KEY_MONITOR_ID + "=\"" + monitorId + "\"", app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX) ? null : MONITOR_METRIC_KEY + "=\"" + metric + "\"" ).filter(Objects::nonNull).collect(Collectors.joining(",")); Map> instanceValuesMap = new HashMap<>(8); @@ -362,6 +370,12 @@ public Map> getHistoryMetricData(String instance, String app @Override public Map> getHistoryIntervalMetricData(String instance, String app, String metrics, String metric, String history) { + return getHistoryIntervalMetricData(null, instance, app, metrics, metric, history); + } + + @Override + public Map> getHistoryIntervalMetricData(Long monitorId, String instance, String app, + String metrics, String metric, String history) { if (!serverAvailable) { log.error(""" @@ -393,7 +407,9 @@ public Map> getHistoryIntervalMetricData(String instance, St } String timeSeriesSelector = Stream.of( LABEL_KEY_NAME + "=\"" + labelName + "\"", - LABEL_KEY_INSTANCE + "=\"" + instance + "\"", + monitorId == null + ? LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + : LABEL_KEY_MONITOR_ID + "=\"" + monitorId + "\"", app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX) ? null : MONITOR_METRIC_KEY + "=\"" + metric + "\"" ).filter(Objects::nonNull).collect(Collectors.joining(",")); Map> instanceValuesMap = new HashMap<>(8); diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorage.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorage.java index 5714e42186e..8c7ca02857d 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorage.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorage.java @@ -264,12 +264,21 @@ public void destroy() { @Override public Map> getHistoryMetricData(String instance, String app, String metrics, String metric, String history) { + return getHistoryMetricData(null, instance, app, metrics, metric, history); + } + + @Override + public Map> getHistoryMetricData(Long monitorId, String instance, String app, String metrics, + String metric, String history) { String labelName = metrics + SPILT + metric; if (app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX)) { labelName = metrics; } + String monitorSelector = monitorId == null + ? LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + : LABEL_KEY_MONITOR_ID + "=\"" + monitorId + "\""; String timeSeriesSelector = LABEL_KEY_NAME + "=\"" + labelName + "\"" - + "," + LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + + "," + monitorSelector + (app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX) ? "" : "," + MONITOR_METRIC_KEY + "=\"" + metric + "\""); Map> instanceValuesMap = new HashMap<>(8); try { @@ -336,6 +345,12 @@ public Map> getHistoryMetricData(String instance, String app @Override public Map> getHistoryIntervalMetricData(String instance, String app, String metrics, String metric, String history) { + return getHistoryIntervalMetricData(null, instance, app, metrics, metric, history); + } + + @Override + public Map> getHistoryIntervalMetricData(Long monitorId, String instance, String app, + String metrics, String metric, String history) { if (!serverAvailable) { log.error(""" @@ -365,8 +380,11 @@ public Map> getHistoryIntervalMetricData(String instance, St if (app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX)) { labelName = metrics; } + String monitorSelector = monitorId == null + ? LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + : LABEL_KEY_MONITOR_ID + "=\"" + monitorId + "\""; String timeSeriesSelector = LABEL_KEY_NAME + "=\"" + labelName + "\"" - + "," + LABEL_KEY_INSTANCE + "=\"" + instance + "\"" + + "," + monitorSelector + (app.startsWith(CommonConstants.PROMETHEUS_APP_PREFIX) ? "" : "," + MONITOR_METRIC_KEY + "=\"" + metric + "\""); Map> instanceValuesMap = new HashMap<>(8); try { diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/controller/MetricsDataControllerTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/controller/MetricsDataControllerTest.java index 674f2b5d333..dd2eb304f0d 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/controller/MetricsDataControllerTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/controller/MetricsDataControllerTest.java @@ -116,6 +116,7 @@ void getMetricsData() throws Exception { @Test void getMetricHistoryData() throws Exception { + final long monitorId = 599733946907392L; final String instance = "127.0.0.1:8081"; final String app = "linux"; final String metrics = "cpu"; @@ -133,6 +134,7 @@ void getMetricHistoryData() throws Exception { params.add("label", label); params.add("history", history); params.add("interval", String.valueOf(interval)); + params.add("monitorId", String.valueOf(monitorId)); when(metricsDataService.getWarehouseStorageServerStatus()).thenReturn(false); this.mockMvc.perform(MockMvcRequestBuilders.get(getUrl).params(params)) @@ -156,7 +158,8 @@ void getMetricHistoryData() throws Exception { .field(Field.builder().name(metric).type(CommonConstants.TYPE_NUMBER).build()) .build(); when(metricsDataService.getWarehouseStorageServerStatus()).thenReturn(true); - lenient().when(metricsDataService.getMetricHistoryData(eq(instance), eq(app), eq(metrics), eq(metric), eq(history), eq(interval))) + lenient().when(metricsDataService.getMetricHistoryData(eq(monitorId), eq(instance), eq(app), eq(metrics), + eq(metric), eq(history), eq(interval))) .thenReturn(metricsHistoryData); this.mockMvc.perform(MockMvcRequestBuilders.get(getUrl).params(params)) .andExpect(status().isOk()) diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImplTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImplTest.java index a5aea4718a7..c4a73faec10 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImplTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/service/impl/MetricsDataServiceImplTest.java @@ -105,6 +105,20 @@ void testMissingRangeFallsBackToTheDefault() { verify(historyDataReader).getHistoryMetricData("127.0.0.1", "linux", "cpu", "usage", "6h"); } + @Test + void testMonitorIdIsForwardedToHistoryStorage() { + long monitorId = 599733946907392L; + when(historyDataReader.getHistoryMetricData( + monitorId, "127.0.0.1", "linux", "cpu", "usage", "6h")) + .thenReturn(Map.of("", List.of(new Value("1", 1L)))); + + metricsDataService.getMetricHistoryData( + monitorId, "127.0.0.1", "linux", "cpu", "usage", "6h", false); + + verify(historyDataReader).getHistoryMetricData( + monitorId, "127.0.0.1", "linux", "cpu", "usage", "6h"); + } + @Test void testIdentifierEscapingTheQuotingIsRejected() { // a backtick closes a tdengine identifier, a double quote closes a questdb one diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorageTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorageTest.java new file mode 100644 index 00000000000..f5179ca0009 --- /dev/null +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsClusterDataStorageTest.java @@ -0,0 +1,105 @@ +/* + * 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.hertzbeat.warehouse.store.history.tsdb.vm; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.net.URI; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.http.HttpEntity; +import org.springframework.http.HttpMethod; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.client.RestTemplate; + +/** + * Test case for {@link VictoriaMetricsClusterDataStorage}. + */ +@ExtendWith(MockitoExtension.class) +class VictoriaMetricsClusterDataStorageTest { + + private static final long MONITOR_ID = 599733946907392L; + private static final String INSTANCE = "hdp-hadoop2:10003"; + + @Mock + private RestTemplate restTemplate; + + private VictoriaMetricsClusterDataStorage dataStorage; + + @BeforeEach + void setUp() { + VictoriaMetricsInsertProperties insertProperties = new VictoriaMetricsInsertProperties( + "http://localhost:8480", "", "", 1, 0); + VictoriaMetricsSelectProperties selectProperties = new VictoriaMetricsSelectProperties( + "http://localhost:8481", "", ""); + VictoriaMetricsClusterProperties clusterProperties = new VictoriaMetricsClusterProperties( + true, "42", insertProperties, selectProperties); + when(restTemplate.exchange( + anyString(), eq(HttpMethod.GET), any(HttpEntity.class), eq(String.class))) + .thenReturn(ResponseEntity.ok("{\"status\":\"success\"}")); + dataStorage = new VictoriaMetricsClusterDataStorage(clusterProperties, restTemplate); + } + + @Test + void shouldQueryHistoryByMonitorIdInsteadOfInstance() { + when(restTemplate.exchange( + any(URI.class), eq(HttpMethod.GET), any(HttpEntity.class), eq(String.class))) + .thenReturn(ResponseEntity.ok("")); + + dataStorage.getHistoryMetricData( + MONITOR_ID, INSTANCE, "flink", "taskmanager", "value", "6h"); + + ArgumentCaptor uriCaptor = ArgumentCaptor.forClass(URI.class); + verify(restTemplate).exchange( + uriCaptor.capture(), eq(HttpMethod.GET), any(HttpEntity.class), eq(String.class)); + assertUsesMonitorId(uriCaptor.getValue()); + } + + @Test + void shouldQueryIntervalHistoryByMonitorIdInsteadOfInstance() { + when(restTemplate.exchange( + any(URI.class), eq(HttpMethod.GET), any(HttpEntity.class), eq(PromQlQueryContent.class))) + .thenReturn(new ResponseEntity<>(HttpStatus.OK)); + + dataStorage.getHistoryIntervalMetricData( + MONITOR_ID, INSTANCE, "flink", "taskmanager", "value", "1w"); + + ArgumentCaptor uriCaptor = ArgumentCaptor.forClass(URI.class); + verify(restTemplate, times(4)).exchange( + uriCaptor.capture(), eq(HttpMethod.GET), any(HttpEntity.class), eq(PromQlQueryContent.class)); + assertThat(uriCaptor.getAllValues()).allSatisfy(this::assertUsesMonitorId); + } + + private void assertUsesMonitorId(URI uri) { + assertThat(uri.getQuery()) + .contains("__monitor_id__=\"" + MONITOR_ID + "\"") + .doesNotContain("instance=\"" + INSTANCE + "\""); + } +} diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorageTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorageTest.java index f11fdb905dc..ea40750dff3 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorageTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/vm/VictoriaMetricsDataStorageTest.java @@ -17,12 +17,19 @@ package org.apache.hertzbeat.warehouse.store.history.tsdb.vm; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.startsWith; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -import static org.assertj.core.api.Assertions.assertThat; + +import java.net.URI; +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.arrow.vector.types.pojo.ArrowType; import org.apache.arrow.vector.types.pojo.Field; @@ -38,6 +45,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; @@ -50,10 +58,6 @@ import org.springframework.http.ResponseEntity; import org.springframework.web.client.RestTemplate; -import java.util.List; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; - /** * Test case for {@link VictoriaMetricsDataStorage} */ @@ -184,6 +188,50 @@ void testMultiThreadSaveDataBySize() { .isGreaterThanOrEqualTo(threadCount * writeSize / bufferSize)); } + @Test + void shouldQueryHistoryByMonitorIdInsteadOfInstance() { + long monitorId = 599733946907392L; + when(restTemplate.exchange( + any(URI.class), + eq(HttpMethod.GET), + any(HttpEntity.class), + eq(String.class) + )).thenReturn(ResponseEntity.ok("")); + victoriaMetricsDataStorage = new VictoriaMetricsDataStorage(victoriaMetricsProperties, restTemplate); + + victoriaMetricsDataStorage.getHistoryMetricData( + monitorId, "hdp-hadoop2:10003", "flink", "taskmanager", "value", "6h"); + + ArgumentCaptor uriCaptor = ArgumentCaptor.forClass(URI.class); + verify(restTemplate).exchange( + uriCaptor.capture(), eq(HttpMethod.GET), any(HttpEntity.class), eq(String.class)); + assertThat(uriCaptor.getValue().getQuery()) + .contains("__monitor_id__=\"" + monitorId + "\"") + .doesNotContain("instance=\"hdp-hadoop2:10003\""); + } + + @Test + void shouldQueryIntervalHistoryByMonitorIdInsteadOfInstance() { + long monitorId = 599733946907392L; + when(restTemplate.exchange( + any(URI.class), + eq(HttpMethod.GET), + any(HttpEntity.class), + eq(PromQlQueryContent.class) + )).thenReturn(new ResponseEntity<>(HttpStatus.OK)); + victoriaMetricsDataStorage = new VictoriaMetricsDataStorage(victoriaMetricsProperties, restTemplate); + + victoriaMetricsDataStorage.getHistoryIntervalMetricData( + monitorId, "hdp-hadoop2:10003", "flink", "taskmanager", "value", "1w"); + + ArgumentCaptor uriCaptor = ArgumentCaptor.forClass(URI.class); + verify(restTemplate, times(4)).exchange( + uriCaptor.capture(), eq(HttpMethod.GET), any(HttpEntity.class), eq(PromQlQueryContent.class)); + assertThat(uriCaptor.getAllValues()).allSatisfy(uri -> assertThat(uri.getQuery()) + .contains("__monitor_id__=\"" + monitorId + "\"") + .doesNotContain("instance=\"hdp-hadoop2:10003\"")); + } + @AfterEach void stop() { if (victoriaMetricsDataStorage != null) { diff --git a/web-app/src/app/routes/monitor/monitor-data-chart/monitor-data-chart.component.ts b/web-app/src/app/routes/monitor/monitor-data-chart/monitor-data-chart.component.ts index d98c6962778..659870c13e9 100644 --- a/web-app/src/app/routes/monitor/monitor-data-chart/monitor-data-chart.component.ts +++ b/web-app/src/app/routes/monitor/monitor-data-chart/monitor-data-chart.component.ts @@ -267,6 +267,7 @@ export class MonitorDataChartComponent implements OnInit, OnDestroy { this.loading = `${this.i18nSvc.fanyi('monitor.detail.chart.data-loading')}`; let metricData$ = this.monitorSvc .getMonitorMetricHistoryData( + this.monitorId, this.instance, this.app == 'prometheus' ? `_prometheus_${this.monitorName}` : this.app, this.metrics, diff --git a/web-app/src/app/service/monitor.service.spec.ts b/web-app/src/app/service/monitor.service.spec.ts index c431e2b85c1..f16268ddfbe 100644 --- a/web-app/src/app/service/monitor.service.spec.ts +++ b/web-app/src/app/service/monitor.service.spec.ts @@ -17,6 +17,7 @@ * under the License. */ +import { HttpTestingController } from '@angular/common/http/testing'; import { TestBed } from '@angular/core/testing'; import { configureHttpServiceTest } from '@testing'; @@ -24,13 +25,27 @@ import { MonitorService } from './monitor.service'; describe('MonitorService', () => { let service: MonitorService; + let http: HttpTestingController; beforeEach(() => { configureHttpServiceTest(); service = TestBed.inject(MonitorService); + http = TestBed.inject(HttpTestingController); }); + afterEach(() => http.verify()); + it('should be created', () => { expect(service).toBeTruthy(); }); + + it('should query metric history with monitor id', () => { + service.getMonitorMetricHistoryData(599733946907392, 'hdp-hadoop2:10003', 'flink', 'taskmanager', 'value', '6h', false).subscribe(); + + const request = http.expectOne(req => req.url === '/monitor/hdp-hadoop2:10003/metric/flink.taskmanager.value'); + expect(request.request.params.get('monitorId')).toBe('599733946907392'); + expect(request.request.params.get('history')).toBe('6h'); + expect(request.request.params.get('interval')).toBe('false'); + request.flush({ code: 0 }); + }); }); diff --git a/web-app/src/app/service/monitor.service.ts b/web-app/src/app/service/monitor.service.ts index 5919ee2779a..218f1e23fd7 100644 --- a/web-app/src/app/service/monitor.service.ts +++ b/web-app/src/app/service/monitor.service.ts @@ -164,6 +164,7 @@ export class MonitorService { } public getMonitorMetricHistoryData( + monitorId: number, instance: string, app: string, metrics: string, @@ -174,6 +175,7 @@ export class MonitorService { let metricFull = `${app}.${metrics}.${metric}`; let httpParams = new HttpParams(); httpParams = httpParams.appendAll({ + monitorId: monitorId, history: history, interval: interval });