From 8d03fda68de9a6ca5033cacc7e86fbd9d64037fa Mon Sep 17 00:00:00 2001 From: Logic Date: Thu, 30 Jul 2026 16:43:20 +0800 Subject: [PATCH 1/2] maintenance: preserve VictoriaMetrics label identity --- .../vm/VictoriaMetricsClusterDataStorage.java | 6 +--- .../tsdb/vm/VictoriaMetricsDataStorage.java | 24 +++++++++++--- .../vm/VictoriaMetricsDataStorageTest.java | 32 +++++++++++++++++++ 3 files changed, 53 insertions(+), 9 deletions(-) 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..4ae29064d7f 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 @@ -69,7 +69,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; @@ -243,10 +242,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()}) 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..7a1e1bbbd1b 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,6 +32,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; @@ -98,6 +99,13 @@ 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 RESERVED_LABEL_KEYS = Set.of( + LABEL_KEY_NAME, + LABEL_KEY_JOB, + LABEL_KEY_INSTANCE, + 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; @@ -226,10 +234,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 +260,17 @@ public void saveData(CollectRep.MetricsData metricsData) { sendVictoriaMetrics(contentList); } + static void addCustomizedLabels(Map labels, Map customizedLabels) { + if (ObjectUtils.isEmpty(customizedLabels)) { + return; + } + customizedLabels.forEach((key, value) -> { + if (!RESERVED_LABEL_KEYS.contains(key)) { + labels.put(key, value); + } + }); + } + @Override public void destroy() { if (metricsFlushTimer != null && !metricsFlushTimer.isStop()) { 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..e36dea8d45a 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 @@ -50,7 +50,9 @@ import org.springframework.http.ResponseEntity; import org.springframework.web.client.RestTemplate; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -184,6 +186,36 @@ void testMultiThreadSaveDataBySize() { .isGreaterThanOrEqualTo(threadCount * writeSize / bufferSize)); } + @Test + void customizedLabelsDoNotReplaceStorageIdentity() { + Map labels = new HashMap<>(Map.of( + "__name__", "cpu_usage", + "job", "linux", + "instance", "server-01", + "__monitor_id__", "42", + "__metrics__", "cpu", + "__metric__", "usage")); + Map customizedLabels = Map.of( + "__name__", "other_metric", + "job", "other_job", + "instance", "server-02", + "__monitor_id__", "99", + "__metrics__", "memory", + "__metric__", "free", + "region", "west"); + + VictoriaMetricsDataStorage.addCustomizedLabels(labels, customizedLabels); + + assertThat(labels) + .containsEntry("__name__", "cpu_usage") + .containsEntry("job", "linux") + .containsEntry("instance", "server-01") + .containsEntry("__monitor_id__", "42") + .containsEntry("__metrics__", "cpu") + .containsEntry("__metric__", "usage") + .containsEntry("region", "west"); + } + @AfterEach void stop() { if (victoriaMetricsDataStorage != null) { From 0d24ce02874d297fd84bd1878eae9c67c8b5817f Mon Sep 17 00:00:00 2001 From: Logic Date: Thu, 30 Jul 2026 23:23:00 +0800 Subject: [PATCH 2/2] maintenance: define VictoriaMetrics label collisions --- .../vm/VictoriaMetricsClusterDataStorage.java | 7 ++ .../tsdb/vm/VictoriaMetricsDataStorage.java | 31 +++++-- .../vm/VictoriaMetricsDataStorageTest.java | 80 ++++++++++++------- home/docs/start/victoria-metrics-init.md | 18 +++++ 4 files changed, 100 insertions(+), 36 deletions(-) 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 4ae29064d7f..73fa2a9191f 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 @@ -185,6 +185,13 @@ public void saveData(CollectRep.MetricsData metricsData) { metricsData.getId(), metricsData.getApp(), metricsData.getMetrics()); return; } + var managedLabelCollisions = + VictoriaMetricsDataStorage.findManagedLabelCollisions(metricsData.getLabels()); + if (!managedLabelCollisions.isEmpty()) { + log.error("[warehouse victoria-metrics] reject metrics data {} because custom labels contain " + + "HertzBeat-managed keys {}.", metricsData.getId(), managedLabelCollisions); + return; + } Map defaultLabels = Maps.newHashMapWithExpectedSize(8); defaultLabels.put(MONITOR_METRICS_KEY, metricsData.getMetrics()); boolean isPrometheusAuto; 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 7a1e1bbbd1b..4d2f09080fc 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 @@ -33,6 +33,7 @@ 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; @@ -99,10 +100,8 @@ 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 RESERVED_LABEL_KEYS = Set.of( + private static final Set MANAGED_LABEL_KEYS = Set.of( LABEL_KEY_NAME, - LABEL_KEY_JOB, - LABEL_KEY_INSTANCE, LABEL_KEY_MONITOR_ID, MONITOR_METRICS_KEY, MONITOR_METRIC_KEY); @@ -178,6 +177,12 @@ public void saveData(CollectRep.MetricsData metricsData) { metricsData.getId(), metricsData.getApp(), metricsData.getMetrics()); return; } + Set managedLabelCollisions = findManagedLabelCollisions(metricsData.getLabels()); + if (!managedLabelCollisions.isEmpty()) { + log.error("[warehouse victoria-metrics] reject metrics data {} because custom labels contain " + + "HertzBeat-managed keys {}.", metricsData.getId(), managedLabelCollisions); + return; + } Map defaultLabels = Maps.newHashMapWithExpectedSize(8); defaultLabels.put(MONITOR_METRICS_KEY, metricsData.getMetrics()); boolean isPrometheusAuto = false; @@ -264,11 +269,21 @@ static void addCustomizedLabels(Map labels, Map if (ObjectUtils.isEmpty(customizedLabels)) { return; } - customizedLabels.forEach((key, value) -> { - if (!RESERVED_LABEL_KEYS.contains(key)) { - labels.put(key, value); - } - }); + Set managedLabelCollisions = findManagedLabelCollisions(customizedLabels); + if (!managedLabelCollisions.isEmpty()) { + throw new IllegalArgumentException( + "Custom labels contain HertzBeat-managed keys " + managedLabelCollisions); + } + labels.putAll(customizedLabels); + } + + 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 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 e36dea8d45a..af3ff6fad66 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,18 +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.HashMap; 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 { @@ -75,6 +78,7 @@ class VictoriaMetricsDataStorageTest { private VictoriaMetricsDataStorage victoriaMetricsDataStorage; private final AtomicInteger postForEntityCount = new AtomicInteger(0); + private final AtomicReference lastPayload = new AtomicReference<>(); @BeforeEach void setUp() { @@ -99,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; }); } @@ -187,35 +195,49 @@ void testMultiThreadSaveDataBySize() { } @Test - void customizedLabelsDoNotReplaceStorageIdentity() { - Map labels = new HashMap<>(Map.of( - "__name__", "cpu_usage", - "job", "linux", - "instance", "server-01", - "__monitor_id__", "42", - "__metrics__", "cpu", - "__metric__", "usage")); - Map customizedLabels = Map.of( - "__name__", "other_metric", - "job", "other_job", - "instance", "server-02", - "__monitor_id__", "99", - "__metrics__", "memory", - "__metric__", "free", - "region", "west"); - - VictoriaMetricsDataStorage.addCustomizedLabels(labels, customizedLabels); - - assertThat(labels) - .containsEntry("__name__", "cpu_usage") - .containsEntry("job", "linux") - .containsEntry("instance", "server-01") - .containsEntry("__monitor_id__", "42") - .containsEntry("__metrics__", "cpu") - .containsEntry("__metric__", "usage") + 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(output.getAll()) + .contains("__name__") + .contains("__monitor_id__") + .doesNotContain("custom-name") + .doesNotContain("custom-monitor"); + } + @AfterEach void stop() { if (victoriaMetricsDataStorage != null) { @@ -232,6 +254,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?