This is an automated email from the ASF dual-hosted git repository. zqr10159 pushed a commit to branch 2.0.0 in repository https://gitbox.apache.org/repos/asf/hertzbeat.git
commit 72e7fcb0c721e1cbd330712e49ec6fa07a8af306 Author: Logic <[email protected]> AuthorDate: Fri Aug 28 10:08:07 2026 +0800 Add native metric system dimensions --- .../MetricsCollectExecutionContextTest.java | 14 ++ .../entity/metric/NativeMetricSystemContext.java | 89 +++++++ .../metric/NativeMetricSystemDimensions.java | 55 +++++ ...reptimeNativeMetricSystemDimensionsE2eTest.java | 257 +++++++++++++++++++++ hertzbeat-manager/pom.xml | 4 + .../ManagerNativeMetricSystemContextResolver.java | 114 +++++++++ ...nagerNativeMetricSystemContextResolverTest.java | 132 +++++++++++ .../tsdb/greptime/GreptimeDbDataStorage.java | 92 ++++++-- .../tsdb/greptime/GreptimeMetricSchemaCache.java | 84 ++++--- .../NativeMetricSystemContextResolver.java | 36 +++ .../tsdb/greptime/GreptimeDbDataStorageTest.java | 158 ++++++++++++- .../greptime/GreptimeMetricSchemaCacheTest.java | 47 ++++ 12 files changed, 1029 insertions(+), 53 deletions(-) diff --git a/hertzbeat-collector/hertzbeat-collector-collector/src/test/java/org/apache/hertzbeat/collector/dispatch/MetricsCollectExecutionContextTest.java b/hertzbeat-collector/hertzbeat-collector-collector/src/test/java/org/apache/hertzbeat/collector/dispatch/MetricsCollectExecutionContextTest.java index e1cf4aac10..a227f95a14 100644 --- a/hertzbeat-collector/hertzbeat-collector-collector/src/test/java/org/apache/hertzbeat/collector/dispatch/MetricsCollectExecutionContextTest.java +++ b/hertzbeat-collector/hertzbeat-collector-collector/src/test/java/org/apache/hertzbeat/collector/dispatch/MetricsCollectExecutionContextTest.java @@ -18,8 +18,10 @@ package org.apache.hertzbeat.collector.dispatch; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import org.apache.hertzbeat.common.constants.MetricDataConstants; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; import org.apache.hertzbeat.common.entity.message.CollectRep; import org.junit.jupiter.api.Test; @@ -28,16 +30,28 @@ class MetricsCollectExecutionContextTest { @Test void addsCompactExecutionContextToArrowSchemaMetadata() { CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder() + .setId(42L) .setApp("linux") .setMetrics("cpu") .setTime(1_700L) .setCode(CollectRep.Code.SUCCESS); + builder.addMetadata(MetricDataConstants.INSTANCE, "db.internal:3306"); + builder.addMetadata("hertzbeat.workspace.id", "spoof-workspace"); + builder.addMetadata(MetricDataConstants.ENTITY_ID, "999"); + builder.addMetadata("hertzbeat.entity.type", "spoof-type"); MetricsCollect.addExecutionContext(builder, 1_000L, "collector-arm-1"); try (CollectRep.MetricsData metricsData = builder.build()) { assertEquals("1000", metricsData.getMetadataValue(MetricDataConstants.COLLECTION_STARTED_AT)); assertEquals("collector-arm-1", metricsData.getMetadataValue(MetricDataConstants.COLLECTOR_ID)); + NativeMetricSystemContext context = NativeMetricSystemContext.from(metricsData); + assertEquals(42L, context.monitorId()); + assertEquals("collector-arm-1", context.collectorId()); + assertEquals("db.internal:3306", context.instance()); + assertNull(context.workspaceId()); + assertNull(context.entityId()); + assertNull(context.entityType()); } } } diff --git a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemContext.java b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemContext.java new file mode 100644 index 0000000000..de9547e8a3 --- /dev/null +++ b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemContext.java @@ -0,0 +1,89 @@ +/* + * 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.common.entity.metric; + +import org.apache.hertzbeat.common.constants.MetricDataConstants; +import org.apache.hertzbeat.common.entity.message.CollectRep; + +/** + * Authoritative HertzBeat-owned identity carried beside one native metric sample. + * Missing optional authority remains {@code null}; it is never converted to a fake zero or unknown value. + */ +public record NativeMetricSystemContext( + String workspaceId, + Long entityId, + String entityType, + Long monitorId, + String collectorId, + String instance) { + + public NativeMetricSystemContext { + workspaceId = trimToNull(workspaceId); + entityId = positiveOrNull(entityId); + entityType = trimToNull(entityType); + monitorId = positiveOrNull(monitorId); + collectorId = trimToNull(collectorId); + instance = trimToNull(instance); + } + + /** + * Resolve the context that is intrinsic to a Collector response before manager entity enrichment. + * + * @param metricsData native metric response + * @return intrinsic system context + */ + public static NativeMetricSystemContext from(CollectRep.MetricsData metricsData) { + if (metricsData == null) { + return new NativeMetricSystemContext(null, null, null, null, null, null); + } + return new NativeMetricSystemContext( + null, + null, + null, + metricsData.getId(), + metricsData.getMetadataValue(MetricDataConstants.COLLECTOR_ID), + metricsData.getInstance()); + } + + /** + * Return the fixed Greptime tag values in {@link NativeMetricSystemDimensions#TAG_NAMES} order. + * + * @return fixed-order tag values + */ + public Object[] tagValues() { + return new Object[] { + workspaceId, + entityId == null ? null : String.valueOf(entityId), + entityType, + monitorId == null ? null : String.valueOf(monitorId), + collectorId, + instance + }; + } + + private static Long positiveOrNull(Long value) { + return value != null && value > 0 ? value : null; + } + + private static String trimToNull(String value) { + if (value == null || value.isBlank()) { + return null; + } + return value.trim(); + } +} diff --git a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemDimensions.java b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemDimensions.java new file mode 100644 index 0000000000..6fa9b7e9ed --- /dev/null +++ b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/metric/NativeMetricSystemDimensions.java @@ -0,0 +1,55 @@ +/* + * 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.common.entity.metric; + +import java.util.List; +import java.util.Locale; +import java.util.Set; + +/** + * Physical Greptime tag names owned by HertzBeat for native metric samples. + */ +public final class NativeMetricSystemDimensions { + + public static final String WORKSPACE_ID = "hertzbeat_workspace_id"; + public static final String ENTITY_ID = "hertzbeat_entity_id"; + public static final String ENTITY_TYPE = "hertzbeat_entity_type"; + public static final String MONITOR_ID = "hertzbeat_monitor_id"; + public static final String COLLECTOR_ID = "hertzbeat_collector_id"; + public static final String INSTANCE = "instance"; + public static final String TIMESTAMP = "ts"; + + public static final List<String> TAG_NAMES = List.of( + WORKSPACE_ID, ENTITY_ID, ENTITY_TYPE, MONITOR_ID, COLLECTOR_ID, INSTANCE); + + private static final Set<String> RESERVED_NAMES = Set.of( + WORKSPACE_ID, ENTITY_ID, ENTITY_TYPE, MONITOR_ID, COLLECTOR_ID, INSTANCE, TIMESTAMP); + + private NativeMetricSystemDimensions() { + } + + /** + * Whether a Collector field name would collide with a HertzBeat-owned metric column. + * + * @param fieldName Collector field name + * @return true when the field must not be projected as an application column + */ + public static boolean isReserved(String fieldName) { + return fieldName != null && RESERVED_NAMES.contains(fieldName.trim().toLowerCase(Locale.ROOT)); + } +} diff --git a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeNativeMetricSystemDimensionsE2eTest.java b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeNativeMetricSystemDimensionsE2eTest.java new file mode 100644 index 0000000000..6b0fbea78b --- /dev/null +++ b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeNativeMetricSystemDimensionsE2eTest.java @@ -0,0 +1,257 @@ +/* + * 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.observability.storage; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.regex.Matcher; +import java.util.regex.Pattern; +import org.apache.hertzbeat.common.constants.CommonConstants; +import org.apache.hertzbeat.common.constants.MetricDataConstants; +import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; +import org.apache.hertzbeat.warehouse.db.GreptimeQueryGuard; +import org.apache.hertzbeat.warehouse.db.GreptimeSqlQueryExecutor; +import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.GreptimeDbDataStorage; +import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.GreptimeProperties; +import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.NativeMetricSystemContextResolver; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.support.StaticListableBeanFactory; +import org.springframework.web.client.RestTemplate; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +/** Proves the production native metric writer owns entity dimensions as Greptime tags. */ +@Testcontainers +class GreptimeNativeMetricSystemDimensionsE2eTest { + + private static final String GREPTIME_IMAGE = "greptime/greptimedb:v1.1.4"; + private static final int GREPTIME_HTTP_PORT = 4000; + private static final int GREPTIME_GRPC_PORT = 4001; + private static final String TABLE = "linux_native_system_dimensions"; + private static final long COLLECTION_TIME = 1_712_733_600_123L; + private static final Pattern PRIMARY_KEY_PATTERN = Pattern.compile( + "primary\\s+key\\s*\\(([^)]*)\\)", Pattern.CASE_INSENSITIVE | Pattern.DOTALL); + + @Container + @SuppressWarnings("resource") + private static final GenericContainer<?> GREPTIME = new GenericContainer<>(DockerImageName.parse(GREPTIME_IMAGE)) + .withExposedPorts(GREPTIME_HTTP_PORT, GREPTIME_GRPC_PORT) + .withCommand("standalone", "start", + "--http-addr", "0.0.0.0:" + GREPTIME_HTTP_PORT, + "--rpc-bind-addr", "0.0.0.0:" + GREPTIME_GRPC_PORT) + .waitingFor(Wait.forListeningPorts(GREPTIME_HTTP_PORT, GREPTIME_GRPC_PORT)) + .withStartupTimeout(Duration.ofSeconds(120)); + + private GreptimeQueryGuard queryGuard; + private GreptimeDbDataStorage storage; + + @AfterEach + void closeResources() { + if (storage != null) { + storage.destroy(); + } + if (queryGuard != null) { + queryGuard.close(); + } + } + + @Test + void writesAuthoritativeDimensionsAsTagsAndRejectsAllCollectorCollisions() { + GreptimeProperties properties = properties(); + RestTemplate restTemplate = new RestTemplate(); + queryGuard = new GreptimeQueryGuard(4, Duration.ofSeconds(10), Duration.ofMillis(100)); + GreptimeSqlQueryExecutor sql = new GreptimeSqlQueryExecutor(properties, restTemplate, queryGuard); + StaticListableBeanFactory beanFactory = new StaticListableBeanFactory(); + beanFactory.addBean("nativeMetricSystemContextResolver", authoritativeContextResolver()); + storage = new GreptimeDbDataStorage(properties, restTemplate, sql, queryGuard, + beanFactory.getBeanProvider(NativeMetricSystemContextResolver.class)); + + try (CollectRep.MetricsData metricsData = nativeMetrics(); + CollectRep.MetricsData missingAuthority = nativeMetricsWithoutOptionalAuthority()) { + storage.saveData(metricsData); + storage.saveData(missingAuthority); + } + + await().atMost(Duration.ofSeconds(20)).pollInterval(Duration.ofMillis(200)).untilAsserted(() -> + assertThat(sql.executeStrict("SELECT COUNT(*) AS row_count FROM " + TABLE).getFirst()) + .containsEntry("row_count", 3)); + + List<Map<String, Object>> description = sql.executeStrict("DESC " + TABLE); + assertSemanticType(description, "hertzbeat_workspace_id", "tag"); + assertSemanticType(description, "hertzbeat_entity_id", "tag"); + assertSemanticType(description, "hertzbeat_entity_type", "tag"); + assertSemanticType(description, "hertzbeat_monitor_id", "tag"); + assertSemanticType(description, "hertzbeat_collector_id", "tag"); + assertSemanticType(description, "instance", "tag"); + assertSemanticType(description, "ts", "timestamp"); + assertSemanticType(description, "usage", "field"); + + String createTable = sql.executeStrict("SHOW CREATE TABLE " + TABLE).stream() + .flatMap(row -> row.values().stream()) + .map(String::valueOf) + .reduce("", (left, right) -> left + "\n" + right); + Matcher primaryKey = PRIMARY_KEY_PATTERN.matcher(createTable); + assertThat(primaryKey.find()).as(createTable).isTrue(); + String primaryKeyColumns = primaryKey.group(1).toLowerCase(Locale.ROOT); + assertThat(primaryKeyColumns) + .contains("hertzbeat_workspace_id") + .contains("hertzbeat_entity_id") + .contains("hertzbeat_entity_type") + .contains("hertzbeat_monitor_id") + .contains("hertzbeat_collector_id") + .contains("instance") + .doesNotContain("usage"); + + List<Map<String, Object>> rows = sql.executeStrict("SELECT hertzbeat_workspace_id, " + + "hertzbeat_entity_id, hertzbeat_entity_type, hertzbeat_monitor_id, " + + "hertzbeat_collector_id, instance, usage, device FROM " + TABLE + " ORDER BY usage DESC"); + assertThat(rows).hasSize(3); + assertThat(rows.getFirst()) + .containsEntry("hertzbeat_workspace_id", "team-a") + .containsEntry("hertzbeat_entity_id", "99") + .containsEntry("hertzbeat_entity_type", "database") + .containsEntry("hertzbeat_monitor_id", "42") + .containsEntry("hertzbeat_collector_id", "collector-arm-1") + .containsEntry("instance", "db.internal:3306") + .containsEntry("device", "sda"); + assertThat(((Number) rows.getFirst().get("usage")).doubleValue()).isEqualTo(85.5D); + assertThat(rows.get(1)) + .containsEntry("hertzbeat_workspace_id", "team-a") + .containsEntry("hertzbeat_entity_id", "99") + .containsEntry("hertzbeat_monitor_id", "42") + .containsEntry("device", "sdb"); + assertThat(((Number) rows.get(1).get("usage")).doubleValue()).isEqualTo(73.25D); + Map<String, Object> missingAuthority = rows.get(2); + assertThat(missingAuthority) + .containsKeys("hertzbeat_workspace_id", "hertzbeat_entity_id", "hertzbeat_entity_type", + "hertzbeat_collector_id", "instance") + .containsEntry("hertzbeat_monitor_id", "43") + .containsEntry("device", "sdc"); + assertThat(missingAuthority.get("hertzbeat_workspace_id")).isNull(); + assertThat(missingAuthority.get("hertzbeat_entity_id")).isNull(); + assertThat(missingAuthority.get("hertzbeat_entity_type")).isNull(); + assertThat(missingAuthority.get("hertzbeat_collector_id")).isNull(); + assertThat(missingAuthority.get("instance")).isNull(); + assertThat(((Number) missingAuthority.get("usage")).doubleValue()).isEqualTo(12.5D); + } + + private void assertSemanticType(List<Map<String, Object>> description, String column, String semanticType) { + String row = description.stream() + .map(Map::toString) + .filter(value -> value.toLowerCase(Locale.ROOT).contains(column.toLowerCase(Locale.ROOT))) + .findFirst() + .orElseThrow(() -> new AssertionError("Missing Greptime column " + column + ": " + description)); + assertThat(row.toLowerCase(Locale.ROOT)).contains(semanticType); + } + + private GreptimeProperties properties() { + return new GreptimeProperties( + true, + GREPTIME.getHost() + ":" + GREPTIME.getMappedPort(GREPTIME_GRPC_PORT), + "http://" + GREPTIME.getHost() + ":" + GREPTIME.getMappedPort(GREPTIME_HTTP_PORT), + "public", + "", + "", + null); + } + + private CollectRep.MetricsData nativeMetrics() { + List<CollectRep.Field> fields = new ArrayList<>(); + for (String name : List.of( + "hertzbeat_workspace_id", + "hertzbeat_entity_id", + "hertzbeat_entity_type", + "hertzbeat_monitor_id", + "hertzbeat_collector_id", + "instance", + "ts")) { + fields.add(field(name, CommonConstants.TYPE_STRING, true)); + } + fields.add(field("usage", CommonConstants.TYPE_NUMBER, false)); + fields.add(field("device", CommonConstants.TYPE_STRING, true)); + + CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder() + .setId(42L) + .setApp("linux") + .setMetrics("native_system_dimensions") + .setTime(COLLECTION_TIME) + .setCode(CollectRep.Code.SUCCESS) + .addMetadata("hertzbeat.workspace.id", "spoof-workspace") + .addMetadata(MetricDataConstants.ENTITY_ID, "999") + .addMetadata("hertzbeat.entity.type", "spoof-type") + .addMetadata(MetricDataConstants.COLLECTOR_ID, "collector-arm-1") + .addMetadata(MetricDataConstants.INSTANCE, "db.internal:3306"); + builder.addAllFields(fields); + builder.addValueRow(row("spoof-workspace", "spoof-entity", "spoof-type", "spoof-monitor", + "spoof-collector", "spoof-instance", "spoof-ts", "85.5", "sda")); + builder.addValueRow(row("spoof-workspace-2", "spoof-entity-2", "spoof-type-2", "spoof-monitor-2", + "spoof-collector-2", "spoof-instance-2", "spoof-ts-2", "73.25", "sdb")); + return builder.build(); + } + + private CollectRep.MetricsData nativeMetricsWithoutOptionalAuthority() { + CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder() + .setId(43L) + .setApp("linux") + .setMetrics("native_system_dimensions") + .setTime(COLLECTION_TIME + 1) + .setCode(CollectRep.Code.SUCCESS) + .addMetadata("hertzbeat.workspace.id", "spoof-unbound-workspace") + .addMetadata(MetricDataConstants.ENTITY_ID, "1000") + .addMetadata("hertzbeat.entity.type", "spoof-unbound-type"); + builder.addAllFields(List.of( + field("usage", CommonConstants.TYPE_NUMBER, false), + field("device", CommonConstants.TYPE_STRING, true))); + builder.addValueRow(row("12.5", "sdc")); + return builder.build(); + } + + private NativeMetricSystemContextResolver authoritativeContextResolver() { + return (metricsData, intrinsic) -> { + if (intrinsic == null || !Long.valueOf(42L).equals(intrinsic.monitorId())) { + return intrinsic; + } + return new NativeMetricSystemContext( + "team-a", + 99L, + "database", + intrinsic.monitorId(), + intrinsic.collectorId(), + intrinsic.instance()); + }; + } + + private CollectRep.Field field(String name, int type, boolean label) { + return CollectRep.Field.newBuilder().setName(name).setType(type).setLabel(label).build(); + } + + private CollectRep.ValueRow row(String... columns) { + return CollectRep.ValueRow.newBuilder().setColumns(List.of(columns)).build(); + } +} diff --git a/hertzbeat-manager/pom.xml b/hertzbeat-manager/pom.xml index 82ac33eb7b..766b601757 100644 --- a/hertzbeat-manager/pom.xml +++ b/hertzbeat-manager/pom.xml @@ -50,6 +50,10 @@ <groupId>org.apache.hertzbeat</groupId> <artifactId>hertzbeat-warehouse</artifactId> </dependency> + <dependency> + <groupId>com.github.ben-manes.caffeine</groupId> + <artifactId>caffeine</artifactId> + </dependency> <!-- alerter --> <dependency> <groupId>org.apache.hertzbeat</groupId> diff --git a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolver.java b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolver.java new file mode 100644 index 0000000000..1bd0cae14e --- /dev/null +++ b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolver.java @@ -0,0 +1,114 @@ +/* + * 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.manager.service.entity; + +import com.github.benmanes.caffeine.cache.Cache; +import com.github.benmanes.caffeine.cache.Caffeine; +import java.time.Duration; +import java.util.List; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicLong; +import lombok.extern.slf4j.Slf4j; +import org.apache.hertzbeat.common.entity.manager.EntityMonitorBind; +import org.apache.hertzbeat.common.entity.manager.ObserveEntity; +import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; +import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.NativeMetricSystemContextResolver; +import org.springframework.stereotype.Service; + +/** Resolves the single active monitor-to-entity authority for native metric persistence. */ +@Slf4j +@Service +public class ManagerNativeMetricSystemContextResolver implements NativeMetricSystemContextResolver { + + private static final String ACTIVE = "active"; + private static final long MAXIMUM_MONITOR_AUTHORITIES = 100_000; + private static final Duration AUTHORITY_REFRESH_INTERVAL = Duration.ofSeconds(5); + private static final AtomicLong AUTHORITY_LOOKUP_FAILURE_COUNT = new AtomicLong(); + + private final EntityMonitorBindQueryService bindQueryService; + private final EntityWorkspaceQueryService entityQueryService; + private final Cache<Long, EntityAuthority> authorityCache = Caffeine.newBuilder() + .maximumSize(MAXIMUM_MONITOR_AUTHORITIES) + .expireAfterWrite(AUTHORITY_REFRESH_INTERVAL) + .build(); + + public ManagerNativeMetricSystemContextResolver(EntityMonitorBindQueryService bindQueryService, + EntityWorkspaceQueryService entityQueryService) { + this.bindQueryService = bindQueryService; + this.entityQueryService = entityQueryService; + } + + @Override + public NativeMetricSystemContext resolve( + CollectRep.MetricsData metricsData, NativeMetricSystemContext intrinsic) { + if (intrinsic == null || intrinsic.monitorId() == null) { + return intrinsic; + } + EntityAuthority authority = authorityCache.get(intrinsic.monitorId(), this::loadAuthority); + if (authority == null || authority.entityId() == null) { + return intrinsic; + } + return new NativeMetricSystemContext( + authority.workspaceId(), + authority.entityId(), + authority.entityType(), + intrinsic.monitorId(), + intrinsic.collectorId(), + intrinsic.instance()); + } + + private EntityAuthority loadAuthority(Long monitorId) { + try { + return loadAuthorityFromMetadata(monitorId); + } catch (RuntimeException exception) { + long failureCount = AUTHORITY_LOOKUP_FAILURE_COUNT.incrementAndGet(); + if (Long.bitCount(failureCount) == 1) { + log.warn("Native metric authority lookup failed for monitor {}; failures={}; exception={}", + monitorId, failureCount, exception.getClass().getName()); + } + return EntityAuthority.ABSENT; + } + } + + private EntityAuthority loadAuthorityFromMetadata(Long monitorId) { + List<EntityMonitorBind> bindings = bindQueryService.findMonitorBindsByMonitorId(monitorId); + List<EntityMonitorBind> activeBindings = bindings == null + ? List.of() + : bindings.stream() + .filter(Objects::nonNull) + .filter(binding -> Objects.equals(monitorId, binding.getMonitorId())) + .filter(binding -> binding.getEntityId() != null) + .filter(binding -> ACTIVE.equalsIgnoreCase(binding.getStatus())) + .toList(); + if (activeBindings.size() != 1) { + return EntityAuthority.ABSENT; + } + EntityMonitorBind binding = activeBindings.getFirst(); + ObserveEntity entity = entityQueryService.findEntityById(binding.getEntityId()).orElse(null); + if (entity == null || !Objects.equals(binding.getEntityId(), entity.getId())) { + return EntityAuthority.ABSENT; + } + return new EntityAuthority(entity.getWorkspaceId(), entity.getId(), entity.getType()); + } + + private record EntityAuthority(String workspaceId, Long entityId, String entityType) { + + private static final EntityAuthority ABSENT = new EntityAuthority(null, null, null); + } +} diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolverTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolverTest.java new file mode 100644 index 0000000000..c99b5f4dcd --- /dev/null +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/service/entity/ManagerNativeMetricSystemContextResolverTest.java @@ -0,0 +1,132 @@ +/* + * 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.manager.service.entity; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.List; +import java.util.Optional; +import org.apache.hertzbeat.common.entity.manager.EntityMonitorBind; +import org.apache.hertzbeat.common.entity.manager.ObserveEntity; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class ManagerNativeMetricSystemContextResolverTest { + + private EntityMonitorBindQueryService bindQueryService; + private EntityWorkspaceQueryService entityQueryService; + private ManagerNativeMetricSystemContextResolver resolver; + + @BeforeEach + void setUp() { + bindQueryService = mock(EntityMonitorBindQueryService.class); + entityQueryService = mock(EntityWorkspaceQueryService.class); + resolver = new ManagerNativeMetricSystemContextResolver(bindQueryService, entityQueryService); + } + + @Test + void enrichesCollectorContextFromSingleActiveEntityAuthority() { + NativeMetricSystemContext intrinsic = new NativeMetricSystemContext( + null, null, null, 42L, "collector-arm-1", "db.internal:3306"); + when(bindQueryService.findMonitorBindsByMonitorId(42L)).thenReturn(List.of( + EntityMonitorBind.builder() + .monitorId(42L) + .entityId(99L) + .status("active") + .build())); + when(entityQueryService.findEntityById(99L)).thenReturn(Optional.of( + ObserveEntity.builder() + .id(99L) + .workspaceId("team-a") + .type("database") + .build())); + + NativeMetricSystemContext result = resolver.resolve(null, intrinsic); + + assertEquals("team-a", result.workspaceId()); + assertEquals(99L, result.entityId()); + assertEquals("database", result.entityType()); + assertEquals(42L, result.monitorId()); + assertEquals("collector-arm-1", result.collectorId()); + assertEquals("db.internal:3306", result.instance()); + + assertEquals(result, resolver.resolve(null, intrinsic)); + verify(bindQueryService, times(1)).findMonitorBindsByMonitorId(42L); + verify(entityQueryService, times(1)).findEntityById(99L); + } + + @Test + void ignoresInactiveHistoryWhenResolvingSingleActiveEntityAuthority() { + NativeMetricSystemContext intrinsic = new NativeMetricSystemContext( + null, null, null, 42L, "collector-arm-1", "db.internal:3306"); + when(bindQueryService.findMonitorBindsByMonitorId(42L)).thenReturn(List.of( + EntityMonitorBind.builder().monitorId(42L).entityId(98L).status("inactive").build(), + EntityMonitorBind.builder().monitorId(42L).entityId(99L).status("active").build())); + when(entityQueryService.findEntityById(99L)).thenReturn(Optional.of( + ObserveEntity.builder() + .id(99L) + .workspaceId("team-a") + .type("database") + .build())); + + NativeMetricSystemContext result = resolver.resolve(null, intrinsic); + + assertEquals("team-a", result.workspaceId()); + assertEquals(99L, result.entityId()); + assertEquals("database", result.entityType()); + assertEquals(42L, result.monitorId()); + assertEquals("collector-arm-1", result.collectorId()); + assertEquals("db.internal:3306", result.instance()); + verify(entityQueryService).findEntityById(99L); + } + + @Test + void cachesAbsentAuthorityWhenBindingQueryFails() { + NativeMetricSystemContext intrinsic = new NativeMetricSystemContext( + null, null, null, 42L, "collector-arm-1", "db.internal:3306"); + when(bindQueryService.findMonitorBindsByMonitorId(42L)) + .thenThrow(new IllegalStateException("metadata unavailable")); + + assertEquals(intrinsic, resolver.resolve(null, intrinsic)); + assertEquals(intrinsic, resolver.resolve(null, intrinsic)); + verify(bindQueryService, times(1)).findMonitorBindsByMonitorId(42L); + } + + @Test + void leavesOptionalEntityAuthorityAbsentWhenBindingIsAmbiguous() { + NativeMetricSystemContext intrinsic = new NativeMetricSystemContext( + null, null, null, 42L, null, "db.internal:3306"); + when(bindQueryService.findMonitorBindsByMonitorId(42L)).thenReturn(List.of( + EntityMonitorBind.builder().monitorId(42L).entityId(99L).status("active").build(), + EntityMonitorBind.builder().monitorId(42L).entityId(100L).status("active").build())); + + NativeMetricSystemContext result = resolver.resolve(null, intrinsic); + + assertNull(result.workspaceId()); + assertNull(result.entityId()); + assertNull(result.entityType()); + assertNull(result.collectorId()); + assertEquals(42L, result.monitorId()); + } +} 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 9762bcb983..9538315a18 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 @@ -67,6 +67,7 @@ import org.apache.hertzbeat.common.entity.dto.Value; import org.apache.hertzbeat.common.entity.event.CollectionExecutionEvent; import org.apache.hertzbeat.common.entity.log.LogEntry; import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.common.runtime.ConditionalOnNormalBusinessRuntime; import org.apache.hertzbeat.common.util.Base64Util; @@ -79,6 +80,7 @@ import org.apache.hertzbeat.warehouse.store.history.tsdb.AbstractHistoryDataStor import org.apache.hertzbeat.warehouse.store.history.tsdb.HistoryDataReader.ServerAvailability; import org.apache.hertzbeat.warehouse.store.history.tsdb.vm.PromQlQueryContent; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.http.HttpEntity; @@ -147,21 +149,42 @@ public class GreptimeDbDataStorage extends AbstractHistoryDataStorage { private final GreptimeSqlQueryExecutor greptimeSqlQueryExecutor; private final GreptimeQueryGuard queryGuard; private final GreptimeServerAvailabilityProbe serverAvailabilityProbe; + private final NativeMetricSystemContextResolver nativeMetricSystemContextResolver; - @Autowired public GreptimeDbDataStorage( GreptimeProperties greptimeProperties, @Qualifier(WarehouseConstants.GREPTIME_QUERY_REST_TEMPLATE) RestTemplate restTemplate, GreptimeSqlQueryExecutor greptimeSqlQueryExecutor, GreptimeQueryGuard queryGuard) { this(greptimeProperties, restTemplate, greptimeSqlQueryExecutor, queryGuard, - createServerAvailabilityProbe(greptimeProperties)); + createServerAvailabilityProbe(greptimeProperties), (metricsData, intrinsic) -> intrinsic); + } + + @Autowired + public GreptimeDbDataStorage( + GreptimeProperties greptimeProperties, + @Qualifier(WarehouseConstants.GREPTIME_QUERY_REST_TEMPLATE) RestTemplate restTemplate, + GreptimeSqlQueryExecutor greptimeSqlQueryExecutor, + GreptimeQueryGuard queryGuard, + ObjectProvider<NativeMetricSystemContextResolver> nativeMetricSystemContextResolverProvider) { + this(greptimeProperties, restTemplate, greptimeSqlQueryExecutor, queryGuard, + createServerAvailabilityProbe(greptimeProperties), + nativeMetricSystemContextResolverProvider.getIfAvailable(() -> (metricsData, intrinsic) -> intrinsic)); } GreptimeDbDataStorage(GreptimeProperties greptimeProperties, RestTemplate restTemplate, GreptimeSqlQueryExecutor greptimeSqlQueryExecutor, GreptimeQueryGuard queryGuard, GreptimeServerAvailabilityProbe serverAvailabilityProbe) { + this(greptimeProperties, restTemplate, greptimeSqlQueryExecutor, queryGuard, serverAvailabilityProbe, + (metricsData, intrinsic) -> intrinsic); + } + + GreptimeDbDataStorage(GreptimeProperties greptimeProperties, RestTemplate restTemplate, + GreptimeSqlQueryExecutor greptimeSqlQueryExecutor, + GreptimeQueryGuard queryGuard, + GreptimeServerAvailabilityProbe serverAvailabilityProbe, + NativeMetricSystemContextResolver nativeMetricSystemContextResolver) { if (greptimeProperties == null) { log.error("init error, please config Warehouse GreptimeDB props in application.yml"); throw new IllegalArgumentException("please config Warehouse GreptimeDB props"); @@ -171,6 +194,7 @@ public class GreptimeDbDataStorage extends AbstractHistoryDataStorage { this.greptimeSqlQueryExecutor = greptimeSqlQueryExecutor; this.queryGuard = Objects.requireNonNull(queryGuard); this.serverAvailabilityProbe = Objects.requireNonNull(serverAvailabilityProbe); + this.nativeMetricSystemContextResolver = Objects.requireNonNull(nativeMetricSystemContextResolver); serverAvailable = initGreptimeDbClient(greptimeProperties); if (serverAvailable) { applyDatabaseTtlIfConfigured(greptimeProperties); @@ -253,36 +277,35 @@ public class GreptimeDbDataStorage extends AbstractHistoryDataStorage { metricsData.getId(), metricsData.getMetrics()); return; } - String instance = metricsData.getInstance(); String app = metricsData.getApp(); String tableName = getTableName(app, metricsData.getMetrics()); List<CollectRep.Field> fields = metricsData.getFields(); - Table table = Table.from(metricSchemaCache.getOrCreate(tableName, fields)); + GreptimeMetricSchemaCache.ResolvedSchema resolvedSchema = metricSchemaCache.resolve(tableName, fields); + Table table = Table.from(resolvedSchema.schema()); + if (resolvedSchema.schemaChanged() && !resolvedSchema.rejectedNames().isEmpty()) { + log.warn("[warehouse greptime] ignored Collector fields that collide with system dimensions: {}", + resolvedSchema.rejectedNames()); + } long collectionTime = metricsData.getTime(); long timestamp = collectionTime > 0 ? collectionTime : System.currentTimeMillis(); - Object[] values = new Object[2 + fields.size()]; - values[0] = instance; - values[1] = timestamp; + Object[] systemTagValues = resolveNativeMetricSystemContext(metricsData).tagValues(); + int fieldOffset = systemTagValues.length + 1; RowWrapper rowWrapper = metricsData.readRow(); while (rowWrapper.hasNextRow()) { rowWrapper = rowWrapper.nextRow(); - int index = 0; + Object[] values = new Object[fieldOffset + resolvedSchema.sourceIndexes().size()]; + System.arraycopy(systemTagValues, 0, values, 0, systemTagValues.length); + values[systemTagValues.length] = timestamp; + int sourceIndex = 0; + int acceptedIndex = 0; while (rowWrapper.hasNextCell()) { ArrowCell cell = rowWrapper.nextCell(); - if (CommonConstants.NULL_VALUE.equals(cell.getValue())) { - values[2 + index] = null; - } else { - Boolean label = cell.getMetadataAsBoolean(MetricDataConstants.LABEL); - Byte type = cell.getMetadataAsByte(MetricDataConstants.TYPE); - if (label) { - values[2 + index] = cell.getValue(); - } else if (type == CommonConstants.TYPE_NUMBER) { - values[2 + index] = Double.parseDouble(cell.getValue()); - } else if (type == CommonConstants.TYPE_STRING) { - values[2 + index] = cell.getValue(); - } + if (acceptedIndex < resolvedSchema.sourceIndexes().size() + && resolvedSchema.sourceIndexes().get(acceptedIndex) == sourceIndex) { + values[fieldOffset + acceptedIndex] = metricCellValue(cell); + acceptedIndex++; } - index++; + sourceIndex++; } table.addRow(values); @@ -301,6 +324,33 @@ public class GreptimeDbDataStorage extends AbstractHistoryDataStorage { } } + private NativeMetricSystemContext resolveNativeMetricSystemContext(CollectRep.MetricsData metricsData) { + NativeMetricSystemContext intrinsic = NativeMetricSystemContext.from(metricsData); + try { + NativeMetricSystemContext resolved = nativeMetricSystemContextResolver.resolve(metricsData, intrinsic); + return resolved == null ? intrinsic : resolved; + } catch (RuntimeException exception) { + log.warn("[warehouse greptime] native metric entity authority unavailable for monitor {}: {}", + metricsData.getId(), exception.getClass().getSimpleName()); + return intrinsic; + } + } + + private Object metricCellValue(ArrowCell cell) { + if (CommonConstants.NULL_VALUE.equals(cell.getValue())) { + return null; + } + Boolean label = cell.getMetadataAsBoolean(MetricDataConstants.LABEL); + Byte type = cell.getMetadataAsByte(MetricDataConstants.TYPE); + if (Boolean.TRUE.equals(label) || type != null && type == CommonConstants.TYPE_STRING) { + return cell.getValue(); + } + if (type != null && type == CommonConstants.TYPE_NUMBER) { + return Double.parseDouble(cell.getValue()); + } + return null; + } + @Override public boolean supportsCollectionExecutionEvents() { return true; diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCache.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCache.java index 430ac12576..895fd1faae 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCache.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCache.java @@ -21,8 +21,11 @@ import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; import io.greptime.models.DataType; import io.greptime.models.TableSchema; +import java.util.ArrayList; import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.hertzbeat.common.constants.CommonConstants; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemDimensions; import org.apache.hertzbeat.common.entity.message.CollectRep; /** @@ -45,26 +48,52 @@ final class GreptimeMetricSchemaCache { } TableSchema getOrCreate(String tableName, List<CollectRep.Field> fields) { + return resolve(tableName, fields).schema(); + } + + ResolvedSchema resolve(String tableName, List<CollectRep.Field> fields) { + List<IndexedMetricColumn> accepted = new ArrayList<>(fields.size()); + List<String> rejectedNames = new ArrayList<>(); + for (int index = 0; index < fields.size(); index++) { + CollectRep.Field field = fields.get(index); + if (field == null) { + continue; + } + if (NativeMetricSystemDimensions.isReserved(field.getName())) { + rejectedNames.add(field.getName()); + continue; + } + if (!isSupported(field)) { + continue; + } + accepted.add(new IndexedMetricColumn(index, MetricColumn.from(field))); + } + List<MetricColumn> columns = accepted.stream().map(IndexedMetricColumn::column).toList(); CachedSchema cached = schemas.getIfPresent(tableName); - if (cached != null && cached.matches(fields)) { - return cached.schema(); + if (cached != null && cached.matches(columns)) { + return resolved(cached.schema(), accepted, rejectedNames, false); } + AtomicBoolean schemaChanged = new AtomicBoolean(); CachedSchema resolved = schemas.asMap().compute(tableName, (key, current) -> { - if (current != null && current.matches(fields)) { + if (current != null && current.matches(columns)) { return current; } - return createSchema(key, fields); + schemaChanged.set(true); + return createSchema(key, columns); }); - return resolved.schema(); + return resolved(resolved.schema(), accepted, rejectedNames, schemaChanged.get()); + } + + private static ResolvedSchema resolved(TableSchema schema, List<IndexedMetricColumn> columns, + List<String> rejectedNames, boolean schemaChanged) { + return new ResolvedSchema(schema, columns.stream().map(IndexedMetricColumn::sourceIndex).toList(), + List.copyOf(rejectedNames), schemaChanged); } - private static CachedSchema createSchema(String tableName, List<CollectRep.Field> fields) { - List<MetricColumn> columns = fields.stream() - .map(MetricColumn::from) - .toList(); - TableSchema.Builder builder = TableSchema.newBuilder(tableName) - .addTag("instance", DataType.String) - .addTimestamp("ts", DataType.TimestampMillisecond); + private static CachedSchema createSchema(String tableName, List<MetricColumn> columns) { + TableSchema.Builder builder = TableSchema.newBuilder(tableName); + NativeMetricSystemDimensions.TAG_NAMES.forEach(name -> builder.addTag(name, DataType.String)); + builder.addTimestamp(NativeMetricSystemDimensions.TIMESTAMP, DataType.TimestampMillisecond); for (MetricColumn column : columns) { if (column.label()) { builder.addTag(column.name(), DataType.String); @@ -77,31 +106,30 @@ final class GreptimeMetricSchemaCache { return new CachedSchema(columns, builder.build()); } + private static boolean isSupported(CollectRep.Field field) { + return field.getLabel() + || field.getType() == CommonConstants.TYPE_NUMBER + || field.getType() == CommonConstants.TYPE_STRING; + } + private record CachedSchema(List<MetricColumn> columns, TableSchema schema) { - private boolean matches(List<CollectRep.Field> fields) { - if (columns.size() != fields.size()) { - return false; - } - for (int index = 0; index < fields.size(); index++) { - if (!columns.get(index).matches(fields.get(index))) { - return false; - } - } - return true; + private boolean matches(List<MetricColumn> candidateColumns) { + return columns.equals(candidateColumns); } } + record ResolvedSchema(TableSchema schema, List<Integer> sourceIndexes, List<String> rejectedNames, + boolean schemaChanged) { + } + + private record IndexedMetricColumn(int sourceIndex, MetricColumn column) { + } + private record MetricColumn(String name, int type, boolean label) { private static MetricColumn from(CollectRep.Field field) { return new MetricColumn(field.getName(), field.getType(), field.getLabel()); } - - private boolean matches(CollectRep.Field field) { - return name.equals(field.getName()) - && type == field.getType() - && label == field.getLabel(); - } } } diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/NativeMetricSystemContextResolver.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/NativeMetricSystemContextResolver.java new file mode 100644 index 0000000000..a1a373eb20 --- /dev/null +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/NativeMetricSystemContextResolver.java @@ -0,0 +1,36 @@ +/* + * 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.greptime; + +import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; + +/** Resolves manager-owned entity authority immediately before one native metric write. */ +@FunctionalInterface +public interface NativeMetricSystemContextResolver { + + /** + * Resolve manager-owned dimensions without replacing Collector-owned values. + * + * @param metricsData source metric sample + * @param intrinsic Collector-owned context + * @return resolved system context + */ + NativeMetricSystemContext resolve( + CollectRep.MetricsData metricsData, NativeMetricSystemContext intrinsic); +} 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 984a24b773..d95bb80cf3 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,6 +38,7 @@ import io.greptime.models.Err; import io.greptime.models.Result; import io.greptime.models.Table; import io.greptime.models.WriteOk; +import io.greptime.v1.Common.SemanticType; import io.greptime.v1.RowData; import java.net.URI; import java.time.Duration; @@ -55,6 +56,7 @@ import org.apache.hertzbeat.common.entity.dto.Value; import org.apache.hertzbeat.common.entity.event.CollectionExecutionEvent; import org.apache.hertzbeat.common.entity.log.LogEntry; import org.apache.hertzbeat.common.entity.message.CollectRep; +import org.apache.hertzbeat.common.entity.metric.NativeMetricSystemContext; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.warehouse.constants.WarehouseConstants; import org.apache.hertzbeat.warehouse.db.GreptimeQueryGuard; @@ -237,10 +239,10 @@ class GreptimeDbDataStorageTest { Table writtenTable = tableCaptor.getValue(); List<RowData.Value> values = writtenTable.intoRowInsertRequest() .getRows().getRows(0).getValuesList(); - assertEquals("server-1", values.get(0).getStringValue()); - assertEquals(1_712_733_600_123L, values.get(1).getTimestampMillisecondValue()); - assertEquals(85.5, values.get(2).getF64Value()); - assertEquals("server1", values.get(3).getStringValue()); + assertEquals("1", values.get(3).getStringValue()); + assertEquals("server-1", values.get(5).getStringValue()); + assertEquals(1_712_733_600_123L, values.get(6).getTimestampMillisecondValue()); + assertEquals(85.5, values.get(7).getF64Value()); verify(metricsData, never()).getValues(); verify(metricsData.readRow(), never()).cellStream(); @@ -260,6 +262,69 @@ class GreptimeDbDataStorageTest { } } + @Test + void writesFixedSystemTagsAndDropsEveryCollidingCollectorFieldAcrossRows() { + try (MockedStatic<GreptimeDB> mockedStatic = mockStatic(GreptimeDB.class)) { + mockedStatic.when(() -> GreptimeDB.create(any())).thenReturn(greptimeDb); + @SuppressWarnings("unchecked") + Result<WriteOk, Err> mockResult = mock(Result.class); + when(mockResult.isOk()).thenReturn(true); + when(greptimeDb.write(any(Table.class))) + .thenReturn(CompletableFuture.completedFuture(mockResult)); + NativeMetricSystemContext context = new NativeMetricSystemContext( + "team-a", 99L, "database", 42L, null, "db.internal:3306"); + greptimeDbDataStorage = new GreptimeDbDataStorage( + greptimeProperties, + restTemplate, + greptimeSqlQueryExecutor, + queryGuard, + serverAvailabilityProbe, + (metricsData, intrinsic) -> context); + CollectRep.MetricsData metricsData = collidingSystemDimensionMetrics(); + + greptimeDbDataStorage.saveData(metricsData); + + ArgumentCaptor<Table> tableCaptor = ArgumentCaptor.forClass(Table.class); + verify(greptimeDb).write(tableCaptor.capture()); + var rows = tableCaptor.getValue().intoRowInsertRequest().getRows(); + assertEquals(List.of( + "hertzbeat_workspace_id", + "hertzbeat_entity_id", + "hertzbeat_entity_type", + "hertzbeat_monitor_id", + "hertzbeat_collector_id", + "instance", + "ts", + "usage", + "device"), rows.getSchemaList().stream().map(RowData.ColumnSchema::getColumnName).toList()); + assertEquals(List.of( + SemanticType.TAG, + SemanticType.TAG, + SemanticType.TAG, + SemanticType.TAG, + SemanticType.TAG, + SemanticType.TAG), + rows.getSchemaList().subList(0, 6).stream() + .map(RowData.ColumnSchema::getSemanticType) + .toList()); + assertEquals(2, rows.getRowsCount()); + List<RowData.Value> first = rows.getRows(0).getValuesList(); + assertEquals("team-a", first.get(0).getStringValue()); + assertEquals("99", first.get(1).getStringValue()); + assertEquals("database", first.get(2).getStringValue()); + assertEquals("42", first.get(3).getStringValue()); + assertEquals(RowData.Value.ValueDataCase.VALUEDATA_NOT_SET, first.get(4).getValueDataCase()); + assertEquals("db.internal:3306", first.get(5).getStringValue()); + assertEquals(1_712_733_600_123L, first.get(6).getTimestampMillisecondValue()); + assertEquals(85.5, first.get(7).getF64Value()); + assertEquals("sda", first.get(8).getStringValue()); + List<RowData.Value> second = rows.getRows(1).getValuesList(); + assertEquals("team-a", second.get(0).getStringValue()); + assertEquals(73.25, second.get(7).getF64Value()); + assertEquals("sdb", second.get(8).getStringValue()); + } + } + @Test void testSaveCollectionExecutionEventsAsOneBatch() { try (MockedStatic<GreptimeDB> mockedStatic = mockStatic(GreptimeDB.class)) { @@ -336,6 +401,31 @@ class GreptimeDbDataStorageTest { assertEquals("85.5", result.values().iterator().next().get(0).getOrigin()); } + @Test + void testGetHistoryMetricDataPreservesEntitySystemDimensions() { + greptimeDbDataStorage = new GreptimeDbDataStorage( + greptimeProperties, restTemplate, greptimeSqlQueryExecutor, queryGuard); + PromQlQueryContent content = createMockPromQlQueryContent(); + Map<String, String> metric = content.getData().getResult().getFirst().getMetric(); + metric.put("hertzbeat_workspace_id", "team-a"); + metric.put("hertzbeat_entity_id", "99"); + metric.put("hertzbeat_entity_type", "database"); + metric.put("hertzbeat_monitor_id", "42"); + metric.put("hertzbeat_collector_id", "collector-arm-1"); + when(restTemplate.exchange(any(), eq(HttpMethod.GET), any(HttpEntity.class), + eq(PromQlQueryContent.class))).thenReturn(new ResponseEntity<>(content, HttpStatus.OK)); + + Map<String, List<Value>> result = greptimeDbDataStorage.getHistoryMetricData( + "db.internal:3306", "linux", "cpu", "usage", "6h"); + + String seriesKey = result.keySet().iterator().next(); + assertTrue(seriesKey.contains("\"hertzbeat_workspace_id\":\"team-a\"")); + assertTrue(seriesKey.contains("\"hertzbeat_entity_id\":\"99\"")); + assertTrue(seriesKey.contains("\"hertzbeat_entity_type\":\"database\"")); + assertTrue(seriesKey.contains("\"hertzbeat_monitor_id\":\"42\"")); + assertTrue(seriesKey.contains("\"hertzbeat_collector_id\":\"collector-arm-1\"")); + } + @Test void testGetHistoryMetricDataUsesAbsoluteRangeAndStep() { greptimeDbDataStorage = new GreptimeDbDataStorage(greptimeProperties, restTemplate, greptimeSqlQueryExecutor, queryGuard); @@ -1073,6 +1163,66 @@ class GreptimeDbDataStorageTest { return mockMetricsData; } + private CollectRep.MetricsData collidingSystemDimensionMetrics() { + List<CollectRep.Field> fields = new ArrayList<>(); + for (String name : List.of( + "hertzbeat_workspace_id", + "hertzbeat_entity_id", + "hertzbeat_entity_type", + "hertzbeat_monitor_id", + "hertzbeat_collector_id", + "instance", + "ts")) { + fields.add(CollectRep.Field.newBuilder() + .setName(name) + .setType(CommonConstants.TYPE_STRING) + .setLabel(true) + .build()); + } + fields.add(CollectRep.Field.newBuilder() + .setName("usage") + .setType(CommonConstants.TYPE_NUMBER) + .setLabel(false) + .build()); + fields.add(CollectRep.Field.newBuilder() + .setName("device") + .setType(CommonConstants.TYPE_STRING) + .setLabel(true) + .build()); + CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder() + .setId(42L) + .setApp("linux") + .setMetrics("cpu") + .setTime(1_712_733_600_123L) + .setCode(CollectRep.Code.SUCCESS); + builder.addAllFields(fields); + builder.addValueRow(CollectRep.ValueRow.newBuilder() + .setColumns(List.of( + "spoof-workspace", + "spoof-entity", + "spoof-type", + "spoof-monitor", + "spoof-collector", + "spoof-instance", + "spoof-ts", + "85.5", + "sda")) + .build()); + builder.addValueRow(CollectRep.ValueRow.newBuilder() + .setColumns(List.of( + "spoof-workspace-2", + "spoof-entity-2", + "spoof-type-2", + "spoof-monitor-2", + "spoof-collector-2", + "spoof-instance-2", + "spoof-ts-2", + "73.25", + "sdb")) + .build()); + return builder.build(); + } + private PromQlQueryContent createMockPromQlQueryContent() { PromQlQueryContent content = new PromQlQueryContent(); PromQlQueryContent.ContentData data = new PromQlQueryContent.ContentData(); diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCacheTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCacheTest.java index b2d7d4898d..94ab4add96 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCacheTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/store/history/tsdb/greptime/GreptimeMetricSchemaCacheTest.java @@ -19,6 +19,9 @@ package org.apache.hertzbeat.warehouse.store.history.tsdb.greptime; import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import io.greptime.models.TableSchema; import java.util.List; @@ -32,6 +35,34 @@ import org.junit.jupiter.api.Test; class GreptimeMetricSchemaCacheTest { + @Test + void reservesSystemDimensionColumnsAndRetainsSourceAlignment() { + GreptimeMetricSchemaCache cache = new GreptimeMetricSchemaCache(8); + List<CollectRep.Field> fields = List.of( + field("hertzbeat_workspace_id", CommonConstants.TYPE_STRING, true), + field("usage", CommonConstants.TYPE_NUMBER, false), + field("INSTANCE", CommonConstants.TYPE_STRING, true), + field("unsupported", Integer.MAX_VALUE, false), + field("device", CommonConstants.TYPE_STRING, true)); + + GreptimeMetricSchemaCache.ResolvedSchema resolved = cache.resolve("linux_cpu", fields); + + assertEquals(List.of(1, 4), resolved.sourceIndexes()); + assertEquals(List.of("hertzbeat_workspace_id", "INSTANCE"), resolved.rejectedNames()); + assertTrue(resolved.schemaChanged()); + assertEquals(List.of( + "hertzbeat_workspace_id", + "hertzbeat_entity_id", + "hertzbeat_entity_type", + "hertzbeat_monitor_id", + "hertzbeat_collector_id", + "instance", + "ts", + "usage", + "device"), resolved.schema().getColumnNames()); + assertFalse(cache.resolve("linux_cpu", fields).schemaChanged()); + } + @Test void reusesEquivalentSchemaAndRebuildsAfterEvolution() { GreptimeMetricSchemaCache cache = new GreptimeMetricSchemaCache(8); @@ -46,6 +77,22 @@ class GreptimeMetricSchemaCacheTest { assertNotSame(initial, cache.getOrCreate("linux_cpu", List.of(changedUsage, device))); } + @Test + void recalculatesSourceIndexesWhenEquivalentSchemaHasCollisionsInDifferentPositions() { + GreptimeMetricSchemaCache cache = new GreptimeMetricSchemaCache(8); + CollectRep.Field usage = field("usage", CommonConstants.TYPE_NUMBER, false); + CollectRep.Field device = field("device", CommonConstants.TYPE_STRING, true); + + GreptimeMetricSchemaCache.ResolvedSchema first = cache.resolve("linux_cpu", List.of( + field("instance", CommonConstants.TYPE_STRING, true), usage, device)); + GreptimeMetricSchemaCache.ResolvedSchema second = cache.resolve("linux_cpu", List.of( + usage, field("HERTZBEAT_ENTITY_ID", CommonConstants.TYPE_STRING, true), device)); + + assertSame(first.schema(), second.schema()); + assertEquals(List.of(1, 2), first.sourceIndexes()); + assertEquals(List.of(0, 2), second.sourceIndexes()); + } + @Test void keepsSchemasForDifferentTablesIndependent() { GreptimeMetricSchemaCache cache = new GreptimeMetricSchemaCache(8); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
