Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,9 @@ public ResponseEntity<Message<MetricsHistoryData>> 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"));
Expand All @@ -98,7 +100,9 @@ public ResponseEntity<Message<MetricsHistoryData>> 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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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";
}
Expand All @@ -162,9 +173,13 @@ public MetricsHistoryData getMetricHistoryData(String instance, String app, Stri
validateInstance(instance);
Map<String, List<Value>> 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("{}"));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,22 @@ default boolean supportsLogQuery() {
*/
Map<String, List<Value>> 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<String, List<Value>> 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
Expand All @@ -79,6 +95,22 @@ default boolean supportsLogQuery() {
*/
Map<String, List<Value>> 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<String, List<Value>> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -332,13 +332,21 @@ public void destroy() {
@Override
public Map<String, List<Value>> getHistoryMetricData(String instance, String app, String metrics, String metric,
String history) {
return getHistoryMetricData(null, instance, app, metrics, metric, history);
}

@Override
public Map<String, List<Value>> 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<String, List<Value>> instanceValuesMap = new HashMap<>(8);
Expand Down Expand Up @@ -408,6 +416,12 @@ public Map<String, List<Value>> getHistoryMetricData(String instance, String app
@Override
public Map<String, List<Value>> getHistoryIntervalMetricData(String instance, String app, String metrics,
String metric, String history) {
return getHistoryIntervalMetricData(null, instance, app, metrics, metric, history);
}

@Override
public Map<String, List<Value>> getHistoryIntervalMetricData(Long monitorId, String instance, String app,
String metrics, String metric, String history) {
if (!serverAvailable) {
log.error("""

Expand Down Expand Up @@ -439,7 +453,9 @@ public Map<String, List<Value>> 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<String, List<Value>> instanceValuesMap = new HashMap<>(8);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -349,12 +349,21 @@ public void destroy() {

@Override
public Map<String, List<Value>> getHistoryMetricData(String instance, String app, String metrics, String metric, String history) {
return getHistoryMetricData(null, instance, app, metrics, metric, history);
}

@Override
public Map<String, List<Value>> 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<String, List<Value>> instanceValuesMap = new HashMap<>(8);
try {
Expand Down Expand Up @@ -421,6 +430,12 @@ public Map<String, List<Value>> getHistoryMetricData(String instance, String app
@Override
public Map<String, List<Value>> getHistoryIntervalMetricData(String instance, String app, String metrics,
String metric, String history) {
return getHistoryIntervalMetricData(null, instance, app, metrics, metric, history);
}

@Override
public Map<String, List<Value>> getHistoryIntervalMetricData(Long monitorId, String instance, String app,
String metrics, String metric, String history) {
if (!serverAvailable) {
log.error("""

Expand Down Expand Up @@ -450,8 +465,11 @@ public Map<String, List<Value>> 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<String, List<Value>> instanceValuesMap = new HashMap<>(8);
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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))
Expand All @@ -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())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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<URI> 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<URI> 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();
Expand Down Expand Up @@ -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 + "\"");
}
}
Loading
Loading