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 b9844a85390..a4d7bbe4872 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 @@ -86,7 +86,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")); @@ -98,7 +100,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 ea2222cc31b..379460fa2e0 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 @@ -58,4 +58,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 a1c81ebc8ea..3085a4e07ed 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 @@ -152,6 +152,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"; } @@ -162,9 +173,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 fd821a7bfd5..07f789fa2d3 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 @@ -66,6 +66,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 @@ -79,6 +95,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 e808f62202c..7f14e861c52 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 @@ -332,13 +332,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); @@ -408,6 +416,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(""" @@ -439,7 +453,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 7f199ff4f7d..bfd5c2a0a74 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 @@ -349,12 +349,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 { @@ -421,6 +430,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(""" @@ -450,8 +465,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 ce26d541f71..94e2ba7d5c3 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 @@ -123,6 +123,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"; @@ -140,6 +141,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)) @@ -163,7 +165,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 index 960b28ea893..b102f64acbf 100644 --- 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 @@ -22,8 +22,11 @@ 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 java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; @@ -32,9 +35,11 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; + import org.apache.hertzbeat.common.timer.TimerTask; 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.mockito.junit.jupiter.MockitoSettings; @@ -53,9 +58,54 @@ @MockitoSettings(strictness = Strictness.LENIENT) class VictoriaMetricsClusterDataStorageTest { + private static final long MONITOR_ID = 599733946907392L; + private static final String INSTANCE = "hdp-hadoop2:10003"; + @Mock private RestTemplate restTemplate; + @Test + void shouldQueryHistoryByMonitorIdInsteadOfInstance() { + mockHealthCheck(); + when(restTemplate.exchange( + any(URI.class), eq(HttpMethod.GET), any(HttpEntity.class), eq(String.class))) + .thenReturn(ResponseEntity.ok("")); + VictoriaMetricsClusterDataStorage storage = createStorage(1, 0); + + try { + storage.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()); + } finally { + storage.destroy(); + } + } + + @Test + void shouldQueryIntervalHistoryByMonitorIdInsteadOfInstance() { + mockHealthCheck(); + when(restTemplate.exchange( + any(URI.class), eq(HttpMethod.GET), any(HttpEntity.class), eq(PromQlQueryContent.class))) + .thenReturn(new ResponseEntity<>(HttpStatus.OK)); + VictoriaMetricsClusterDataStorage storage = createStorage(1, 0); + + try { + storage.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); + } finally { + storage.destroy(); + } + } + @Test void flushesDataAddedWhileAnImmediateFlushIsRunning() throws Exception { mockHealthCheck(); @@ -246,4 +296,10 @@ private static void saveOneMetric(VictoriaMetricsClusterDataStorage storage) { private static long lineCount(String body) { return body.lines().filter(line -> !line.isBlank()).count(); } + + 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 31528e84cc4..67d2997bf38 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,26 @@ 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.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import org.apache.arrow.vector.types.pojo.ArrowType; import org.apache.arrow.vector.types.pojo.Field; @@ -38,33 +52,20 @@ 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.Mockito; import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoSettings; import org.mockito.quality.Strictness; - +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; import org.springframework.http.HttpEntity; import org.springframework.http.HttpMethod; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; -import org.springframework.boot.test.system.CapturedOutput; -import org.springframework.boot.test.system.OutputCaptureExtension; -import org.springframework.web.client.RestTemplate; - -import java.util.List; -import java.util.Map; -import java.util.concurrent.CopyOnWriteArrayList; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.atomic.AtomicReference; - import org.springframework.test.util.ReflectionTestUtils; +import org.springframework.web.client.RestTemplate; /** * Test case for {@link VictoriaMetricsDataStorage} @@ -201,6 +202,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\"")); + } + @Test void failedSingleNodeFlushRetainsTheBatchAndRetriesQuickly() { when(victoriaMetricsProperties.insert()).thenReturn(new VictoriaMetricsProperties.InsertConfig( 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 });