diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorage.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorage.java index d0132de5b14..d6c500768f4 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorage.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorage.java @@ -45,9 +45,12 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import java.util.function.Function; import java.util.stream.Collectors; @@ -109,6 +112,7 @@ public class GreptimeDbDataStorage extends AbstractHistoryDataStorage { private final RestTemplate restTemplate; private final GreptimeSqlQueryExecutor greptimeSqlQueryExecutor; + private final AtomicLong rejectedLabelCollisionCount = new AtomicLong(); public GreptimeDbDataStorage(GreptimeProperties greptimeProperties, @Qualifier(WarehouseConstants.GREPTIME_QUERY_REST_TEMPLATE) @@ -160,6 +164,14 @@ public void saveData(CollectRep.MetricsData metricsData) { List fields = metricsData.getFields(); Map customLabels = metricsData.getLabels(); List fieldNames = fields.stream().map(CollectRep.Field::getName).collect(Collectors.toList()); + Set labelCollisions = findLabelCollisions(customLabels, fieldNames); + if (!labelCollisions.isEmpty()) { + long rejectedCount = rejectedLabelCollisionCount.incrementAndGet(); + log.error("[warehouse greptime] reject metrics data {} because custom labels contain " + + "storage-managed keys {}; cumulative rejected batches: {}.", + metricsData.getId(), labelCollisions, rejectedCount); + return; + } fields.forEach(field -> { if (field.getLabel()) { tableSchemaBuilder.addTag(field.getName(), DataType.String); @@ -233,6 +245,23 @@ public void saveData(CollectRep.MetricsData metricsData) { } } + private Set findLabelCollisions(Map customLabels, List fieldNames) { + if (customLabels == null || customLabels.isEmpty()) { + return Set.of(); + } + Set collisions = new TreeSet<>(); + for (String key : customLabels.keySet()) { + if (LABEL_KEY_INSTANCE.equals(key) || LABEL_KEY_TS.equals(key) || fieldNames.contains(key)) { + collisions.add(key); + } + } + return collisions; + } + + long getRejectedLabelCollisionCount() { + return rejectedLabelCollisionCount.get(); + } + @Override public Map> getHistoryMetricData(String instance, String app, String metrics, String metric, String history) { 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..d7af656f55d 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 @@ -36,6 +36,7 @@ import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -69,7 +70,6 @@ import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.stereotype.Component; -import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; import org.springframework.web.client.RestTemplate; import org.springframework.web.util.UriComponentsBuilder; @@ -107,6 +107,7 @@ public class VictoriaMetricsClusterDataStorage extends AbstractHistoryDataStorag private final VictoriaMetricsSelectProperties vmSelectProps; private final RestTemplate restTemplate; private final BlockingQueue metricsBufferQueue; + private final AtomicLong rejectedLabelCollisionCount = new AtomicLong(); private HashedWheelTimer metricsFlushTimer = null; private MetricsFlushTask metricsFlushtask = null; @@ -186,6 +187,15 @@ public void saveData(CollectRep.MetricsData metricsData) { metricsData.getId(), metricsData.getApp(), metricsData.getMetrics()); return; } + var managedLabelCollisions = + VictoriaMetricsDataStorage.findManagedLabelCollisions(metricsData.getLabels()); + if (!managedLabelCollisions.isEmpty()) { + long rejectedCount = rejectedLabelCollisionCount.incrementAndGet(); + log.error("[warehouse victoria-metrics] reject metrics data {} because custom labels contain " + + "HertzBeat-managed keys {}; cumulative rejected batches: {}.", + metricsData.getId(), managedLabelCollisions, rejectedCount); + return; + } Map defaultLabels = Maps.newHashMapWithExpectedSize(8); defaultLabels.put(MONITOR_METRICS_KEY, metricsData.getMetrics()); boolean isPrometheusAuto; @@ -243,10 +253,7 @@ public void saveData(CollectRep.MetricsData metricsData) { } labels.put(LABEL_KEY_MONITOR_ID, String.valueOf(metricsData.getId())); // add customized labels as identifier - var customizedLabels = metricsData.getLabels(); - if (!ObjectUtils.isEmpty(customizedLabels)) { - labels.putAll(customizedLabels); - } + VictoriaMetricsDataStorage.addCustomizedLabels(labels, metricsData.getLabels()); VictoriaMetricsDataStorage.VictoriaMetricsContent content = VictoriaMetricsDataStorage.VictoriaMetricsContent.builder() .metric(new HashMap<>(labels)) .values(new Double[]{entry.getValue()}) @@ -276,6 +283,10 @@ public void saveData(CollectRep.MetricsData metricsData) { } } + long getRejectedLabelCollisionCount() { + return rejectedLabelCollisionCount.get(); + } + @Override public void destroy() { if (metricsFlushTimer != null && !metricsFlushTimer.isStop()) { 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..3fcd9563f79 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 @@ -32,10 +32,13 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.zip.GZIPOutputStream; import com.google.common.collect.Maps; @@ -98,12 +101,18 @@ public class VictoriaMetricsDataStorage extends AbstractHistoryDataStorage { private static final String SPILT = "_"; private static final String MONITOR_METRICS_KEY = "__metrics__"; private static final String MONITOR_METRIC_KEY = "__metric__"; + private static final Set MANAGED_LABEL_KEYS = Set.of( + LABEL_KEY_NAME, + LABEL_KEY_MONITOR_ID, + MONITOR_METRICS_KEY, + MONITOR_METRIC_KEY); private static final long MAX_WAIT_MS = 500L; private static final int MAX_RETRIES = 3; private final VictoriaMetricsProperties victoriaMetricsProp; private final RestTemplate restTemplate; private final BlockingQueue metricsBufferQueue; + private final AtomicLong rejectedLabelCollisionCount = new AtomicLong(); private HashedWheelTimer metricsFlushTimer = null; private final VictoriaMetricsProperties.InsertConfig insertConfig; @@ -170,6 +179,14 @@ public void saveData(CollectRep.MetricsData metricsData) { metricsData.getId(), metricsData.getApp(), metricsData.getMetrics()); return; } + Set managedLabelCollisions = findManagedLabelCollisions(metricsData.getLabels()); + if (!managedLabelCollisions.isEmpty()) { + long rejectedCount = rejectedLabelCollisionCount.incrementAndGet(); + log.error("[warehouse victoria-metrics] reject metrics data {} because custom labels contain " + + "HertzBeat-managed keys {}; cumulative rejected batches: {}.", + metricsData.getId(), managedLabelCollisions, rejectedCount); + return; + } Map defaultLabels = Maps.newHashMapWithExpectedSize(8); defaultLabels.put(MONITOR_METRICS_KEY, metricsData.getMetrics()); boolean isPrometheusAuto = false; @@ -226,10 +243,7 @@ public void saveData(CollectRep.MetricsData metricsData) { } labels.put(LABEL_KEY_MONITOR_ID, String.valueOf(metricsData.getId())); // add customized labels as identifier - var customizedLabels = metricsData.getLabels(); - if (!ObjectUtils.isEmpty(customizedLabels)) { - labels.putAll(customizedLabels); - } + addCustomizedLabels(labels, metricsData.getLabels()); VictoriaMetricsContent content = VictoriaMetricsContent.builder() .metric(new HashMap<>(labels)) .values(new Double[]{entry.getValue()}) @@ -255,6 +269,26 @@ public void saveData(CollectRep.MetricsData metricsData) { sendVictoriaMetrics(contentList); } + static void addCustomizedLabels(Map labels, Map customizedLabels) { + if (ObjectUtils.isEmpty(customizedLabels)) { + return; + } + labels.putAll(customizedLabels); + } + + long getRejectedLabelCollisionCount() { + return rejectedLabelCollisionCount.get(); + } + + static Set findManagedLabelCollisions(Map customizedLabels) { + if (ObjectUtils.isEmpty(customizedLabels)) { + return Set.of(); + } + Set collisions = new TreeSet<>(customizedLabels.keySet()); + collisions.retainAll(MANAGED_LABEL_KEYS); + return collisions; + } + @Override public void destroy() { if (metricsFlushTimer != null && !metricsFlushTimer.isStop()) { diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorageTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorageTest.java index 3de76446b6c..925c168ddbf 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorageTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeDbDataStorageTest.java @@ -38,7 +38,6 @@ import io.greptime.models.Result; import io.greptime.models.Table; import io.greptime.models.WriteOk; -import io.greptime.v1.Common; import io.greptime.v1.RowData; import java.lang.reflect.Field; import java.util.ArrayList; @@ -48,7 +47,6 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; -import java.util.stream.Collectors; import org.apache.hertzbeat.common.constants.CommonConstants; import org.apache.hertzbeat.common.entity.arrow.ArrowCell; import org.apache.hertzbeat.common.entity.arrow.RowWrapper; @@ -132,7 +130,6 @@ void testSaveData() { when(mockResult.isOk()).thenReturn(true); CompletableFuture> mockFuture = CompletableFuture.completedFuture(mockResult); when(greptimeDb.write(any(Table.class))).thenReturn(mockFuture); - greptimeDbDataStorage = new GreptimeDbDataStorage(greptimeProperties, restTemplate, greptimeSqlQueryExecutor); // Test with valid metrics data @@ -156,15 +153,10 @@ void testSaveData() { } @Test - void testSaveDataWithCustomLabels() throws Exception { + void testSaveDataRejectsCustomLabelCollisions() { try (MockedStatic mockedStatic = mockStatic(GreptimeDB.class)) { mockedStatic.when(() -> GreptimeDB.create(any())).thenReturn(greptimeDb); - Result mockResult = mock(Result.class); - when(mockResult.isOk()).thenReturn(true); - CompletableFuture> mockFuture = CompletableFuture.completedFuture(mockResult); - when(greptimeDb.write(any(Table.class))).thenReturn(mockFuture); - greptimeDbDataStorage = new GreptimeDbDataStorage(greptimeProperties, restTemplate, greptimeSqlQueryExecutor); CollectRep.MetricsData metricsData = createMockMetricsData(true); @@ -176,31 +168,33 @@ void testSaveDataWithCustomLabels() throws Exception { when(metricsData.getLabels()).thenReturn(customLabels); when(metricsData.getInstance()).thenReturn("server1"); + greptimeDbDataStorage.saveData(metricsData); + + verify(greptimeDb, never()).write(any(Table.class)); + assertEquals(1, greptimeDbDataStorage.getRejectedLabelCollisionCount()); + } + } + + @Test + void testSaveDataAcceptsNonConflictingCustomLabels() throws Exception { + try (MockedStatic mockedStatic = mockStatic(GreptimeDB.class)) { + mockedStatic.when(() -> GreptimeDB.create(any())).thenReturn(greptimeDb); + @SuppressWarnings("unchecked") + Result mockResult = mock(Result.class); + when(mockResult.isOk()).thenReturn(true); + when(greptimeDb.write(any(Table.class))) + .thenReturn(CompletableFuture.completedFuture(mockResult)); + greptimeDbDataStorage = new GreptimeDbDataStorage( + greptimeProperties, restTemplate, greptimeSqlQueryExecutor); + CollectRep.MetricsData metricsData = createMockMetricsData(true); + when(metricsData.getLabels()).thenReturn(Map.of("env", "prod")); + ArgumentCaptor tableCaptor = ArgumentCaptor.forClass(Table.class); greptimeDbDataStorage.saveData(metricsData); verify(greptimeDb).write(tableCaptor.capture()); - Table capturedTable = tableCaptor.getValue(); - - List columnSchemas = getColumnSchemas(capturedTable); - List columnNames = columnSchemas.stream() - .map(RowData.ColumnSchema::getColumnName) - .collect(Collectors.toList()); - assertEquals(5, columnNames.size()); - assertEquals(2, Collections.frequency(columnNames, "instance")); - assertEquals(1, Collections.frequency(columnNames, "ts")); - assertEquals(1, Collections.frequency(columnNames, "usage")); - - List envColumnSchemas = columnSchemas.stream() - .filter(columnSchema -> "env".equals(columnSchema.getColumnName())) - .collect(Collectors.toList()); - assertEquals(1, envColumnSchemas.size()); - assertEquals(Common.SemanticType.TAG, envColumnSchemas.get(0).getSemanticType()); - - List rows = getRows(capturedTable); - assertEquals(1, rows.size()); - RowData.Row row = rows.get(0); - assertEquals("prod", row.getValuesList().get(4).getStringValue()); + List schemas = getColumnSchemas(tableCaptor.getValue()); + assertTrue(schemas.stream().anyMatch(schema -> "env".equals(schema.getColumnName()))); } } 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..a7feb59954f 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 @@ -32,6 +32,7 @@ import org.apache.hertzbeat.common.entity.arrow.ArrowCell; import org.apache.hertzbeat.common.entity.arrow.RowWrapper; import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.util.JsonUtil; import org.awaitility.Awaitility; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -48,16 +49,20 @@ 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.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; /** * Test case for {@link VictoriaMetricsDataStorage} */ -@ExtendWith(MockitoExtension.class) +@ExtendWith({MockitoExtension.class, OutputCaptureExtension.class}) @MockitoSettings(strictness = Strictness.LENIENT) class VictoriaMetricsDataStorageTest { @@ -73,6 +78,7 @@ class VictoriaMetricsDataStorageTest { private VictoriaMetricsDataStorage victoriaMetricsDataStorage; private final AtomicInteger postForEntityCount = new AtomicInteger(0); + private final AtomicReference lastPayload = new AtomicReference<>(); @BeforeEach void setUp() { @@ -97,6 +103,10 @@ void setUp() { eq(String.class) )).thenAnswer(invocation -> { postForEntityCount.incrementAndGet(); + HttpEntity httpEntity = invocation.getArgument(1); + if (httpEntity.getBody() instanceof String payload) { + lastPayload.set(payload); + } return responseEntity; }); } @@ -184,6 +194,51 @@ void testMultiThreadSaveDataBySize() { .isGreaterThanOrEqualTo(threadCount * writeSize / bufferSize)); } + @Test + void existingJobAndInstanceLabelsKeepTheirSeriesIdentity() { + when(victoriaMetricsProperties.insert()).thenReturn(new VictoriaMetricsProperties.InsertConfig( + 1, Integer.MAX_VALUE, new VictoriaMetricsProperties.Compression(false))); + CollectRep.MetricsData metricsData = generateMockedMetricsData(); + when(metricsData.getLabels()).thenReturn(Map.of( + "job", "custom-job", + "instance", "custom-instance", + "region", "west")); + victoriaMetricsDataStorage = new VictoriaMetricsDataStorage(victoriaMetricsProperties, restTemplate); + + victoriaMetricsDataStorage.saveData(metricsData); + + Awaitility.await() + .atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> assertThat(postForEntityCount.get()).isEqualTo(1)); + VictoriaMetricsDataStorage.VictoriaMetricsContent content = + JsonUtil.fromJson(lastPayload.get().trim(), VictoriaMetricsDataStorage.VictoriaMetricsContent.class); + assertThat(content.getMetric()) + .containsEntry("job", "custom-job") + .containsEntry("instance", "custom-instance") + .containsEntry("region", "west"); + } + + @Test + void managedLabelCollisionsRejectTheBatchWithDiagnostics(CapturedOutput output) { + when(victoriaMetricsProperties.insert()).thenReturn(new VictoriaMetricsProperties.InsertConfig( + 1, Integer.MAX_VALUE, new VictoriaMetricsProperties.Compression(false))); + CollectRep.MetricsData metricsData = generateMockedMetricsData(); + when(metricsData.getLabels()).thenReturn(Map.of( + "__name__", "custom-name", + "__monitor_id__", "custom-monitor")); + victoriaMetricsDataStorage = new VictoriaMetricsDataStorage(victoriaMetricsProperties, restTemplate); + + victoriaMetricsDataStorage.saveData(metricsData); + + assertThat(postForEntityCount.get()).isZero(); + assertThat(victoriaMetricsDataStorage.getRejectedLabelCollisionCount()).isEqualTo(1); + assertThat(output.getAll()) + .contains("__name__") + .contains("__monitor_id__") + .doesNotContain("custom-name") + .doesNotContain("custom-monitor"); + } + @AfterEach void stop() { if (victoriaMetricsDataStorage != null) { @@ -200,6 +255,8 @@ public static CollectRep.MetricsData generateMockedMetricsData() { when(mockMetricsData.getTime()).thenReturn(System.currentTimeMillis()); when(mockMetricsData.getCode()).thenReturn(CollectRep.Code.SUCCESS); when(mockMetricsData.getApp()).thenReturn("app"); + when(mockMetricsData.getInstance()).thenReturn("storage-instance"); + when(mockMetricsData.getLabels()).thenReturn(Map.of()); CollectRep.ValueRow mockValueRow = Mockito.mock(CollectRep.ValueRow.class); List columnValues = List.of("server-test-01", "68.7"); diff --git a/home/docs/start/victoria-metrics-init.md b/home/docs/start/victoria-metrics-init.md index eb0f3265a91..05a5c37d03c 100644 --- a/home/docs/start/victoria-metrics-init.md +++ b/home/docs/start/victoria-metrics-init.md @@ -148,6 +148,24 @@ warehouse: Once configured, restart HertzBeat to connect to the VictoriaMetrics cluster. +### Custom Label Collision Policy + +Monitor custom labels keep their existing Prometheus semantics when HertzBeat +writes to VictoriaMetrics: + +- `job`, `instance`, and ordinary custom labels continue to use the configured + custom values. Upgrading does not rename these labels or move new samples to + a different label set. +- `__name__`, `__monitor_id__`, `__metrics__`, and `__metric__` are managed by + HertzBeat and cannot be used as monitor custom-label keys. If one is present, + HertzBeat rejects that metrics batch and logs the conflicting key names + without logging their values. + +Before upgrading, inspect monitor custom labels and rename any of the four +HertzBeat-managed keys. Existing VictoriaMetrics series are not rewritten. +No migration is needed for monitors that use `job`, `instance`, or other +custom labels. + ### FAQ 1. Do both the time series databases need to be configured? Can they both be used?