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 21aa26e418524a24beaee6e397669f5f0e403745 Author: Logic <[email protected]> AuthorDate: Fri Aug 28 10:07:50 2026 +0800 Rework native trace storage contract --- .../alert/service/DataSourceServiceTest.java | 4 +- .../storage/GreptimeTraceApmFlowProofE2eTest.java | 241 ++++----- .../forwarder/GreptimeApmFlowInitializer.java | 108 +--- .../service/impl/EntityTraceQueryServiceImpl.java | 60 +-- .../greptime/flows/hertzbeat_apm_red_1m.sql | 14 +- .../main/resources/greptime/tables/hzb_traces.sql | 15 +- .../forwarder/GreptimeApmFlowInitializerTest.java | 41 +- .../GreptimeTraceTableInitializerTest.java | 71 ++- .../impl/EntityTraceQueryServiceImplTest.java | 53 +- .../repository/GreptimeTraceQueryRepository.java | 181 ++++--- .../GreptimeTraceQueryRepositoryTest.java | 559 ++++++++++++--------- script/dev/seed-trace-rich-demo.sh | 165 +++--- 12 files changed, 802 insertions(+), 710 deletions(-) diff --git a/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/service/DataSourceServiceTest.java b/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/service/DataSourceServiceTest.java index fad8bed326..d445904754 100644 --- a/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/service/DataSourceServiceTest.java +++ b/hertzbeat-alerter/src/test/java/org/apache/hertzbeat/alert/service/DataSourceServiceTest.java @@ -885,7 +885,7 @@ class DataSourceServiceTest { } @Test - void queryTraceScopeAllowsRawTraceJsonResourceFilters() { + void queryTraceScopeAllowsRawTraceFlattenedResourceFilters() { List<Map<String, Object>> sqlData = List.of(new HashMap<>(Map.of("__value__", 1))); QueryExecutor mockExecutor = Mockito.mock(QueryExecutor.class); when(mockExecutor.support("sql")).thenReturn(true); @@ -895,7 +895,7 @@ class DataSourceServiceTest { List<Map<String, Object>> result = dataSourceService.query("sql", "SELECT service_name, span_name AS operation, span_kind, COUNT(*) AS __value__ FROM hzb_traces " + "WHERE service_name = 'checkout' " - + "AND json_get_string(resource_attributes, '$[\"service.version\"]') = '1.2.3' " + + "AND \"resource_attributes.service.version\" = '1.2.3' " + "AND span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') " + "GROUP BY service_name, span_name, span_kind HAVING __value__ > 0", TRACE_ALERT_THRESHOLD_TYPE_PERIODIC); diff --git a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeTraceApmFlowProofE2eTest.java b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeTraceApmFlowProofE2eTest.java index ac35ab7b41..8b95f2efb2 100644 --- a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeTraceApmFlowProofE2eTest.java +++ b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeTraceApmFlowProofE2eTest.java @@ -22,6 +22,17 @@ import static org.awaitility.Awaitility.await; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.protobuf.ByteString; +import io.opentelemetry.proto.collector.trace.v1.ExportTraceServiceRequest; +import io.opentelemetry.proto.collector.trace.v1.ExportTraceServiceResponse; +import io.opentelemetry.proto.common.v1.AnyValue; +import io.opentelemetry.proto.common.v1.KeyValue; +import io.opentelemetry.proto.resource.v1.Resource; +import io.opentelemetry.proto.trace.v1.ResourceSpans; +import io.opentelemetry.proto.trace.v1.ScopeSpans; +import io.opentelemetry.proto.trace.v1.Span; +import io.opentelemetry.proto.trace.v1.Status; +import java.io.InputStream; import java.net.URI; import java.net.URLEncoder; import java.net.http.HttpClient; @@ -31,7 +42,9 @@ import java.nio.charset.StandardCharsets; import java.time.Duration; import java.time.Instant; import java.time.temporal.ChronoUnit; +import java.util.HexFormat; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import org.junit.jupiter.api.Test; import org.testcontainers.containers.GenericContainer; @@ -41,24 +54,31 @@ import org.testcontainers.junit.jupiter.Testcontainers; import org.testcontainers.utility.DockerImageName; /** - * Proves Greptime Flow can derive HertzBeat APM RED alert data from native trace rows. + * Proves Greptime's native OTLP trace pipeline writes the production HertzBeat trace schema and RED Flow. */ @Testcontainers class GreptimeTraceApmFlowProofE2eTest { - private static final String GREPTIME_IMAGE = "greptime/greptimedb:latest"; + 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 TRACE_TABLE = "hzb_traces"; - private static final String APM_FLOW = "hertzbeat_apm_red_1m_flow"; private static final String APM_TABLE = "hertzbeat_apm_red_1m"; private static final String SERVICE_NAME = "checkout"; private static final String OPERATION = "GET /checkout"; - private static final String SPAN_KIND = "SPAN_KIND_SERVER"; + private static final String ROLLUP_SPAN_KIND = "SERVER"; private static final String WORKSPACE_ID = "workspace-trace-flow-proof"; private static final String ENTITY_ID = "entity-trace-flow-proof"; + private static final String ENTITY_TYPE = "service"; private static final String ENVIRONMENT = "prod"; private static final String SERVICE_NAMESPACE = "payments"; + private static final String SERVICE_INSTANCE_ID = "checkout-proof-7d9"; + private static final String COLLECTOR_ID = "collector-trace-flow-proof"; + private static final String RESOURCE_LONG_TAIL_VALUE = "2026.08-proof"; + private static final String SPAN_LONG_TAIL_VALUE = "postgresql"; + private static final String TRACE_ID = "0123456789abcdef0123456789abcdef"; + private static final String SPAN_ID = "0123456789abcdef"; + private static final long DURATION_NANOS = 1_200_000_000L; private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); @Container @@ -75,32 +95,66 @@ class GreptimeTraceApmFlowProofE2eTest { private final long windowNanos = Instant.now().truncatedTo(ChronoUnit.MINUTES).toEpochMilli() * 1_000_000L; @Test - void greptimeFlowCanDeriveApmRedRowsFromNativeTraceRows() throws Exception { - executeSql(traceTableDdl()); - executeSql(apmSinkDdl()); - executeSql(apmFlowSql()); - insertTraceRows(); + void greptimeOtlpPipelinePopulatesFlattenedTraceColumnsAndRedFlow() throws Exception { + executeSqlScript("greptime/tables/hzb_traces.sql"); + executeSqlScript("greptime/flows/hertzbeat_apm_red_1m.sql"); + ExportTraceServiceResponse response = exportTrace(); + + assertThat(response.hasPartialSuccess()).isFalse(); + + await().atMost(Duration.ofSeconds(45)).pollInterval(Duration.ofSeconds(1)).untilAsserted(() -> { + Map<String, Object> row = querySingleRow("SELECT service_name, trace_id, span_id, span_kind, " + + "span_status_code, duration_nano, " + + "\"resource_attributes.hertzbeat.workspace_id\" AS workspace_id, " + + "\"resource_attributes.hertzbeat.entity_id\" AS entity_id, " + + "\"resource_attributes.hertzbeat.entity_type\" AS entity_type, " + + "\"resource_attributes.service.namespace\" AS service_namespace, " + + "\"resource_attributes.service.instance.id\" AS service_instance_id, " + + "\"resource_attributes.deployment.environment.name\" AS deployment_environment, " + + "\"resource_attributes.hertzbeat.collector.id\" AS collector_id, " + + "\"resource_attributes.service.version\" AS resource_long_tail, " + + "\"span_attributes.db.system\" AS span_long_tail " + + "FROM " + TRACE_TABLE + " WHERE trace_id = '" + TRACE_ID + "'"); + + assertThat(row) + .containsEntry("service_name", SERVICE_NAME) + .containsEntry("trace_id", TRACE_ID) + .containsEntry("span_id", SPAN_ID) + .containsEntry("span_kind", "SPAN_KIND_SERVER") + .containsEntry("span_status_code", "STATUS_CODE_ERROR") + .containsEntry("workspace_id", WORKSPACE_ID) + .containsEntry("entity_id", ENTITY_ID) + .containsEntry("entity_type", ENTITY_TYPE) + .containsEntry("service_namespace", SERVICE_NAMESPACE) + .containsEntry("service_instance_id", SERVICE_INSTANCE_ID) + .containsEntry("deployment_environment", ENVIRONMENT) + .containsEntry("collector_id", COLLECTOR_ID) + .containsEntry("resource_long_tail", RESOURCE_LONG_TAIL_VALUE) + .containsEntry("span_long_tail", SPAN_LONG_TAIL_VALUE); + assertThat(number(row.get("duration_nano"))).isEqualTo(DURATION_NANOS); + }); await().atMost(Duration.ofSeconds(45)).pollInterval(Duration.ofSeconds(1)).untilAsserted(() -> { Map<String, Object> row = querySingleRow("SELECT service_name, operation, span_kind, workspace_id, " - + "entity_id, deployment_environment, service_namespace, calls_total, error_total, " + + "entity_id, entity_type, deployment_environment, service_namespace, calls_total, error_total, " + "duration_sum_nano, duration_count FROM " + APM_TABLE + " WHERE service_name = '" + SERVICE_NAME + "'" + " AND operation = '" + OPERATION + "'" - + " AND span_kind = '" + SPAN_KIND + "'"); + + " AND span_kind = '" + ROLLUP_SPAN_KIND + "'"); assertThat(row) .containsEntry("service_name", SERVICE_NAME) .containsEntry("operation", OPERATION) - .containsEntry("span_kind", SPAN_KIND) + .containsEntry("span_kind", ROLLUP_SPAN_KIND) .containsEntry("workspace_id", WORKSPACE_ID) .containsEntry("entity_id", ENTITY_ID) + .containsEntry("entity_type", ENTITY_TYPE) .containsEntry("deployment_environment", ENVIRONMENT) .containsEntry("service_namespace", SERVICE_NAMESPACE); - assertThat(number(row.get("calls_total"))).isEqualTo(2L); + assertThat(number(row.get("calls_total"))).isEqualTo(1L); assertThat(number(row.get("error_total"))).isEqualTo(1L); - assertThat(number(row.get("duration_sum_nano"))).isEqualTo(1_300_000_000L); - assertThat(number(row.get("duration_count"))).isEqualTo(2L); + assertThat(number(row.get("duration_sum_nano"))).isEqualTo(DURATION_NANOS); + assertThat(number(row.get("duration_count"))).isEqualTo(1L); Map<String, Object> sketchRow = querySingleRow("SELECT " + "uddsketch_calc(0.95, uddsketch_merge(128, 0.01, duration_sketch)) AS p95_nano " @@ -112,116 +166,65 @@ class GreptimeTraceApmFlowProofE2eTest { }); } - private String traceTableDdl() { - return "CREATE TABLE IF NOT EXISTS " + TRACE_TABLE + " (" - + "\"timestamp\" TIMESTAMP(9) TIME INDEX," - + "\"timestamp_end\" TIMESTAMP(9) NULL," - + "\"duration_nano\" BIGINT UNSIGNED NULL," - + "\"trace_id\" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM')," - + "\"span_id\" STRING NULL," - + "\"parent_span_id\" STRING NULL," - + "\"span_kind\" STRING NULL," - + "\"span_name\" STRING NULL," - + "\"span_status_code\" STRING NULL," - + "\"span_status_message\" STRING NULL," - + "\"trace_state\" STRING NULL," - + "\"scope_name\" STRING NULL," - + "\"scope_version\" STRING NULL," - + "\"service_name\" STRING NULL," - + "\"resource_attributes\" JSON NULL," - + "\"span_attributes\" JSON NULL," - + "\"span_events\" JSON NULL," - + "\"span_links\" JSON NULL," - + "PRIMARY KEY(\"service_name\"))" - + " WITH (append_mode = true, table_data_model = 'greptime_trace_v1')"; - } - - private String apmSinkDdl() { - return "CREATE TABLE IF NOT EXISTS " + APM_TABLE + " (" - + "time_window TIMESTAMP(9) TIME INDEX," - + "service_name STRING," - + "operation STRING," - + "span_kind STRING," - + "workspace_id STRING NULL," - + "entity_id STRING NULL," - + "deployment_environment STRING NULL," - + "service_namespace STRING NULL," - + "calls_total BIGINT," - + "error_total BIGINT," - + "duration_sum_nano BIGINT," - + "duration_count BIGINT," - + "duration_sketch BINARY," - + "PRIMARY KEY(service_name, operation, span_kind, workspace_id, entity_id, " - + "deployment_environment, service_namespace))"; - } + private ExportTraceServiceResponse exportTrace() throws Exception { + Span span = Span.newBuilder() + .setTraceId(ByteString.copyFrom(HexFormat.of().parseHex(TRACE_ID))) + .setSpanId(ByteString.copyFrom(HexFormat.of().parseHex(SPAN_ID))) + .setName(OPERATION) + .setKind(Span.SpanKind.SPAN_KIND_SERVER) + .setStartTimeUnixNano(windowNanos + 1_000_000_000L) + .setEndTimeUnixNano(windowNanos + 1_000_000_000L + DURATION_NANOS) + .setStatus(Status.newBuilder().setCode(Status.StatusCode.STATUS_CODE_ERROR)) + .addAttributes(attribute("db.system", SPAN_LONG_TAIL_VALUE)) + .build(); + Resource resource = Resource.newBuilder().addAllAttributes(List.of( + attribute("service.name", SERVICE_NAME), + attribute("hertzbeat.workspace_id", WORKSPACE_ID), + attribute("hertzbeat.entity_id", ENTITY_ID), + attribute("hertzbeat.entity_type", ENTITY_TYPE), + attribute("service.namespace", SERVICE_NAMESPACE), + attribute("service.instance.id", SERVICE_INSTANCE_ID), + attribute("deployment.environment.name", ENVIRONMENT), + attribute("hertzbeat.collector.id", COLLECTOR_ID), + attribute("service.version", RESOURCE_LONG_TAIL_VALUE))).build(); + ExportTraceServiceRequest request = ExportTraceServiceRequest.newBuilder() + .addResourceSpans(ResourceSpans.newBuilder() + .setResource(resource) + .addScopeSpans(ScopeSpans.newBuilder().addSpans(span))) + .build(); - private String apmFlowSql() { - return "CREATE FLOW IF NOT EXISTS " + APM_FLOW + " " - + "SINK TO " + APM_TABLE + " " - + "EXPIRE AFTER '6 hours'::INTERVAL " - + "AS SELECT " - + "date_bin('1 minute'::INTERVAL, \"timestamp\") AS time_window," - + "service_name," - + "span_name AS operation," - + "span_kind," - + "json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') AS workspace_id," - + "json_get_string(resource_attributes, '$[\"hertzbeat.entity_id\"]') AS entity_id," - + "json_get_string(resource_attributes, '$[\"deployment.environment\"]') AS deployment_environment," - + "json_get_string(resource_attributes, '$[\"service.namespace\"]') AS service_namespace," - + "COUNT(*) AS calls_total," - + "SUM(CASE WHEN span_status_code = 'STATUS_CODE_ERROR' THEN 1 ELSE 0 END) AS error_total," - + "SUM(duration_nano) AS duration_sum_nano," - + "COUNT(duration_nano) AS duration_count," - + "uddsketch_state(128, 0.01, duration_nano) AS duration_sketch " - + "FROM " + TRACE_TABLE + " " - + "WHERE span_kind IN ('SPAN_KIND_SERVER', 'SERVER', 'SPAN_KIND_CONSUMER', 'CONSUMER') " - + "GROUP BY time_window, service_name, operation, span_kind, workspace_id, entity_id, " - + "deployment_environment, service_namespace"; - } + HttpResponse<byte[]> response = httpClient.send(HttpRequest.newBuilder() + .uri(URI.create(greptimeEndpoint() + "/v1/otlp/v1/traces")) + .header("Content-Type", "application/x-protobuf") + .header("X-Greptime-DB-Name", "public") + .header("X-Greptime-Trace-Table-Name", TRACE_TABLE) + .header("X-Greptime-Pipeline-Name", "greptime_trace_v1") + .POST(HttpRequest.BodyPublishers.ofByteArray(request.toByteArray())) + .build(), HttpResponse.BodyHandlers.ofByteArray()); - private void insertTraceRows() throws Exception { - executeSql("INSERT INTO " + TRACE_TABLE + " (\"timestamp\", \"timestamp_end\", duration_nano, trace_id, " - + "span_id, parent_span_id, span_kind, span_name, span_status_code, span_status_message, " - + "trace_state, scope_name, scope_version, service_name, resource_attributes, span_attributes, " - + "span_events, span_links) VALUES " - + traceRow(windowNanos + 1_000_000_000L, 100_000_000L, "trace-flow-ok-1", "span-flow-ok-1", - "STATUS_CODE_OK") - + "," - + traceRow(windowNanos + 2_000_000_000L, 1_200_000_000L, "trace-flow-error-1", - "span-flow-error-1", "STATUS_CODE_ERROR") - + "," - + traceRow(windowNanos + 3_000_000_000L, 50_000_000L, "trace-flow-client-1", - "span-flow-client-1", "STATUS_CODE_ERROR").replace("'" + SPAN_KIND + "'", "'SPAN_KIND_CLIENT'")); + assertThat(response.statusCode()) + .as("OTLP response body: %s", new String(response.body(), StandardCharsets.UTF_8)) + .isBetween(200, 299); + return ExportTraceServiceResponse.parseFrom(response.body()); } - private String traceRow(long startNanos, long durationNanos, String traceId, String spanId, String statusCode) { - return "(" - + "to_timestamp_nanos(" + startNanos + ")," - + "to_timestamp_nanos(" + (startNanos + durationNanos) + ")," - + durationNanos + "," - + "'" + traceId + "'," - + "'" + spanId + "'," - + "''," - + "'" + SPAN_KIND + "'," - + "'" + OPERATION + "'," - + "'" + statusCode + "'," - + "''," - + "''," - + "'io.opentelemetry'," - + "'1.0.0'," - + "'" + SERVICE_NAME + "'," - + "parse_json('" + resourceAttributesJson() + "')," - + "parse_json('{\"http.route\":\"/checkout\"}')," - + "parse_json('[]')," - + "parse_json('[]')" - + ")"; + private KeyValue attribute(String key, String value) { + return KeyValue.newBuilder() + .setKey(key) + .setValue(AnyValue.newBuilder().setStringValue(value)) + .build(); } - private String resourceAttributesJson() { - return "{\"hertzbeat.workspace_id\":\"" + WORKSPACE_ID - + "\",\"hertzbeat.entity_id\":\"" + ENTITY_ID - + "\",\"deployment.environment\":\"" + ENVIRONMENT - + "\",\"service.namespace\":\"" + SERVICE_NAMESPACE + "\"}"; + private void executeSqlScript(String resourceName) throws Exception { + try (InputStream input = Thread.currentThread().getContextClassLoader().getResourceAsStream(resourceName)) { + assertThat(input).as(resourceName).isNotNull(); + String sql = new String(input.readAllBytes(), StandardCharsets.UTF_8); + for (String statement : sql.split(";")) { + if (!statement.isBlank()) { + executeSql(statement.strip()); + } + } + } } private Map<String, Object> querySingleRow(String sql) throws Exception { diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializer.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializer.java index 6ab86e13cd..b48c8b9f2b 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializer.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializer.java @@ -17,15 +17,10 @@ package org.apache.hertzbeat.observability.ingestion.forwarder; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import java.io.IOException; import java.nio.charset.StandardCharsets; import java.util.Arrays; -import java.util.Collections; -import java.util.LinkedHashSet; import java.util.List; -import java.util.Set; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.apache.hertzbeat.common.runtime.ConditionalOnNormalBusinessRuntime; @@ -43,7 +38,6 @@ import org.springframework.http.HttpMethod; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.stereotype.Component; -import org.springframework.web.client.HttpClientErrorException; import org.springframework.web.client.RestTemplate; import org.springframework.web.util.UriUtils; @@ -60,10 +54,6 @@ public class GreptimeApmFlowInitializer { private static final String SQL_PATH = "/v1/sql"; private static final String DEFAULT_GREPTIME_DB_NAME = "public"; - private static final String TRACE_TABLE = "hzb_traces"; - private static final String RESOURCE_ATTRIBUTES_COLUMN = "resource_attributes"; - private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); - private final RestTemplate restTemplate; private final ObjectProvider<GreptimeProperties> greptimePropertiesProvider; private final OtlpIngestionRetryService retryService; @@ -92,26 +82,9 @@ public class GreptimeApmFlowInitializer { } try { List<String> statements = readFlowStatements(); - Set<String> traceTableColumns = null; for (String statement : statements) { - try { - if (!executeStatement(greptimeProperties, statement)) { - return; - } - } catch (HttpClientErrorException.BadRequest ex) { - if (!isTraceResourceAttributeSchemaFailure(statement, ex)) { - throw ex; - } - if (traceTableColumns == null) { - traceTableColumns = traceTableColumns(greptimeProperties); - } - String adaptedStatement = adaptTraceResourceAttributeExpressions(statement, traceTableColumns); - if (statement.equals(adaptedStatement)) { - throw ex; - } - if (!executeStatement(greptimeProperties, adaptedStatement)) { - return; - } + if (!executeStatement(greptimeProperties, statement)) { + return; } } log.info("[observability greptime-apm-flow] initialized Greptime APM RED Flow."); @@ -169,83 +142,6 @@ public class GreptimeApmFlowInitializer { .toList(); } - private boolean isTraceResourceAttributeSchemaFailure(String statement, HttpClientErrorException ex) { - if (!StringUtils.containsIgnoreCase(statement, "CREATE FLOW IF NOT EXISTS hertzbeat_apm_red_1m_flow")) { - return false; - } - String responseBody = ex.getResponseBodyAsString(StandardCharsets.UTF_8); - return StringUtils.containsIgnoreCase(ex.getMessage(), "resource_attributes") - || StringUtils.containsIgnoreCase(responseBody, "resource_attributes"); - } - - private Set<String> traceTableColumns(GreptimeProperties greptimeProperties) { - try { - ResponseEntity<String> response = restTemplate.exchange( - endpoint(greptimeProperties), - HttpMethod.POST, - sqlRequest(greptimeProperties, "DESC " + TRACE_TABLE), - String.class); - if (response == null || !response.getStatusCode().is2xxSuccessful() - || StringUtils.isBlank(response.getBody())) { - return Collections.singleton(RESOURCE_ATTRIBUTES_COLUMN); - } - return parseTraceTableColumns(response.getBody()); - } catch (RuntimeException | IOException ex) { - log.debug("[observability greptime-apm-flow] failed to inspect Greptime trace schema: {}", - ex.getMessage(), ex); - return Collections.singleton(RESOURCE_ATTRIBUTES_COLUMN); - } - } - - private Set<String> parseTraceTableColumns(String responseBody) throws IOException { - Set<String> columns = new LinkedHashSet<>(); - JsonNode root = OBJECT_MAPPER.readTree(responseBody); - for (JsonNode output : root.path("output")) { - JsonNode rows = output.path("records").path("rows"); - if (!rows.isArray()) { - continue; - } - for (JsonNode row : rows) { - if (row.isArray() && !row.isEmpty() && StringUtils.isNotBlank(row.get(0).asText())) { - columns.add(row.get(0).asText()); - } - } - } - if (columns.isEmpty()) { - return Collections.singleton(RESOURCE_ATTRIBUTES_COLUMN); - } - return Collections.unmodifiableSet(columns); - } - - private String adaptTraceResourceAttributeExpressions(String statement, Set<String> traceTableColumns) { - String adaptedStatement = statement; - for (String key : List.of( - "hertzbeat.workspace_id", - "hertzbeat.entity_id", - "deployment.environment.name", - "service.namespace")) { - adaptedStatement = adaptedStatement.replace( - "json_get_string(resource_attributes, '$[\"" + key + "\"]')", - resourceAttributeExpression(traceTableColumns, key)); - } - return adaptedStatement; - } - - private String resourceAttributeExpression(Set<String> traceTableColumns, String key) { - String flattenedColumn = RESOURCE_ATTRIBUTES_COLUMN + "." + key; - if (traceTableColumns.contains(flattenedColumn)) { - return quoteIdentifier(flattenedColumn); - } - if (traceTableColumns.contains(RESOURCE_ATTRIBUTES_COLUMN)) { - return "json_get_string(resource_attributes, '$[\"" + key.replace("\"", "\\\"") + "\"]')"; - } - return "NULL"; - } - - private String quoteIdentifier(String column) { - return "\"" + column.replace("\"", "\"\"") + "\""; - } - private void addAuthenticationHeader(HttpHeaders headers, GreptimeProperties greptimeProperties) { String username = StringUtils.trimToNull(greptimeProperties.username()); String password = StringUtils.trimToNull(greptimeProperties.password()); diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImpl.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImpl.java index aeec3a2701..b71ce17aea 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImpl.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImpl.java @@ -1848,7 +1848,7 @@ public class EntityTraceQueryServiceImpl implements EntityTraceQueryService { } private TraceListItemDto toTraceListItem(Map<String, Object> row) { - Map<String, String> resourceAttributes = parseAttributes(row.get("resource_attributes"), "resource_attributes.", row); + Map<String, String> resourceAttributes = parseAttributes("resource_attributes.", row); String serviceName = defaultText(readText(row, "service_name"), resourceAttributes.get("service.name")); String serviceNamespace = defaultText(readText(row, "service_namespace"), resourceAttributes.get("service.namespace")); @@ -1885,8 +1885,8 @@ public class EntityTraceQueryServiceImpl implements EntityTraceQueryService { } private TraceSpanNodeDto toSpanNode(Map<String, Object> row) { - Map<String, String> resourceAttributes = parseAttributes(row.get("resource_attributes"), "resource_attributes.", row); - Map<String, String> spanAttributes = parseAttributes(row.get("span_attributes"), "span_attributes.", row); + Map<String, String> resourceAttributes = parseAttributes("resource_attributes.", row); + Map<String, String> spanAttributes = parseAttributes("span_attributes.", row); String status = normalizeStatus(readText(row, "span_status_code")); TraceSpanNodeDto span = new TraceSpanNodeDto(); span.setTraceId(readText(row, "trace_id")); @@ -2054,19 +2054,8 @@ public class EntityTraceQueryServiceImpl implements EntityTraceQueryService { return value == null ? null : Math.max(0, value); } - private Map<String, String> parseAttributes(Object rawValue, String prefix, Map<String, Object> row) { + private Map<String, String> parseAttributes(String prefix, Map<String, Object> row) { Map<String, String> values = new LinkedHashMap<>(); - if (rawValue instanceof String rawText && StringUtils.hasText(rawText)) { - try { - Object parsed = JSON_MAPPER.readValue(rawText, new TypeReference<>() { - }); - collectTraceAttributes(values, parsed); - } catch (Exception ignored) { - // Keep fallback scan below. - } - } else { - collectTraceAttributes(values, rawValue); - } row.forEach((key, value) -> { String normalizedKey = trimText(key); if (normalizedKey == null || !normalizedKey.startsWith(prefix)) { @@ -2081,47 +2070,6 @@ public class EntityTraceQueryServiceImpl implements EntityTraceQueryService { return values; } - private void collectTraceAttributes(Map<String, String> values, Object rawValue) { - if (rawValue instanceof Map<?, ?> rawMap) { - Object key = rawMap.get("key"); - if (key != null && rawMap.containsKey("value")) { - putTraceAttribute(values, key, rawMap.get("value")); - return; - } - rawMap.forEach((attributeKey, attributeValue) -> putTraceAttribute(values, attributeKey, attributeValue)); - return; - } - if (rawValue instanceof Iterable<?> items) { - items.forEach(item -> collectTraceAttributes(values, item)); - } - } - - private void putTraceAttribute(Map<String, String> values, Object key, Object value) { - String normalizedKey = trimText(Objects.toString(key, null)); - String normalizedValue = traceAttributeValue(value); - if (normalizedKey != null && normalizedValue != null) { - values.put(normalizedKey, normalizedValue); - } - } - - private String traceAttributeValue(Object value) { - if (value == null) { - return null; - } - if (value instanceof Map<?, ?> rawMap) { - for (String key : List.of("stringValue", "string_value", "intValue", "int_value", - "doubleValue", "double_value", "boolValue", "bool_value")) { - if (rawMap.containsKey(key)) { - return traceAttributeValue(rawMap.get(key)); - } - } - if (rawMap.containsKey("value")) { - return traceAttributeValue(rawMap.get("value")); - } - } - return trimText(Objects.toString(value, null)); - } - private String resolveEntityTitle(ObservedEntityContext entityContext) { if (entityContext == null || entityContext.getEntity() == null) { return "entity"; diff --git a/hertzbeat-observability/src/main/resources/greptime/flows/hertzbeat_apm_red_1m.sql b/hertzbeat-observability/src/main/resources/greptime/flows/hertzbeat_apm_red_1m.sql index 25499c4bf6..f4f8a208cc 100644 --- a/hertzbeat-observability/src/main/resources/greptime/flows/hertzbeat_apm_red_1m.sql +++ b/hertzbeat-observability/src/main/resources/greptime/flows/hertzbeat_apm_red_1m.sql @@ -5,6 +5,7 @@ CREATE TABLE IF NOT EXISTS hertzbeat_apm_red_1m ( span_kind STRING, workspace_id STRING NULL, entity_id STRING NULL, + entity_type STRING NULL, deployment_environment STRING NULL, service_namespace STRING NULL, calls_total BIGINT, @@ -12,7 +13,7 @@ CREATE TABLE IF NOT EXISTS hertzbeat_apm_red_1m ( duration_sum_nano BIGINT, duration_count BIGINT, duration_sketch BINARY, - PRIMARY KEY(service_name, operation, span_kind, workspace_id, entity_id, deployment_environment, service_namespace) + PRIMARY KEY(service_name, operation, span_kind, workspace_id, entity_id, entity_type, deployment_environment, service_namespace) ); CREATE FLOW IF NOT EXISTS hertzbeat_apm_red_1m_flow @@ -27,10 +28,11 @@ AS SELECT WHEN span_kind IN ('SPAN_KIND_CONSUMER', 'CONSUMER') THEN 'CONSUMER' ELSE 'UNKNOWN' END AS span_kind, - json_get_string(resource_attributes, '$["hertzbeat.workspace_id"]') AS workspace_id, - json_get_string(resource_attributes, '$["hertzbeat.entity_id"]') AS entity_id, - json_get_string(resource_attributes, '$["deployment.environment.name"]') AS deployment_environment, - json_get_string(resource_attributes, '$["service.namespace"]') AS service_namespace, + "resource_attributes.hertzbeat.workspace_id" AS workspace_id, + "resource_attributes.hertzbeat.entity_id" AS entity_id, + "resource_attributes.hertzbeat.entity_type" AS entity_type, + "resource_attributes.deployment.environment.name" AS deployment_environment, + "resource_attributes.service.namespace" AS service_namespace, COUNT(*) AS calls_total, SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') THEN 1 ELSE 0 END) AS error_total, COALESCE(SUM(duration_nano), 0) AS duration_sum_nano, @@ -38,4 +40,4 @@ AS SELECT uddsketch_state(128, 0.01, duration_nano) AS duration_sketch FROM hzb_traces WHERE span_kind IN ('SPAN_KIND_SERVER', 'SERVER', 'SPAN_KIND_CONSUMER', 'CONSUMER') -GROUP BY time_window, service_name, operation, span_kind, workspace_id, entity_id, deployment_environment, service_namespace; +GROUP BY time_window, service_name, operation, span_kind, workspace_id, entity_id, entity_type, deployment_environment, service_namespace; diff --git a/hertzbeat-observability/src/main/resources/greptime/tables/hzb_traces.sql b/hertzbeat-observability/src/main/resources/greptime/tables/hzb_traces.sql index 3aa7f0ee36..095f5a8dec 100644 --- a/hertzbeat-observability/src/main/resources/greptime/tables/hzb_traces.sql +++ b/hertzbeat-observability/src/main/resources/greptime/tables/hzb_traces.sql @@ -1,10 +1,10 @@ CREATE TABLE IF NOT EXISTS hzb_traces ( - "timestamp" TIMESTAMP(9) TIME INDEX, + "timestamp" TIMESTAMP(9) NOT NULL TIME INDEX, "timestamp_end" TIMESTAMP(9) NULL, "duration_nano" BIGINT UNSIGNED NULL, "trace_id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), "span_id" STRING NULL, - "parent_span_id" STRING NULL, + "parent_span_id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), "span_kind" STRING NULL, "span_name" STRING NULL, "span_status_code" STRING NULL, @@ -12,9 +12,14 @@ CREATE TABLE IF NOT EXISTS hzb_traces ( "trace_state" STRING NULL, "scope_name" STRING NULL, "scope_version" STRING NULL, - "service_name" STRING NULL, - "resource_attributes" JSON NULL, - "span_attributes" JSON NULL, + "service_name" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.hertzbeat.workspace_id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.hertzbeat.entity_id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.hertzbeat.entity_type" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.service.namespace" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.service.instance.id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.deployment.environment.name" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), + "resource_attributes.hertzbeat.collector.id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), "span_events" JSON NULL, "span_links" JSON NULL, PRIMARY KEY("service_name") diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializerTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializerTest.java index dd691f73c6..ca6f18f9be 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializerTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeApmFlowInitializerTest.java @@ -95,7 +95,6 @@ class GreptimeApmFlowInitializerTest { assertTrue(sqlStatements.get(1).contains("CREATE FLOW IF NOT EXISTS hertzbeat_apm_red_1m_flow")); assertTrue(sqlStatements.get(1).contains("SINK TO hertzbeat_apm_red_1m")); assertTrue(sqlStatements.get(1).contains("EXPIRE AFTER '6 hours'::INTERVAL")); - for (HttpEntity<String> entity : entityCaptor.getAllValues()) { assertEquals(MediaType.APPLICATION_FORM_URLENCODED, entity.getHeaders().getContentType()); assertEquals("Basic " + Base64.getEncoder() @@ -228,7 +227,7 @@ class GreptimeApmFlowInitializerTest { } @Test - void retriesApmFlowWithFlattenedTraceResourceColumnsWhenGreptimeSchemaHasNoJsonColumn() { + void doesNotInspectOrAdaptTheTraceSchemaWhenTheCanonicalFlowIsRejected() { configureGreptimeProperties(true); when(restTemplate.exchange( eq("http://greptime:4000/v1/sql?db=public"), @@ -240,20 +239,14 @@ class GreptimeApmFlowInitializerTest { HttpStatus.BAD_REQUEST, "Bad Request", HttpHeaders.EMPTY, - ("{\"code\":3000,\"error\":\"Failed to plan SQL: No field named resource_attributes. " - + "Did you mean 'hzb_traces.resource_attributes.host.name'?\"}") + "{\"code\":3000,\"error\":\"Failed to plan canonical flow\"}" .getBytes(StandardCharsets.UTF_8), - StandardCharsets.UTF_8)) - .thenReturn(ResponseEntity.ok(""" - {"output":[{"records":{"schema":{"column_schemas":[{"name":"Column","data_type":"String"}]}, - "rows":[["timestamp"],["service_name"],["resource_attributes.host.name"]],"total_rows":3}}]} - """)) - .thenReturn(ResponseEntity.ok("{}")); + StandardCharsets.UTF_8)); initializer.initialize(); ArgumentCaptor<HttpEntity<String>> entityCaptor = ArgumentCaptor.forClass(HttpEntity.class); - verify(restTemplate, times(4)).exchange( + verify(restTemplate, times(2)).exchange( eq("http://greptime:4000/v1/sql?db=public"), eq(HttpMethod.POST), entityCaptor.capture(), @@ -262,15 +255,8 @@ class GreptimeApmFlowInitializerTest { List<String> sqlStatements = entityCaptor.getAllValues().stream() .map(this::decodeSql) .toList(); - assertTrue(sqlStatements.get(2).contains("DESC hzb_traces")); - String adaptedFlowSql = sqlStatements.get(3); - assertTrue(adaptedFlowSql.contains("CREATE FLOW IF NOT EXISTS hertzbeat_apm_red_1m_flow")); - assertTrue(adaptedFlowSql.contains("NULL AS workspace_id")); - assertTrue(adaptedFlowSql.contains("NULL AS entity_id")); - assertTrue(adaptedFlowSql.contains("NULL AS deployment_environment")); - assertTrue(adaptedFlowSql.contains("NULL AS service_namespace")); - assertFalse(adaptedFlowSql.contains("json_get_string(resource_attributes")); - assertFalse(adaptedFlowSql.contains("resource_attributes.host.name")); + assertTrue(sqlStatements.get(1).contains("CREATE FLOW IF NOT EXISTS hertzbeat_apm_red_1m_flow")); + assertFalse(sqlStatements.stream().anyMatch(sql -> sql.startsWith("DESC hzb_traces"))); } @Test @@ -321,11 +307,12 @@ class GreptimeApmFlowInitializerTest { assertTrue(sql.contains("date_bin('1 minute'::INTERVAL, \"timestamp\") AS time_window")); assertTrue(sql.contains("COALESCE(NULLIF(service_name, ''), 'unknown_service') AS service_name")); assertTrue(sql.contains("COALESCE(NULLIF(span_name, ''), 'unknown_operation') AS operation")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]')")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"hertzbeat.entity_id\"]')")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]')")); - assertFalse(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment\"]')")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]')")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.workspace_id\" AS workspace_id")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.entity_id\" AS entity_id")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.entity_type\" AS entity_type")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" AS deployment_environment")); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" AS service_namespace")); + assertFalse(sql.contains("json_get_string(")); assertTrue(sql.contains("span_status_code IN ('STATUS_CODE_ERROR', 'ERROR')")); assertTrue(sql.contains("CASE\n" + " WHEN span_kind IN ('SPAN_KIND_SERVER', 'SERVER') THEN 'SERVER'\n" @@ -336,7 +323,9 @@ class GreptimeApmFlowInitializerTest { assertTrue(sql.contains("uddsketch_state(128, 0.01, duration_nano) AS duration_sketch")); assertTrue(sql.contains("span_kind IN ('SPAN_KIND_SERVER', 'SERVER', 'SPAN_KIND_CONSUMER', 'CONSUMER')")); assertTrue(sql.contains("PRIMARY KEY(service_name, operation, span_kind, workspace_id, entity_id, " - + "deployment_environment, service_namespace)")); + + "entity_type, deployment_environment, service_namespace)")); + assertFalse(sql.contains("service_instance_id")); + assertFalse(sql.contains("collector_id")); assertFalse(lowerSql.contains("drop flow")); assertFalse(lowerSql.contains("drop table")); assertFalse(lowerSql.contains("trace_id")); diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeTraceTableInitializerTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeTraceTableInitializerTest.java index fa4a41c0cd..01de4db3fb 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeTraceTableInitializerTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/forwarder/GreptimeTraceTableInitializerTest.java @@ -19,6 +19,7 @@ package org.apache.hertzbeat.observability.ingestion.forwarder; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.eq; @@ -29,6 +30,8 @@ import static org.mockito.Mockito.when; import java.io.InputStream; import java.net.URLDecoder; import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; import java.util.Base64; import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.GreptimeProperties; import org.junit.jupiter.api.BeforeEach; @@ -86,7 +89,7 @@ class GreptimeTraceTableInitializerTest { String sql = decodeSql(entityCaptor.getValue()); assertTrue(sql.contains("CREATE TABLE IF NOT EXISTS hzb_traces")); - assertTrue(sql.contains("\"timestamp\" TIMESTAMP(9) TIME INDEX")); + assertTrue(sql.contains("\"timestamp\" TIMESTAMP(9) NOT NULL TIME INDEX")); assertTrue(sql.contains("\"trace_id\" STRING NULL SKIPPING INDEX")); assertTrue(sql.contains("WITH (append_mode = true, table_data_model = 'greptime_trace_v1')")); assertEquals(MediaType.APPLICATION_FORM_URLENCODED, entityCaptor.getValue().getHeaders().getContentType()); @@ -99,14 +102,66 @@ class GreptimeTraceTableInitializerTest { String sql = bundledTraceTableSql(); assertTrue(sql.contains("CREATE TABLE IF NOT EXISTS hzb_traces")); - assertTrue(sql.contains("\"timestamp\" TIMESTAMP(9) TIME INDEX")); + assertTrue(sql.contains("\"timestamp\" TIMESTAMP(9) NOT NULL TIME INDEX")); assertTrue(sql.contains("\"duration_nano\" BIGINT UNSIGNED NULL")); assertTrue(sql.contains("\"trace_id\" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM')")); - assertTrue(sql.contains("\"resource_attributes\" JSON NULL")); - assertTrue(sql.contains("PRIMARY KEY(\"service_name\")")); + assertTrue(sql.contains("\"parent_span_id\" STRING NULL SKIPPING INDEX")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.workspace_id\" STRING NULL SKIPPING INDEX " + + "WITH(granularity = '10240', type = 'BLOOM')")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.entity_id\" STRING NULL SKIPPING INDEX")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.entity_type\" STRING NULL SKIPPING INDEX " + + "WITH(granularity = '10240', type = 'BLOOM')")); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" STRING NULL SKIPPING INDEX " + + "WITH(granularity = '10240', type = 'BLOOM')")); + assertTrue(sql.contains("\"resource_attributes.service.instance.id\" STRING NULL SKIPPING INDEX")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" STRING NULL SKIPPING INDEX " + + "WITH(granularity = '10240', type = 'BLOOM')")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.collector.id\" STRING NULL SKIPPING INDEX")); + assertTrue(sql.contains("\"span_events\" JSON NULL")); + assertTrue(sql.contains("\"span_links\" JSON NULL")); + assertFalse(sql.contains("\"resource_attributes\" JSON")); + assertFalse(sql.contains("\"span_attributes\" JSON")); + assertEquals("PRIMARY KEY(\"service_name\")", primaryKeyClause(sql)); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.hertzbeat.workspace_id")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.service.namespace")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.deployment.environment.name")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.hertzbeat.entity_type")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.service.instance.id")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.hertzbeat.collector.id")); + assertFalse(primaryKeyClause(sql).contains("resource_attributes.hertzbeat.entity_id")); assertTrue(sql.contains("WITH (append_mode = true, table_data_model = 'greptime_trace_v1')")); } + @Test + void richTraceSeedUsesProductionSchemaAndNativeFlattenedPipeline() throws Exception { + String script = Files.readString(repositoryRoot().resolve("script/dev/seed-trace-rich-demo.sh")); + + assertTrue(script.contains("greptime/tables/hzb_traces.sql")); + assertTrue(script.contains("X-Greptime-Trace-Table-Name: hzb_traces")); + assertTrue(script.contains("X-Greptime-Pipeline-Name: greptime_trace_v1")); + assertTrue(script.contains("hertzbeat.workspace_id")); + assertTrue(script.contains("hertzbeat.entity_id")); + assertTrue(script.contains("hertzbeat.entity_type")); + assertTrue(script.contains("service.namespace")); + assertTrue(script.contains("deployment.environment.name")); + assertTrue(script.contains("service.version")); + assertTrue(script.contains("db.system")); + assertFalse(script.contains("\"resource_attributes\" JSON")); + assertFalse(script.contains("\"span_attributes\" JSON")); + assertFalse(script.contains("INSERT INTO hzb_traces")); + } + + private Path repositoryRoot() { + Path candidate = Path.of("").toAbsolutePath(); + while (candidate != null) { + if (Files.isRegularFile(candidate.resolve("script/dev/seed-trace-rich-demo.sh"))) { + return candidate; + } + candidate = candidate.getParent(); + } + throw new IllegalStateException("Could not locate HertzBeat repository root"); + } + @Test void trimsAndNormalizesGreptimeEndpointAndDatabaseBeforeExecutingTraceTableSql() { configureGreptimeProperties(true, " http://greptime:4000/// ", " public "); @@ -265,4 +320,12 @@ class GreptimeTraceTableInitializerTest { return new String(inputStream.readAllBytes(), StandardCharsets.UTF_8); } } + + private String primaryKeyClause(String sql) { + int start = sql.indexOf("PRIMARY KEY("); + assertTrue(start >= 0); + int end = sql.indexOf(')', start); + assertTrue(end > start); + return sql.substring(start, end + 1); + } } diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImplTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImplTest.java index 1066320c7e..1c313179a3 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImplTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/traces/service/impl/EntityTraceQueryServiceImplTest.java @@ -332,19 +332,19 @@ class EntityTraceQueryServiceImplTest { Map<String, Object> checkoutRoot = traceRow("trace-checkout", "span-root-1", null, "GET /checkout", "checkout-service", "STATUS_CODE_OK", now - 10_000, 20_000_000L, Map.of("service.name", "checkout-service")); - checkoutRoot.put("span_attributes", Map.of("span.kind", "server")); + putFlattenedAttributes(checkoutRoot, "span_attributes.", Map.of("span.kind", "server")); Map<String, Object> checkoutChild = traceRow("trace-checkout", "span-child-1", "span-root-1", "GET /checkout/{id}", "checkout-service", "STATUS_CODE_OK", now - 9_000, 5_000_000L, Map.of("service.name", "checkout-service")); - checkoutChild.put("span_attributes", Map.of("http.route", "/checkout/{id}")); + putFlattenedAttributes(checkoutChild, "span_attributes.", Map.of("http.route", "/checkout/{id}")); Map<String, Object> inventoryRoot = traceRow("trace-inventory", "span-root-2", null, "GET /inventory", "checkout-service", "STATUS_CODE_ERROR", now - 8_000, 30_000_000L, Map.of("service.name", "checkout-service")); - inventoryRoot.put("span_attributes", Map.of("http.route", "/inventory")); + putFlattenedAttributes(inventoryRoot, "span_attributes.", Map.of("http.route", "/inventory")); Map<String, Object> unknownRoot = traceRow("trace-unknown", "span-root-3", null, "GET /unknown", "checkout-service", "STATUS_CODE_OK", now - 7_000, 10_000_000L, Map.of("service.name", "checkout-service")); - unknownRoot.put("span_attributes", Map.of("span.kind", "server")); + putFlattenedAttributes(unknownRoot, "span_attributes.", Map.of("span.kind", "server")); when(traceQueryRepository.queryRecentTraceRows( eq(1500), eq(start), eq(end), eq("checkout-service"), org.mockito.ArgumentMatchers.isNull(), org.mockito.ArgumentMatchers.isNull(), org.mockito.ArgumentMatchers.isNull(), @@ -939,15 +939,18 @@ class EntityTraceQueryServiceImplTest { Map<String, Object> matchingRow = traceRow("trace-checkout", "span-root-1", null, "GET /checkout", "checkout-service", "STATUS_CODE_OK", now - 10_000, 2_000_000L, Map.of("service.name", "checkout-service", "service.version", "1.2.3")); - matchingRow.put("span_attributes", Map.of("http.route", "/checkout/{id}", "span.kind", "server")); + putFlattenedAttributes(matchingRow, "span_attributes.", + Map.of("http.route", "/checkout/{id}", "span.kind", "server")); Map<String, Object> inventoryRow = traceRow("trace-inventory", "span-root-2", null, "GET /inventory", "checkout-service", "STATUS_CODE_OK", now - 9_000, 2_000_000L, Map.of("service.name", "checkout-service", "service.version", "1.2.3")); - inventoryRow.put("span_attributes", Map.of("http.route", "/inventory", "span.kind", "server")); + putFlattenedAttributes(inventoryRow, "span_attributes.", + Map.of("http.route", "/inventory", "span.kind", "server")); Map<String, Object> databaseRow = traceRow("trace-db", "span-root-3", null, "GET /checkout", "checkout-service", "STATUS_CODE_OK", now - 8_000, 2_000_000L, Map.of("service.name", "checkout-service", "service.version", "1.2.3")); - databaseRow.put("span_attributes", Map.of("http.route", "/checkout/{id}", "db.system", "mysql")); + putFlattenedAttributes(databaseRow, "span_attributes.", + Map.of("http.route", "/checkout/{id}", "db.system", "mysql")); when(traceQueryRepository.queryRecentTraceRows( eq(1500), eq(start), eq(end), eq("checkout-service"), org.mockito.ArgumentMatchers.isNull(), org.mockito.ArgumentMatchers.isNull(), org.mockito.ArgumentMatchers.isNull(), @@ -1498,21 +1501,19 @@ class EntityTraceQueryServiceImplTest { } @Test - void getTraceDetailParsesOtlpKeyValueResourceAttributesForHertzBeatAttribution() { + void getTraceDetailReadsOnlyFlattenedPhysicalAttributeColumns() { long now = System.currentTimeMillis(); Map<String, Object> row = traceRow("trace-attribution", "span-root", null, "POST /checkout", "checkout", "STATUS_CODE_OK", now, 4_000_000L, Map.of()); - row.put("resource_attributes", List.of( - Map.of("key", "service.namespace", "value", Map.of("stringValue", "payments")), - Map.of("key", "deployment.environment.name", "value", Map.of("stringValue", "prod-east")), - Map.of("key", "hertzbeat.entity_id", "value", Map.of("stringValue", "4200")), - Map.of("key", "hertzbeat.entity_name", "value", Map.of("stringValue", "Checkout API")), - Map.of("key", "hertzbeat.collector", "value", Map.of("stringValue", "collector-a")), - Map.of("key", "hertzbeat.template", "value", Map.of("stringValue", "spring-boot")) - )); - row.put("span_attributes", List.of( - Map.of("key", "db.statement", "value", Map.of("stringValue", "select 1")) - )); + row.put("resource_attributes.service.namespace", "payments"); + row.put("resource_attributes.deployment.environment.name", "prod-east"); + row.put("resource_attributes.hertzbeat.entity_id", "4200"); + row.put("resource_attributes.hertzbeat.entity_name", "Checkout API"); + row.put("resource_attributes.hertzbeat.collector.id", "collector-a"); + row.put("resource_attributes.hertzbeat.template", "spring-boot"); + row.put("span_attributes.db.statement", "select 1"); + row.put("resource_attributes", Map.of("legacy.attribute", "must-not-be-read")); + row.put("span_attributes", Map.of("legacy.attribute", "must-not-be-read")); stubDefaultWorkspaceTraceRows("trace-attribution", List.of(row)); TraceDetailDto detail = entityTraceQueryService.getTraceDetail(null, "trace-attribution"); @@ -1520,11 +1521,13 @@ class EntityTraceQueryServiceImplTest { assertNotNull(detail); assertEquals("4200", detail.getResourceAttributes().get("hertzbeat.entity_id")); assertEquals("Checkout API", detail.getResourceAttributes().get("hertzbeat.entity_name")); - assertEquals("collector-a", detail.getResourceAttributes().get("hertzbeat.collector")); + assertEquals("collector-a", detail.getResourceAttributes().get("hertzbeat.collector.id")); assertEquals("spring-boot", detail.getResourceAttributes().get("hertzbeat.template")); assertEquals("payments", detail.getServiceNamespace()); assertEquals("prod-east", detail.getResourceAttributes().get("deployment.environment.name")); assertEquals("select 1", detail.getSpans().getFirst().getSpanAttributes().get("db.statement")); + assertFalse(detail.getResourceAttributes().containsKey("legacy.attribute")); + assertFalse(detail.getSpans().getFirst().getSpanAttributes().containsKey("legacy.attribute")); } @Test @@ -1772,11 +1775,17 @@ class EntityTraceQueryServiceImplTest { row.put("span_status_code", status); row.put("duration_nano", durationNanos); row.put("timestamp", Timestamp.from(Instant.ofEpochMilli(timestampMillis))); - row.put("resource_attributes", resourceAttributes); - row.put("span_attributes", Map.of("db.statement", "select 1")); + putFlattenedAttributes(row, "resource_attributes.", resourceAttributes); + putFlattenedAttributes(row, "span_attributes.", Map.of("db.statement", "select 1")); return row; } + private void putFlattenedAttributes(Map<String, Object> row, + String prefix, + Map<String, String> attributes) { + attributes.forEach((key, value) -> row.put(prefix + key, value)); + } + private Map<String, Object> traceListRow(String traceId, String rootSpanId, String rootSpanName, String serviceName, String serviceNamespace, String status, long timestampMillis, long durationNanos, int errorSpanCount, diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepository.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepository.java index 3ffbeed021..529abdeb75 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepository.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepository.java @@ -20,6 +20,7 @@ package org.apache.hertzbeat.warehouse.repository; import java.nio.charset.StandardCharsets; +import java.time.Duration; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -28,6 +29,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.function.LongSupplier; import lombok.extern.slf4j.Slf4j; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.warehouse.constants.WarehouseConstants; @@ -35,6 +37,7 @@ import org.apache.hertzbeat.warehouse.db.GreptimeSqlQueryExecutor; import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.GreptimeProperties; import org.apache.hertzbeat.warehouse.store.history.tsdb.greptime.GreptimeSqlQueryContent; import org.springframework.beans.factory.ObjectProvider; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.http.HttpEntity; import org.springframework.http.HttpHeaders; @@ -59,19 +62,44 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { private static final String TRACE_SELECT_COLUMNS = "*"; private static final String SELF_TELEMETRY_SERVICE_FILTER = "LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')"; - private static final String RESOURCE_ATTRIBUTES_COLUMN = "resource_attributes"; + private static final List<String> TRACE_LIST_RESOURCE_ATTRIBUTE_KEYS = List.of( + "hertzbeat.workspace_id", + "hertzbeat.entity_id", + "hertzbeat.entity_type", + "service.namespace", + "service.instance.id", + "deployment.environment.name", + "hertzbeat.collector.id"); + private static final Set<String> STABLE_RESOURCE_ATTRIBUTE_KEYS = Set.copyOf(TRACE_LIST_RESOURCE_ATTRIBUTE_KEYS); + private static final int MAX_DISCOVERED_DYNAMIC_ATTRIBUTE_COLUMNS = 4_096; + private static final long DYNAMIC_ATTRIBUTE_SCHEMA_REFRESH_NANOS = Duration.ofSeconds(30).toNanos(); private final ObjectProvider<GreptimeSqlQueryExecutor> greptimeSqlQueryExecutorProvider; private final GreptimeProperties greptimeProperties; private final RestTemplate restTemplate; - private volatile Set<String> traceTableColumns; + private final LongSupplier monotonicNanos; + private final long dynamicAttributeSchemaRefreshNanos; + private volatile DynamicAttributeSchemaSnapshot dynamicAttributeSchemaSnapshot; + @Autowired public GreptimeTraceQueryRepository( ObjectProvider<GreptimeSqlQueryExecutor> greptimeSqlQueryExecutorProvider, GreptimeProperties greptimeProperties, @Qualifier(WarehouseConstants.GREPTIME_QUERY_REST_TEMPLATE) RestTemplate restTemplate) { + this(greptimeSqlQueryExecutorProvider, greptimeProperties, restTemplate, + System::nanoTime, DYNAMIC_ATTRIBUTE_SCHEMA_REFRESH_NANOS); + } + + GreptimeTraceQueryRepository( + ObjectProvider<GreptimeSqlQueryExecutor> greptimeSqlQueryExecutorProvider, + GreptimeProperties greptimeProperties, + RestTemplate restTemplate, + LongSupplier monotonicNanos, + long dynamicAttributeSchemaRefreshNanos) { this.greptimeSqlQueryExecutorProvider = greptimeSqlQueryExecutorProvider; this.greptimeProperties = greptimeProperties; this.restTemplate = restTemplate; + this.monotonicNanos = monotonicNanos; + this.dynamicAttributeSchemaRefreshNanos = Math.max(1L, dynamicAttributeSchemaRefreshNanos); } @Override @@ -194,17 +222,13 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { String errorExpression = "SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') " + "THEN 1 ELSE 0 END)"; String serviceNamespaceExpression = resourceAttributeExpression(null, "service.namespace"); - String serviceNamespaceProjection = StringUtils.hasText(serviceNamespaceExpression) - ? "MAX(" + serviceNamespaceExpression + ")" - : "NULL"; - String resourceAttributesProjection = traceTableColumns().contains(RESOURCE_ATTRIBUTES_COLUMN) - ? "MAX(" + RESOURCE_ATTRIBUTES_COLUMN + ")" - : "NULL"; StringBuilder innerSql = new StringBuilder("SELECT ") .append("trace_id, ") .append("MAX(span_id) AS root_span_id, ") .append("MAX(service_name) AS service_name, ") - .append(serviceNamespaceProjection) + .append("MAX(") + .append(serviceNamespaceExpression) + .append(")") .append(" AS service_namespace, ") .append("MAX(span_name) AS root_span_name, ") .append("MAX(duration_nano) AS duration_nano, ") @@ -214,8 +238,8 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { .append("MIN(timestamp) AS timestamp, ") .append(errorExpression) .append(" AS error_span_count, ") - .append(resourceAttributesProjection) - .append(" AS resource_attributes ") + .append(traceListResourceAttributeProjections()) + .append(' ') .append("FROM ") .append(TRACE_TABLE); List<String> filters = new LinkedList<>(); @@ -912,17 +936,13 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { String expression = spanAttributeExpression(key); return StringUtils.hasText(expression) ? expression + " = '" + escapeSql(value.trim()) + "'" - : null; + : "1 = 0"; } private String spanAttributeExpression(String key) { - Set<String> columns = traceTableColumns(); String normalizedKey = key.trim(); - String flattenedColumn = "span_attributes." + normalizedKey; - if (columns.contains(flattenedColumn)) { - return qualifiedColumn(null, flattenedColumn); - } - return "json_get_string(span_attributes, '$[\"" + escapeJsonPathKey(normalizedKey) + "\"]')"; + String column = "span_attributes." + normalizedKey; + return dynamicAttributeColumnExists(column) ? qualifiedColumn(null, column) : null; } private List<Map<String, Object>> queryRows(String sql) { @@ -1040,18 +1060,7 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { private String workspaceFilter(String alias, String workspaceId) { String normalizedWorkspaceId = workspaceId.trim(); String canonicalWorkspace = resourceAttributeExpression(alias, "hertzbeat.workspace_id"); - if (!StringUtils.hasText(canonicalWorkspace)) { - throw new TelemetryStorageUnavailableException(); - } - String hertzbeatWorkspaceFilter = canonicalWorkspace + " = '" - + escapeSql(normalizedWorkspaceId) + "'"; - String legacyWorkspace = resourceAttributeExpression(alias, "workspace.id"); - if (!StringUtils.hasText(legacyWorkspace)) { - return hertzbeatWorkspaceFilter; - } - String workspaceIdFilter = legacyWorkspace + " = '" + escapeSql(normalizedWorkspaceId) + "'"; - String canonicalMissing = "(" + canonicalWorkspace + " IS NULL OR " + canonicalWorkspace + " = '')"; - return "(" + hertzbeatWorkspaceFilter + " OR (" + canonicalMissing + " AND " + workspaceIdFilter + "))"; + return canonicalWorkspace + " = '" + escapeSql(normalizedWorkspaceId) + "'"; } private String resourceAttributeAnyFilter(String alias, String key, Collection<String> values) { @@ -1213,24 +1222,17 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { private String resourceAttributeFilter(String alias, String key, String value) { String expression = resourceAttributeExpression(alias, key); if (!StringUtils.hasText(expression)) { - return null; + return "1 = 0"; } return expression + " = '" + escapeSql(value.trim()) + "'"; } private String resourceAttributeExpression(String alias, String key) { - Set<String> columns = traceTableColumns(); String normalizedKey = key.trim(); - String flattenedColumn = RESOURCE_ATTRIBUTES_COLUMN + "." + normalizedKey; - if (columns.contains(flattenedColumn)) { - return qualifiedColumn(alias, flattenedColumn); - } - if (columns.contains(RESOURCE_ATTRIBUTES_COLUMN)) { - String column = StringUtils.hasText(alias) - ? alias + "." + RESOURCE_ATTRIBUTES_COLUMN - : RESOURCE_ATTRIBUTES_COLUMN; - return "json_get_string(" + column + ", '$[\"" + escapeJsonPathKey(normalizedKey) + "\"]')"; + String column = "resource_attributes." + normalizedKey; + if (STABLE_RESOURCE_ATTRIBUTE_KEYS.contains(normalizedKey) || dynamicAttributeColumnExists(column)) { + return qualifiedColumn(alias, column); } return null; } @@ -1244,35 +1246,90 @@ public class GreptimeTraceQueryRepository implements TraceQueryRepository { return "\"" + column.replace("\"", "\"\"") + "\""; } - private Set<String> traceTableColumns() { - Set<String> cachedColumns = traceTableColumns; - if (cachedColumns != null) { - return cachedColumns; + private String traceListResourceAttributeProjections() { + return TRACE_LIST_RESOURCE_ATTRIBUTE_KEYS.stream() + .map(key -> { + String column = resourceAttributeExpression(null, key); + return "MAX(" + column + ") AS " + quoteIdentifier("resource_attributes." + key); + }) + .collect(java.util.stream.Collectors.joining(", ")); + } + + /** + * Discovers schema-less long-tail columns with a hard result bound. This is not a + * version-compatibility probe: stable HertzBeat dimensions never depend on it. Missing columns + * trigger a throttled refresh so native pipeline schema evolution becomes visible without + * allowing concurrent or repeated user queries to issue an unbounded number of DESC requests. + */ + private boolean dynamicAttributeColumnExists(String column) { + DynamicAttributeSchemaSnapshot snapshot = dynamicAttributeSchemaSnapshot; + if (snapshot != null && snapshot.columns().contains(column)) { + return true; + } + long now = monotonicNanos.getAsLong(); + if (snapshot == null || refreshDue(snapshot, now)) { + snapshot = refreshDynamicAttributeColumns(now); + } + if (!snapshot.available()) { + throw new TelemetryStorageUnavailableException(); } - Set<String> discoveredColumns = new LinkedHashSet<>(); - for (Map<String, Object> row : queryRows("DESC " + TRACE_TABLE)) { - String column = readText(row, "Column"); - if (StringUtils.hasText(column)) { - discoveredColumns.add(column); + return snapshot.columns().contains(column); + } + + private boolean refreshDue(DynamicAttributeSchemaSnapshot snapshot, long now) { + return now - snapshot.refreshedAtNanos() >= dynamicAttributeSchemaRefreshNanos; + } + + private DynamicAttributeSchemaSnapshot refreshDynamicAttributeColumns(long requestedAtNanos) { + DynamicAttributeSchemaSnapshot cached = dynamicAttributeSchemaSnapshot; + if (cached != null && !refreshDue(cached, requestedAtNanos)) { + return cached; + } + synchronized (this) { + cached = dynamicAttributeSchemaSnapshot; + long refreshedAtNanos = monotonicNanos.getAsLong(); + if (cached != null && !refreshDue(cached, refreshedAtNanos)) { + return cached; + } + try { + Set<String> discovered = new LinkedHashSet<>(); + for (Map<String, Object> row : queryRows("DESC " + TRACE_TABLE)) { + if (discovered.size() >= MAX_DISCOVERED_DYNAMIC_ATTRIBUTE_COLUMNS) { + break; + } + String column = discoveredColumnName(row); + if (StringUtils.hasText(column) + && (column.startsWith("resource_attributes.") || column.startsWith("span_attributes."))) { + discovered.add(column); + } + } + cached = new DynamicAttributeSchemaSnapshot( + Collections.unmodifiableSet(discovered), refreshedAtNanos, true); + dynamicAttributeSchemaSnapshot = cached; + return cached; + } catch (TelemetryStorageUnavailableException ex) { + Set<String> staleColumns = cached == null ? Collections.emptySet() : cached.columns(); + dynamicAttributeSchemaSnapshot = new DynamicAttributeSchemaSnapshot( + staleColumns, refreshedAtNanos, false); + throw ex; } } - if (discoveredColumns.isEmpty()) { - discoveredColumns.add(RESOURCE_ATTRIBUTES_COLUMN); - } - Set<String> immutableColumns = Collections.unmodifiableSet(discoveredColumns); - traceTableColumns = immutableColumns; - return immutableColumns; } - private String readText(Map<String, Object> row, String key) { - if (row == null || !row.containsKey(key)) { + private String discoveredColumnName(Map<String, Object> row) { + if (CollectionUtils.isEmpty(row)) { return null; } - String value = String.valueOf(row.get(key)).trim(); - return StringUtils.hasText(value) ? value : null; + for (Map.Entry<String, Object> entry : row.entrySet()) { + if (("column".equalsIgnoreCase(entry.getKey()) || "column_name".equalsIgnoreCase(entry.getKey())) + && entry.getValue() != null) { + return String.valueOf(entry.getValue()).trim(); + } + } + return null; } - private String escapeJsonPathKey(String key) { - return key.replace("\\", "\\\\").replace("\"", "\\\""); + private record DynamicAttributeSchemaSnapshot(Set<String> columns, long refreshedAtNanos, boolean available) { } + } diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepositoryTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepositoryTest.java index aa15725dac..488bde93b8 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepositoryTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/GreptimeTraceQueryRepositoryTest.java @@ -28,8 +28,8 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.verify; import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import ch.qos.logback.classic.Logger; @@ -38,9 +38,17 @@ import ch.qos.logback.core.read.ListAppender; import java.net.URLDecoder; import java.nio.charset.StandardCharsets; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.warehouse.db.GreptimeSqlQueryExecutor; import org.apache.hertzbeat.warehouse.repository.TraceQueryRepository.TraceRowQuery; @@ -126,23 +134,53 @@ class GreptimeTraceQueryRepositoryTest { assertNotNull(rows); assertEquals(1, rows.size()); ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); String sql = sqlCaptor.getValue(); assertTraceSqlProjectsAttribution(sql); assertTrue(sql.contains("timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("timestamp <= to_timestamp_millis(1710003600000)")); assertTrue(sql.contains("service_name = 'checkout'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.endsWith("ORDER BY timestamp DESC LIMIT 50")); } @Test - void queryRecentTraceRowsPushesWorkspaceAndEntityScopeIntoGreptimeSql() { + void queryRecentTraceRowsUsesOnlyTheCanonicalFlattenedTraceContract() { when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of("trace_id", "trace-1"))); + repository.queryRecentTraceRows( + 20, + 1710000000000L, + 1710003600000L, + "checkout", + "commerce", + "prod", + null, + null, + null, + "team-a", + Map.of("hertzbeat.entity_id", Set.of("entity-1")), + false); + + ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); + String sql = sqlCaptor.getValue(); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" = 'commerce'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.workspace_id\" = 'team-a'")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.entity_id\" = 'entity-1'")); + assertFalse(sql.contains("json_get_string(")); + assertFalse(sql.startsWith("DESC hzb_traces")); + } + + @Test + void queryRecentTraceRowsPushesWorkspaceAndEntityScopeIntoGreptimeSql() { + stubDynamicTraceQuery( + List.of(Map.of("trace_id", "trace-1")), + "resource_attributes.host.name"); + List<Map<String, Object>> rows = repository.queryRecentTraceRows( 75, 1710000000000L, @@ -163,22 +201,16 @@ class GreptimeTraceQueryRepositoryTest { assertNotNull(rows); assertEquals(1, rows.size()); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTraceSqlProjectsAttribution(sql); assertTrue(sql.contains("timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("timestamp <= to_timestamp_millis(1710003600000)")); assertTrue(sql.contains("service_name = 'checkout'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') = 'team-a' " - + "OR ((json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') IS NULL " - + "OR json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') = '') " - + "AND json_get_string(resource_attributes, '$[\"workspace.id\"]') = 'team-a'))")); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" = 'commerce'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertCanonicalWorkspaceFilter(sql, "team-a"); + assertTrue(sql.contains("(\"resource_attributes.host.name\" = 'checkout-1' " + + "OR \"resource_attributes.host.name\" = 'checkout-2')")); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.endsWith("ORDER BY timestamp DESC LIMIT 75")); } @@ -186,47 +218,25 @@ class GreptimeTraceQueryRepositoryTest { @Test void workspaceQueryUsesCanonicalFlattenedColumnWithoutRequiringLegacyAlias() { when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { - String sql = invocation.getArgument(0); - if (sql.startsWith("DESC hzb_traces")) { - return List.of( - Map.of("Column", "timestamp"), - Map.of("Column", "resource_attributes.hertzbeat.workspace_id")); - } - return List.of(Map.of("trace_id", "trace-1")); - }); + when(greptimeSqlQueryExecutor.executeStrict(anyString())) + .thenReturn(List.of(Map.of("trace_id", "trace-1"))); repository.queryRecentTraceRows( 20, 1000L, 2000L, null, null, null, null, null, null, "team-a", Map.of(), false); ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); + String sql = sqlCaptor.getValue(); assertTrue(sql.contains("\"resource_attributes.hertzbeat.workspace_id\" = 'team-a'")); assertFalse(sql.contains("workspace.id")); } - @Test - void workspaceQueryFailsBeforeDataReadWhenCanonicalWorkspaceCannotBeExpressed() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of("Column", "timestamp"))); - - assertThrows(TelemetryStorageUnavailableException.class, () -> repository.queryRecentTraceRows( - 20, 1000L, 2000L, null, null, null, null, null, null, - "team-a", Map.of(), false)); - - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); - assertTrue(sqlCaptor.getValue().startsWith("DESC hzb_traces")); - } - @Test void queryTraceListRowsPushesGroupingPaginationAndTotalCountIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "trace_id", "trace-1", - "total_count", 42L))); + stubDynamicTraceQuery( + List.of(Map.of("trace_id", "trace-1", "total_count", 42L)), + "resource_attributes.host.name"); List<Map<String, Object>> rows = repository.queryTraceListRows( 1710000000000L, @@ -246,27 +256,28 @@ class GreptimeTraceQueryRepositoryTest { assertNotNull(rows); assertEquals(1, rows.size()); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.contains("COUNT(*) OVER () AS total_count")); assertTrue(sql.contains("FROM (SELECT trace_id")); assertTrue(sql.contains("SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') " + "THEN 1 ELSE 0 END) AS error_span_count")); - assertEquals(1, sql.split(" AS resource_attributes", -1).length - 1, - "trace-list SQL must project resource_attributes with exactly one alias"); + assertTrue(sql.contains("MAX(\"resource_attributes.hertzbeat.workspace_id\") " + + "AS \"resource_attributes.hertzbeat.workspace_id\"")); + assertTrue(sql.contains("MAX(\"resource_attributes.hertzbeat.entity_id\") " + + "AS \"resource_attributes.hertzbeat.entity_id\"")); + assertTrue(sql.contains("MAX(\"resource_attributes.hertzbeat.entity_type\") " + + "AS \"resource_attributes.hertzbeat.entity_type\"")); assertTrue(sql.contains("FROM hzb_traces WHERE timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("timestamp <= to_timestamp_millis(1710003600000)")); assertTrue(sql.contains("service_name = 'checkout'")); assertTrue(sql.contains("span_name = 'GET /checkout'")); assertTrue(sql.contains("duration_nano >= 100000000")); assertTrue(sql.contains("duration_nano <= 500000000")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertCanonicalWorkspacePrecedence(sql, "team-a"); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" = 'commerce'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertCanonicalWorkspaceFilter(sql, "team-a"); + assertTrue(sql.contains("(\"resource_attributes.host.name\" = 'checkout-1' " + + "OR \"resource_attributes.host.name\" = 'checkout-2')")); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.contains("GROUP BY trace_id")); assertTrue(sql.endsWith("ORDER BY timestamp DESC LIMIT 20 OFFSET 40")); @@ -298,7 +309,7 @@ class GreptimeTraceQueryRepositoryTest { assertEquals(1, rows.size()); ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); String sql = sqlCaptor.getValue(); assertTrue(sql.contains("span_name = 'POST /checkout'")); assertTrue(sql.contains("duration_nano >= 100000000")); @@ -309,11 +320,12 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceOverviewRowsPushesAggregateFiltersIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "total_trace_count", 42L, - "error_trace_count", 7L, - "latest_observed_at", 1710003600000L))); + stubDynamicTraceQuery( + List.of(Map.of( + "total_trace_count", 42L, + "error_trace_count", 7L, + "latest_observed_at", 1710003600000L)), + "resource_attributes.host.name"); Map<String, Object> overview = repository.queryTraceOverviewRows( 1710000000000L, @@ -331,9 +343,7 @@ class GreptimeTraceQueryRepositoryTest { "entrypoint"); assertEquals(42L, overview.get("total_trace_count")); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.startsWith("SELECT COUNT(*) AS total_trace_count")); assertTrue(sql.contains("SUM(CASE WHEN error_span_count > 0 THEN 1 ELSE 0 END) AS error_trace_count")); assertTrue(sql.contains("MAX(trace_start_time) AS latest_observed_at")); @@ -348,12 +358,7 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("duration_nano <= 500000000")); assertTrue(sql.contains("(parent_span_id IS NULL OR parent_span_id = '' " + "OR UPPER(span_kind) IN ('SPAN_KIND_SERVER', 'SERVER', 'SPAN_KIND_CONSUMER', 'CONSUMER'))")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertCanonicalWorkspacePrecedence(sql, "team-a"); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertCanonicalFlattenedResourceFilters(sql); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.contains("GROUP BY trace_id HAVING " + "SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') THEN 1 ELSE 0 END) > 0")); @@ -362,11 +367,12 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceIdOverviewRowsPushesTraceIdAggregateFiltersIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "total_trace_count", 1L, - "error_trace_count", 1L, - "latest_observed_at", 1710003600000L))); + stubDynamicTraceQuery( + List.of(Map.of( + "total_trace_count", 1L, + "error_trace_count", 1L, + "latest_observed_at", 1710003600000L)), + "resource_attributes.host.name"); Map<String, Object> overview = repository.queryTraceIdOverviewRows( "trace-'filtered", @@ -385,9 +391,7 @@ class GreptimeTraceQueryRepositoryTest { "root"); assertEquals(1L, overview.get("total_trace_count")); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.startsWith("SELECT COUNT(*) AS total_trace_count")); assertTrue(sql.contains("SUM(CASE WHEN error_span_count > 0 THEN 1 ELSE 0 END) AS error_trace_count")); assertTrue(sql.contains("MAX(trace_start_time) AS latest_observed_at")); @@ -400,12 +404,7 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("duration_nano >= 200000000")); assertTrue(sql.contains("duration_nano <= 900000000")); assertTrue(sql.contains("(parent_span_id IS NULL OR parent_span_id = '')")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertCanonicalWorkspacePrecedence(sql, "team-a"); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertCanonicalFlattenedResourceFilters(sql); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.contains("GROUP BY trace_id HAVING " + "SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') THEN 1 ELSE 0 END) > 0")); @@ -414,12 +413,13 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceSummaryRowsPushesAggregateAndLatestTraceIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "total_trace_count", 7L, - "error_trace_count", 2L, - "latest_observed_at", 1710003600000L, - "latest_trace_id", "trace-latest"))); + stubDynamicTraceQuery( + List.of(Map.of( + "total_trace_count", 7L, + "error_trace_count", 2L, + "latest_observed_at", 1710003600000L, + "latest_trace_id", "trace-latest")), + "resource_attributes.host.name"); Map<String, Object> summary = repository.queryTraceSummaryRows( 1710000000000L, @@ -432,9 +432,7 @@ class GreptimeTraceQueryRepositoryTest { true); assertEquals("trace-latest", summary.get("latest_trace_id")); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.startsWith("SELECT summary.total_trace_count")); assertTrue(sql.contains("summary.error_trace_count")); assertTrue(sql.contains("latest.trace_start_time AS latest_observed_at")); @@ -447,12 +445,7 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("FROM hzb_traces WHERE timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("timestamp <= to_timestamp_millis(1710003600000)")); assertTrue(sql.contains("service_name = 'checkout'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertCanonicalWorkspacePrecedence(sql, "team-a"); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertCanonicalFlattenedResourceFilters(sql); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.contains("ORDER BY trace_start_time DESC LIMIT 1")); assertTrue(sql.endsWith(") latest ON TRUE")); @@ -460,13 +453,15 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceGroupByRowsAggregatesByTraceBeforeGroupingFieldValues() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "group_value", "1.2.3", - "trace_count", 12L, - "error_trace_count", 2L, - "latency_avg_ms", 84.5d, - "latency_p95_ms", 210.0d))); + stubDynamicTraceQuery( + List.of(Map.of( + "group_value", "1.2.3", + "trace_count", 12L, + "error_trace_count", 2L, + "latency_avg_ms", 84.5d, + "latency_p95_ms", 210.0d)), + "resource_attributes.host.name", + "resource_attributes.service.version"); List<Map<String, Object>> rows = repository.queryTraceGroupByRows( 1710000000000L, @@ -488,9 +483,7 @@ class GreptimeTraceQueryRepositoryTest { 7); assertEquals("1.2.3", rows.getFirst().get("group_value")); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getValue(); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.startsWith("SELECT group_value, COUNT(*) AS trace_count")); assertTrue(sql.contains("SUM(CASE WHEN error_span_count > 0 THEN 1 ELSE 0 END) AS error_trace_count")); assertTrue(sql.contains("COALESCE(SUM(duration_nano), 0) / NULLIF(COUNT(duration_nano), 0) " @@ -498,7 +491,7 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("uddsketch_calc(0.95, uddsketch_state(128, 0.01, duration_nano)) " + "/ 1000000.0 AS latency_p95_ms")); assertTrue(sql.contains("FROM (SELECT trace_id, " - + "COALESCE(NULLIF(MAX(json_get_string(resource_attributes, '$[\"service.version\"]')), ''), " + + "COALESCE(NULLIF(MAX(\"resource_attributes.service.version\"), ''), " + "'unknown') AS group_value")); assertTrue(sql.contains("MAX(duration_nano) AS duration_nano")); assertTrue(sql.contains("SUM(CASE WHEN span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') " @@ -511,12 +504,7 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("span_name = 'GET /checkout'")); assertTrue(sql.contains("duration_nano >= 100000000")); assertTrue(sql.contains("duration_nano <= 500000000")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') = 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertCanonicalWorkspacePrecedence(sql, "team-a"); - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1' " - + "OR json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-2')")); + assertCanonicalFlattenedResourceFilters(sql); assertTrue(sql.endsWith("GROUP BY group_value HAVING COUNT(*) >= 5 ORDER BY latency_p95_ms DESC LIMIT 7")); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.contains("GROUP BY trace_id HAVING " @@ -527,10 +515,11 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceServiceGraphRowsPushesServiceGraphRedAggregationIntoGreptimeSql() { when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "source_service_name", "checkout-api", - "target_service_name", "payment-api", - "request_count", 2L))); + when(greptimeSqlQueryExecutor.executeStrict(anyString())) + .thenReturn(List.of(Map.of( + "source_service_name", "checkout-api", + "target_service_name", "payment-api", + "request_count", 2L))); List<Map<String, Object>> rows = repository.queryTraceServiceGraphRows( 100, 1710000000000L, 1710003600000L, "prod", List.of("checkout-api", "payment-api"), true); @@ -538,10 +527,10 @@ class GreptimeTraceQueryRepositoryTest { assertNotNull(rows); assertEquals(1, rows.size()); ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); String sql = sqlCaptor.getValue(); assertTrue(sql.startsWith("SELECT parent.service_name AS source_service_name, " - + "child.service_name AS target_service_name")); + + "child.service_name AS target_service_name, COUNT(*) AS request_count")); assertTrue(sql.contains("COUNT(*) AS request_count")); assertTrue(sql.contains("SUM(CASE WHEN child.span_status_code IN ('STATUS_CODE_ERROR', 'ERROR') " + "THEN 1 ELSE 0 END) AS error_count")); @@ -562,10 +551,8 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("child.timestamp <= to_timestamp_millis(1710003600000)")); assertTrue(sql.contains("parent.timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("parent.timestamp <= to_timestamp_millis(1710003600000)")); - assertTrue(sql.contains("json_get_string(child.resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertTrue(sql.contains("json_get_string(parent.resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); + assertTrue(sql.contains("child.\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertTrue(sql.contains("parent.\"resource_attributes.deployment.environment.name\" = 'prod'")); assertTrue(sql.contains("((child.service_name = 'checkout-api' OR child.service_name = 'payment-api') " + "OR (parent.service_name = 'checkout-api' OR parent.service_name = 'payment-api'))")); assertTrue(sql.contains("LOWER(child.service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); @@ -576,71 +563,25 @@ class GreptimeTraceQueryRepositoryTest { } @Test - void scopedServiceGraphFiltersBothSidesByCanonicalWorkspacePrecedence() { + void scopedServiceGraphFiltersBothSidesByCanonicalWorkspaceColumn() { when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of( - "source_service_name", "checkout-api", - "target_service_name", "payment-api", - "request_count", 2L))); + when(greptimeSqlQueryExecutor.executeStrict(anyString())) + .thenReturn(List.of(Map.of( + "source_service_name", "checkout-api", + "target_service_name", "payment-api", + "request_count", 2L))); repository.queryTraceServiceGraphRows( 100, 1710000000000L, 1710003600000L, "prod", "team-a", List.of("checkout-api", "payment-api"), true); ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); - assertTrue(sql.contains("json_get_string(child.resource_attributes, " - + "'$[\"hertzbeat.workspace_id\"]') = 'team-a'")); - assertTrue(sql.contains("json_get_string(parent.resource_attributes, " - + "'$[\"hertzbeat.workspace_id\"]') = 'team-a'")); - assertTrue(sql.contains("json_get_string(child.resource_attributes, '$[\"workspace.id\"]') = 'team-a'")); - assertTrue(sql.contains("json_get_string(parent.resource_attributes, '$[\"workspace.id\"]') = 'team-a'")); - } - - @Test - void scopedServiceGraphIsUnavailableBeforeDataSelectWhenCanonicalWorkspaceCannotBeExpressed() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of( - Map.of("Column", "timestamp"), - Map.of("Column", "service_name"), - Map.of("Column", "resource_attributes.workspace.id"))); - - assertThrows(TelemetryStorageUnavailableException.class, () -> repository.queryTraceServiceGraphRows( - 100, 1710000000000L, 1710003600000L, "prod", "team-a", - List.of("checkout-api", "payment-api"), true)); - - verify(greptimeSqlQueryExecutor).executeStrict("DESC hzb_traces"); - } - - @Test - void queryTraceServiceGraphRowsUsesFlattenedResourceAttributeColumnsWhenGreptimeSchemaHasNoJsonColumn() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { - String sql = invocation.getArgument(0); - if (sql.startsWith("DESC hzb_traces")) { - return List.of( - Map.of("Column", "timestamp"), - Map.of("Column", "service_name"), - Map.of("Column", "resource_attributes.deployment.environment.name"), - Map.of("Column", "resource_attributes.host.name") - ); - } - return List.of(Map.of( - "source_service_name", "checkout-api", - "target_service_name", "payment-api", - "request_count", 2L)); - }); - - repository.queryTraceServiceGraphRows(100, 1710000000000L, 1710003600000L, "prod", true); - - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); - assertTrue(sql.contains("child.\"resource_attributes.deployment.environment.name\" = 'prod'")); - assertTrue(sql.contains("LOWER(child.service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); - assertTrue(sql.contains("LOWER(parent.service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); - assertTrue(!sql.contains("json_get_string(child.resource_attributes")); + verify(greptimeSqlQueryExecutor).executeStrict(sqlCaptor.capture()); + String sql = sqlCaptor.getValue(); + assertTrue(sql.contains("child.\"resource_attributes.hertzbeat.workspace_id\" = 'team-a'")); + assertTrue(sql.contains("parent.\"resource_attributes.hertzbeat.workspace_id\" = 'team-a'")); + assertFalse(sql.contains("workspace.id")); + assertFalse(sql.contains("json_get_string(")); } @Test @@ -726,8 +667,9 @@ class GreptimeTraceQueryRepositoryTest { @Test void queryTraceRowsPushesRouteFiltersIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of("trace_id", "trace-1"))); + stubDynamicTraceQuery( + List.of(Map.of("trace_id", "trace-1")), + "resource_attributes.host.name"); repository.queryTraceRows( "trace-'1", @@ -745,9 +687,7 @@ class GreptimeTraceQueryRepositoryTest { true ); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTraceSqlProjectsAttribution(sql); assertTrue(sql.contains("trace_id = 'trace-''1'")); assertTrue(sql.contains("timestamp >= to_timestamp_millis(1710000000000)")); @@ -756,22 +696,21 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("span_name = 'GET /checkout'")); assertTrue(sql.contains("duration_nano >= 100000000")); assertTrue(sql.contains("duration_nano <= 500000000")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.namespace\"]') " - + "= 'commerce'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"deployment.environment.name\"]') " - + "= 'prod'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') " - + "= 'team-a'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"workspace.id\"]') = 'team-a'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"host.name\"]') = 'checkout-1'")); + assertTrue(sql.contains("\"resource_attributes.service.namespace\" = 'commerce'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertCanonicalWorkspaceFilter(sql, "team-a"); + assertTrue(sql.contains("\"resource_attributes.host.name\" = 'checkout-1'")); + assertFalse(sql.contains("workspace.id")); + assertFalse(sql.contains("json_get_string(")); assertTrue(sql.contains("LOWER(service_name) NOT IN ('hertzbeat', 'apache-hertzbeat')")); assertTrue(sql.endsWith("ORDER BY timestamp ASC LIMIT 25")); } @Test void queryTraceRowsPushesTypedSpanAndAttributeContextIntoGreptimeSql() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of("trace_id", "trace-1"))); + stubDynamicTraceQuery( + List.of(Map.of("trace_id", "trace-1")), + "span_attributes.http.route"); repository.queryTraceRows(new TraceRowQuery( "trace-'1", @@ -792,26 +731,24 @@ class GreptimeTraceQueryRepositoryTest { false), 25); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); + String sql = captureMainSqlAfterDynamicDiscovery(); assertTrue(sql.contains("trace_id = 'trace-''1'")); assertTrue(sql.contains("span_id = 'span-''1'")); assertTrue(sql.contains("timestamp >= to_timestamp_millis(1710000000000)")); assertTrue(sql.contains("timestamp <= to_timestamp_millis(1710003600000)")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"hertzbeat.collector.id\"]') " - + "= 'collector-a'")); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.instance.id\"]') " - + "= 'checkout-7d9'")); - assertTrue(sql.contains("json_get_string(span_attributes, '$[\"http.route\"]') = '/checkout'")); + assertTrue(sql.contains("\"resource_attributes.hertzbeat.collector.id\" = 'collector-a'")); + assertTrue(sql.contains("\"resource_attributes.service.instance.id\" = 'checkout-7d9'")); + assertTrue(sql.contains("\"span_attributes.http.route\" = '/checkout'")); + assertFalse(sql.contains("json_get_string(")); assertTrue(sql.contains("duration_nano >= 100000000")); assertTrue(sql.contains("duration_nano <= 500000000")); } @Test void queryRecentTraceRowsPushesListAttributeContextBeforeApplyingLimit() { - when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); - when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenReturn(List.of(Map.of("trace_id", "trace-1"))); + stubDynamicTraceQuery( + List.of(Map.of("trace_id", "trace-1")), + "span_attributes.http.route"); repository.queryRecentTraceRows(new TraceRowQuery( null, @@ -830,15 +767,147 @@ class GreptimeTraceQueryRepositoryTest { false), 1500); - ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); - verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); - String sql = sqlCaptor.getAllValues().get(1); - assertTrue(sql.contains("json_get_string(resource_attributes, '$[\"service.instance.id\"]') " - + "= 'checkout-7d9'")); - assertTrue(sql.contains("json_get_string(span_attributes, '$[\"http.route\"]') = '/checkout'")); + String sql = captureMainSqlAfterDynamicDiscovery(); + assertTrue(sql.contains("\"resource_attributes.service.instance.id\" = 'checkout-7d9'")); + assertTrue(sql.contains("\"span_attributes.http.route\" = '/checkout'")); + assertFalse(sql.contains("json_get_string(")); assertTrue(sql.endsWith("ORDER BY timestamp DESC LIMIT 1500")); } + @Test + void missingDynamicAttributeColumnsProduceAnHonestEmptyQueryWithoutDroppingFilters() { + stubDynamicTraceQuery(List.of()); + + List<Map<String, Object>> rows = repository.queryRecentTraceRows(new TraceRowQuery( + null, + null, + 1710000000000L, + 1710003600000L, + "checkout", + null, + null, + null, + null, + null, + "team-a", + Map.of("custom.resource.key", Set.of("expected")), + Map.of("custom.span.key", Set.of("expected")), + false), + 20); + + assertTrue(rows.isEmpty()); + String sql = captureMainSqlAfterDynamicDiscovery(); + assertTrue(sql.contains("1 = 0")); + assertFalse(sql.contains("resource_attributes.custom.resource.key")); + assertFalse(sql.contains("span_attributes.custom.span.key")); + } + + @Test + void missingDynamicAttributeColumnIsDiscoveredAfterThrottledSchemaRefresh() { + AtomicLong now = new AtomicLong(); + AtomicInteger schemaProbeCount = new AtomicInteger(); + List<String> executedSql = new java.util.concurrent.CopyOnWriteArrayList<>(); + repository = new GreptimeTraceQueryRepository( + greptimeSqlQueryExecutorProvider, greptimeProperties, restTemplate, now::get, 100L); + when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); + when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { + String sql = invocation.getArgument(0); + executedSql.add(sql); + if (sql.startsWith("DESC hzb_traces")) { + return schemaProbeCount.incrementAndGet() == 1 + ? List.of() + : List.of(Map.of("Column", "resource_attributes.cloud.region")); + } + return List.of(); + }); + + queryByDynamicResourceAttribute("cloud.region"); + now.set(99L); + queryByDynamicResourceAttribute("cloud.region"); + now.set(100L); + queryByDynamicResourceAttribute("cloud.region"); + + assertEquals(2, schemaProbeCount.get()); + List<String> mainQueries = executedSql.stream() + .filter(sql -> !sql.startsWith("DESC hzb_traces")) + .toList(); + assertEquals(3, mainQueries.size()); + assertTrue(mainQueries.get(0).contains("1 = 0")); + assertTrue(mainQueries.get(1).contains("1 = 0")); + assertTrue(mainQueries.get(2).contains("\"resource_attributes.cloud.region\" = 'us-east-1'")); + assertFalse(mainQueries.get(2).contains("1 = 0")); + } + + @Test + void concurrentMissingDynamicAttributeQueriesShareOneSchemaProbeWithinThrottleWindow() throws Exception { + AtomicLong now = new AtomicLong(); + AtomicInteger schemaProbeCount = new AtomicInteger(); + repository = new GreptimeTraceQueryRepository( + greptimeSqlQueryExecutorProvider, greptimeProperties, restTemplate, now::get, 100L); + when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); + when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { + String sql = invocation.getArgument(0); + if (sql.startsWith("DESC hzb_traces")) { + schemaProbeCount.incrementAndGet(); + } + return List.of(); + }); + + int taskCount = 8; + CountDownLatch ready = new CountDownLatch(taskCount); + CountDownLatch start = new CountDownLatch(1); + ExecutorService executor = Executors.newFixedThreadPool(taskCount); + try { + List<Future<?>> futures = new ArrayList<>(); + for (int index = 0; index < taskCount; index++) { + futures.add(executor.submit(() -> { + ready.countDown(); + assertTrue(start.await(5, TimeUnit.SECONDS)); + queryByDynamicResourceAttribute("cloud.region"); + return null; + })); + } + assertTrue(ready.await(5, TimeUnit.SECONDS)); + start.countDown(); + for (Future<?> future : futures) { + future.get(5, TimeUnit.SECONDS); + } + } finally { + executor.shutdownNow(); + } + + assertEquals(1, schemaProbeCount.get()); + } + + @Test + void failedDynamicSchemaProbeIsThrottledAsControlledUnavailable() { + AtomicLong now = new AtomicLong(); + AtomicInteger schemaProbeCount = new AtomicInteger(); + repository = new GreptimeTraceQueryRepository( + greptimeSqlQueryExecutorProvider, greptimeProperties, restTemplate, now::get, 100L); + when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); + when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { + String sql = invocation.getArgument(0); + if (sql.startsWith("DESC hzb_traces")) { + schemaProbeCount.incrementAndGet(); + throw new IllegalStateException("schema unavailable"); + } + return List.of(); + }); + + assertThrows(TelemetryStorageUnavailableException.class, + () -> queryByDynamicResourceAttribute("cloud.region")); + now.set(99L); + assertThrows(TelemetryStorageUnavailableException.class, + () -> queryByDynamicResourceAttribute("cloud.region")); + assertEquals(1, schemaProbeCount.get()); + + now.set(100L); + assertThrows(TelemetryStorageUnavailableException.class, + () -> queryByDynamicResourceAttribute("cloud.region")); + assertEquals(2, schemaProbeCount.get()); + } + @Test void queryFailureLogsOnlyStableCategoryWithoutSqlOrThrowableBody() { String secretSentinel = "Bearer secret-token"; @@ -871,14 +940,54 @@ class GreptimeTraceQueryRepositoryTest { assertTrue(sql.contains("ORDER BY timestamp")); } - private void assertCanonicalWorkspacePrecedence(String sql, String workspaceId) { - assertTrue(sql.contains("(json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') = '" - + workspaceId - + "' OR ((json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') IS NULL " - + "OR json_get_string(resource_attributes, '$[\"hertzbeat.workspace_id\"]') = '') " - + "AND json_get_string(resource_attributes, '$[\"workspace.id\"]') = '" - + workspaceId - + "'))")); + private void assertCanonicalFlattenedResourceFilters(String sql) { + assertTrue(sql.contains("\"resource_attributes.service.namespace\" = 'commerce'")); + assertTrue(sql.contains("\"resource_attributes.deployment.environment.name\" = 'prod'")); + assertCanonicalWorkspaceFilter(sql, "team-a"); + assertTrue(sql.contains("(\"resource_attributes.host.name\" = 'checkout-1' " + + "OR \"resource_attributes.host.name\" = 'checkout-2')")); + assertFalse(sql.contains("json_get_string(")); + } + + private void assertCanonicalWorkspaceFilter(String sql, String workspaceId) { + assertTrue(sql.contains("\"resource_attributes.hertzbeat.workspace_id\" = '" + workspaceId + "'")); + } + + private void queryByDynamicResourceAttribute(String key) { + repository.queryRecentTraceRows(new TraceRowQuery( + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + Map.of(key, Set.of("us-east-1")), + Map.of(), + false), + 20); + } + + private void stubDynamicTraceQuery(List<Map<String, Object>> queryRows, String... dynamicColumns) { + List<Map<String, Object>> schemaRows = Arrays.stream(dynamicColumns) + .map(column -> Map.<String, Object>of("Column", column)) + .toList(); + when(greptimeSqlQueryExecutorProvider.getIfAvailable()).thenReturn(greptimeSqlQueryExecutor); + when(greptimeSqlQueryExecutor.executeStrict(anyString())).thenAnswer(invocation -> { + String sql = invocation.getArgument(0); + return sql.startsWith("DESC hzb_traces") ? schemaRows : queryRows; + }); + } + + private String captureMainSqlAfterDynamicDiscovery() { + ArgumentCaptor<String> sqlCaptor = ArgumentCaptor.forClass(String.class); + verify(greptimeSqlQueryExecutor, times(2)).executeStrict(sqlCaptor.capture()); + assertEquals("DESC hzb_traces", sqlCaptor.getAllValues().getFirst()); + return sqlCaptor.getAllValues().getLast(); } private GreptimeSqlQueryContent sqlResponse() { diff --git a/script/dev/seed-trace-rich-demo.sh b/script/dev/seed-trace-rich-demo.sh index c06d90a39c..eb23177a96 100755 --- a/script/dev/seed-trace-rich-demo.sh +++ b/script/dev/seed-trace-rich-demo.sh @@ -4,6 +4,9 @@ set -euo pipefail GREPTIME_USERNAME="${GREPTIME_USERNAME:-greptime}" GREPTIME_PASSWORD="${GREPTIME_PASSWORD:-greptime}" +SCRIPT_DIR="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd)" +REPOSITORY_ROOT="$(cd -- "${SCRIPT_DIR}/../.." && pwd)" +TRACE_TABLE_SQL_FILE="${REPOSITORY_ROOT}/hertzbeat-observability/src/main/resources/greptime/tables/hzb_traces.sql" resolve_greptime_http() { if [[ -n "${GREPTIME_HTTP:-}" ]]; then @@ -36,15 +39,46 @@ to_ns() { printf '%s\n' "$(( $1 * 1000000 ))" } +hex_to_base64() { + HEX_VALUE="$1" python3 - <<'PY' +import base64 +import os + +print(base64.b64encode(bytes.fromhex(os.environ["HEX_VALUE"])).decode()) +PY +} + GREPTIME_HTTP="$(resolve_greptime_http)" GREPTIME_DB_NAME="${GREPTIME_DB_NAME:-public}" TRACE_WINDOW_END_MS="${TRACE_WINDOW_END_MS:-$(resolve_now_ms)}" TRACE_WINDOW_END_MS="$((TRACE_WINDOW_END_MS / 60000 * 60000))" TRACE_WINDOW_START_MS="${TRACE_WINDOW_START_MS:-$((TRACE_WINDOW_END_MS - 600000))}" TRACE_ROOT_START_MS="${TRACE_ROOT_START_MS:-$((TRACE_WINDOW_START_MS + 180000))}" -TRACE_ID="${TRACE_ID:-trace-ui-rich-demo-${TRACE_WINDOW_END_MS}}" +TRACE_ID="${TRACE_ID:-$(TRACE_WINDOW_END_MS="${TRACE_WINDOW_END_MS}" python3 - <<'PY' +import hashlib +import os + +print(hashlib.sha256(("trace-ui-rich-demo-" + os.environ["TRACE_WINDOW_END_MS"]).encode()).hexdigest()[:32]) +PY +)}" +if [[ ! "${TRACE_ID}" =~ ^[[:xdigit:]]{32}$ ]]; then + printf 'TRACE_ID must be exactly 32 hexadecimal characters\n' >&2 + exit 1 +fi TRACE_ROUTE="http://127.0.0.1:4200/explore?signal=traces&traceId=${TRACE_ID}&start=${TRACE_WINDOW_START_MS}&end=${TRACE_WINDOW_END_MS}" +ROOT_SPAN_ID="0000000000000001" +AUTH_SPAN_ID="0000000000000002" +DB_SPAN_ID="0000000000000003" +CACHE_SPAN_ID="0000000000000004" +TEMPLATE_SPAN_ID="0000000000000005" +TRACE_ID_BASE64="$(hex_to_base64 "${TRACE_ID}")" +ROOT_SPAN_ID_BASE64="$(hex_to_base64 "${ROOT_SPAN_ID}")" +AUTH_SPAN_ID_BASE64="$(hex_to_base64 "${AUTH_SPAN_ID}")" +DB_SPAN_ID_BASE64="$(hex_to_base64 "${DB_SPAN_ID}")" +CACHE_SPAN_ID_BASE64="$(hex_to_base64 "${CACHE_SPAN_ID}")" +TEMPLATE_SPAN_ID_BASE64="$(hex_to_base64 "${TEMPLATE_SPAN_ID}")" + ROOT_START_NS="$(to_ns "${TRACE_ROOT_START_MS}")" ROOT_END_NS="$(to_ns "$((TRACE_ROOT_START_MS + 120))")" ROOT_EVENT_NS="$(to_ns "$((TRACE_ROOT_START_MS + 15))")" @@ -124,37 +158,11 @@ else: PY } -read -r -d '' create_table_sql <<'SQL' || true -CREATE TABLE IF NOT EXISTS "hzb_traces" ( - "timestamp" TIMESTAMP(9) NOT NULL, - "timestamp_end" TIMESTAMP(9) NULL, - "duration_nano" BIGINT UNSIGNED NULL, - "trace_id" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), - "span_id" STRING NULL, - "parent_span_id" STRING NULL, - "span_kind" STRING NULL, - "span_name" STRING NULL, - "span_status_code" STRING NULL, - "span_status_message" STRING NULL, - "trace_state" STRING NULL, - "scope_name" STRING NULL, - "scope_version" STRING NULL, - "service_name" STRING NULL SKIPPING INDEX WITH(granularity = '10240', type = 'BLOOM'), - "resource_attributes" JSON NULL, - "span_attributes" JSON NULL, - "span_events" JSON NULL, - "span_links" JSON NULL, - TIME INDEX ("timestamp"), - PRIMARY KEY ("service_name") -) -ENGINE=mito -WITH( - append_mode = 'true', - table_data_model = 'greptime_trace_v1' -) -SQL - -query_sql_checked "${create_table_sql}" >/tmp/greptime-trace-rich-demo-create.out +if [[ ! -r "${TRACE_TABLE_SQL_FILE}" ]]; then + printf 'Trace schema resource is not readable: %s\n' "${TRACE_TABLE_SQL_FILE}" >&2 + exit 1 +fi +query_sql_checked "$(<"${TRACE_TABLE_SQL_FILE}")" >/tmp/greptime-trace-rich-demo-create.out existing_rows="$(query_count "select count(*) as total from hzb_traces where trace_id = '${TRACE_ID}'")" if [[ "${existing_rows}" != "0" ]]; then @@ -163,50 +171,53 @@ if [[ "${existing_rows}" != "0" ]]; then exit 0 fi -RESOURCE_ATTRIBUTES_JSON='{"service.name":"checkout-service","service.namespace":"storefront","deployment.environment.name":"dev","service.version":"2026.04.01"}' -ROOT_SPAN_ATTRIBUTES_JSON='{"http.route":"/checkout"}' -AUTH_SPAN_ATTRIBUTES_JSON='{"http.route":"/checkout","auth.strategy":"session","user.id":"u-2048"}' -DB_SPAN_ATTRIBUTES_JSON='{"http.route":"/checkout","db.rows":3,"db.system":"mysql","retry.count":1}' -CACHE_SPAN_ATTRIBUTES_JSON='{"http.route":"/checkout","cache.key":"cart:summary","cache.hit":true}' -TEMPLATE_SPAN_ATTRIBUTES_JSON='{"http.route":"/checkout","template.name":"checkout-summary"}' - -ROOT_SPAN_EVENTS_JSON='[{"time_unix_nano":'"${ROOT_EVENT_NS}"',"name":"request.validated","attributes":{"http.route":"/checkout","user.segment":"vip"},"dropped_attributes_count":0}]' -AUTH_SPAN_EVENTS_JSON='[{"time_unix_nano":'"${AUTH_EVENT_NS}"',"name":"auth.user.loaded","attributes":{"auth.strategy":"session","user.id":"u-2048"},"dropped_attributes_count":0}]' -DB_SPAN_EVENTS_JSON='[{"time_unix_nano":'"${DB_EVENT_PREP_NS}"',"name":"db.statement.prepared","attributes":{"db.rows":3,"db.system":"mysql"},"dropped_attributes_count":0},{"time_unix_nano":'"${DB_EVENT_RETRY_NS}"',"name":"db.retry.success","attributes":{"retry.count":1},"dropped_attributes_count":0}]' -CACHE_SPAN_EVENTS_JSON='[{"time_unix_nano":'"${CACHE_EVENT_NS}"',"name":"cache.hit","attributes":{"cache.key":"cart:summary","cache.hit":true},"dropped_attributes_count":0}]' -TEMPLATE_SPAN_EVENTS_JSON='[{"time_unix_nano":'"${TEMPLATE_EVENT_NS}"',"name":"template.partial.rendered","attributes":{"template.name":"checkout-summary"},"dropped_attributes_count":0}]' - -AUTH_SPAN_LINKS_JSON='[{"trace_id":"'"${TRACE_ID}"'","span_id":"span-root-rich","trace_state":"vendor=greptime-demo","attributes":{"link.kind":"follows-from"},"dropped_attributes_count":0}]' -DB_SPAN_LINKS_JSON='[{"trace_id":"'"${TRACE_ID}"'","span_id":"span-root-rich","trace_state":"vendor=greptime-demo","attributes":{"link.kind":"caused-by"},"dropped_attributes_count":0}]' - -read -r -d '' insert_sql <<SQL || true -INSERT INTO hzb_traces ( - timestamp, - timestamp_end, - duration_nano, - trace_id, - span_id, - parent_span_id, - span_kind, - span_name, - span_status_code, - span_status_message, - trace_state, - scope_name, - scope_version, - service_name, - resource_attributes, - span_attributes, - span_events, - span_links -) VALUES - (${ROOT_START_NS}, ${ROOT_END_NS}, 120000000, '${TRACE_ID}', 'span-root-rich', NULL, 'SPAN_KIND_SERVER', 'GET /checkout', 'STATUS_CODE_OK', 'checkout rendered successfully', 'vendor=greptime-demo', 'checkout-api', '1.0.0', 'checkout-service', '${RESOURCE_ATTRIBUTES_JSON}', '${ROOT_SPAN_ATTRIBUTES_JSON}', '${ROOT_SPAN_EVENTS_JSON}', '[]'), - (${AUTH_START_NS}, ${AUTH_END_NS}, 25000000, '${TRACE_ID}', 'span-auth-rich', 'span-root-rich', 'SPAN_KIND_INTERNAL', 'AuthMiddleware', 'STATUS_CODE_OK', 'session token verified', 'vendor=greptime-demo', 'checkout-auth', '1.0.0', 'checkout-service', '${RESOURCE_ATTRIBUTES_JSON}', '${AUTH_SPAN_ATTRIBUTES_JSON}', '${AUTH_SPAN_EVENTS_JSON}', '${AUTH_SPAN_LINKS_JSON}'), - (${DB_START_NS}, ${DB_END_NS}, 55000000, '${TRACE_ID}', 'span-db-rich', 'span-root-rich', 'SPAN_KIND_CLIENT', 'SELECT cart_items', 'STATUS_CODE_ERROR', 'db timeout recovered', 'vendor=greptime-demo', 'checkout-db-instrumentation', '1.0.0', 'checkout-service', '${RESOURCE_ATTRIBUTES_JSON}', '${DB_SPAN_ATTRIBUTES_JSON}', '${DB_SPAN_EVENTS_JSON}', '${DB_SPAN_LINKS_JSON}'), - (${CACHE_START_NS}, ${CACHE_END_NS}, 12000000, '${TRACE_ID}', 'span-cache-rich', 'span-root-rich', 'SPAN_KIND_CLIENT', 'redis GET cart:summary', 'STATUS_CODE_OK', 'cache read completed', 'vendor=greptime-demo', 'checkout-cache', '1.0.0', 'checkout-service', '${RESOURCE_ATTRIBUTES_JSON}', '${CACHE_SPAN_ATTRIBUTES_JSON}', '${CACHE_SPAN_EVENTS_JSON}', '[]'), - (${TEMPLATE_START_NS}, ${TEMPLATE_END_NS}, 6000000, '${TRACE_ID}', 'span-template-rich', 'span-root-rich', 'SPAN_KIND_INTERNAL', 'RenderCheckoutSummary', 'STATUS_CODE_OK', 'template render finished', 'vendor=greptime-demo', 'checkout-renderer', '1.0.0', 'checkout-service', '${RESOURCE_ATTRIBUTES_JSON}', '${TEMPLATE_SPAN_ATTRIBUTES_JSON}', '${TEMPLATE_SPAN_EVENTS_JSON}', '[]'); -SQL - -query_sql_checked "${insert_sql}" >/tmp/greptime-trace-rich-demo.out +read -r -d '' otlp_payload <<JSON || true +{ + "resourceSpans": [{ + "resource": {"attributes": [ + {"key":"service.name","value":{"stringValue":"checkout-service"}}, + {"key":"service.namespace","value":{"stringValue":"storefront"}}, + {"key":"service.instance.id","value":{"stringValue":"checkout-rich-demo"}}, + {"key":"deployment.environment.name","value":{"stringValue":"dev"}}, + {"key":"hertzbeat.workspace_id","value":{"stringValue":"default"}}, + {"key":"hertzbeat.entity_id","value":{"stringValue":"checkout-service"}}, + {"key":"hertzbeat.entity_type","value":{"stringValue":"service"}}, + {"key":"hertzbeat.collector.id","value":{"stringValue":"local-demo"}}, + {"key":"service.version","value":{"stringValue":"2026.04.01"}} + ]}, + "scopeSpans": [{ + "scope": {"name":"checkout-rich-demo","version":"1.0.0"}, + "spans": [ + {"traceId":"${TRACE_ID_BASE64}","spanId":"${ROOT_SPAN_ID_BASE64}","traceState":"vendor=greptime-demo","name":"GET /checkout","kind":"SPAN_KIND_SERVER","startTimeUnixNano":"${ROOT_START_NS}","endTimeUnixNano":"${ROOT_END_NS}","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}}],"events":[{"timeUnixNano":"${ROOT_EVENT_NS}","name":"request.validated","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}},{"key":"user.segment","value":{"stringValue" [...] + {"traceId":"${TRACE_ID_BASE64}","spanId":"${AUTH_SPAN_ID_BASE64}","parentSpanId":"${ROOT_SPAN_ID_BASE64}","traceState":"vendor=greptime-demo","name":"AuthMiddleware","kind":"SPAN_KIND_INTERNAL","startTimeUnixNano":"${AUTH_START_NS}","endTimeUnixNano":"${AUTH_END_NS}","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}},{"key":"auth.strategy","value":{"stringValue":"session"}},{"key":"user.id","value":{"stringValue":"u-2048"}}],"events":[{"timeUnixNano":"${AUTH_E [...] + {"traceId":"${TRACE_ID_BASE64}","spanId":"${DB_SPAN_ID_BASE64}","parentSpanId":"${ROOT_SPAN_ID_BASE64}","traceState":"vendor=greptime-demo","name":"SELECT cart_items","kind":"SPAN_KIND_CLIENT","startTimeUnixNano":"${DB_START_NS}","endTimeUnixNano":"${DB_END_NS}","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}},{"key":"db.rows","value":{"intValue":"3"}},{"key":"db.system","value":{"stringValue":"mysql"}},{"key":"retry.count","value":{"intValue":"1"}}],"events [...] + {"traceId":"${TRACE_ID_BASE64}","spanId":"${CACHE_SPAN_ID_BASE64}","parentSpanId":"${ROOT_SPAN_ID_BASE64}","traceState":"vendor=greptime-demo","name":"redis GET cart:summary","kind":"SPAN_KIND_CLIENT","startTimeUnixNano":"${CACHE_START_NS}","endTimeUnixNano":"${CACHE_END_NS}","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}},{"key":"cache.key","value":{"stringValue":"cart:summary"}},{"key":"cache.hit","value":{"boolValue":true}}],"events":[{"timeUnixNano":"${ [...] + {"traceId":"${TRACE_ID_BASE64}","spanId":"${TEMPLATE_SPAN_ID_BASE64}","parentSpanId":"${ROOT_SPAN_ID_BASE64}","traceState":"vendor=greptime-demo","name":"RenderCheckoutSummary","kind":"SPAN_KIND_INTERNAL","startTimeUnixNano":"${TEMPLATE_START_NS}","endTimeUnixNano":"${TEMPLATE_END_NS}","attributes":[{"key":"http.route","value":{"stringValue":"/checkout"}},{"key":"template.name","value":{"stringValue":"checkout-summary"}}],"events":[{"timeUnixNano":"${TEMPLATE_EVENT_NS}","name":"t [...] + ] + }] + }] +} +JSON + +otlp_response="$(curl --fail-with-body -sS \ + -H "${AUTH_HEADER}" \ + -H 'Content-Type: application/json' \ + -H 'Accept: application/json' \ + -H "X-Greptime-DB-Name: ${GREPTIME_DB_NAME}" \ + -H 'X-Greptime-Trace-Table-Name: hzb_traces' \ + -H 'X-Greptime-Pipeline-Name: greptime_trace_v1' \ + --data-binary "${otlp_payload}" \ + "${GREPTIME_HTTP}/v1/otlp/v1/traces")" +OTLP_RESPONSE_JSON="${otlp_response}" python3 - <<'PY' +import json +import os + +body = json.loads(os.environ["OTLP_RESPONSE_JSON"] or "{}") +partial = body.get("partialSuccess") or {} +if int(partial.get("rejectedSpans", 0)) != 0: + raise SystemExit("Greptime rejected trace spans: " + json.dumps(body)) +PY + +printf '%s\n' "${otlp_response}" >/tmp/greptime-trace-rich-demo-otlp.out printf 'seeded trace demo: %s\n' "${TRACE_ID}" printf 'open: %s\n' "${TRACE_ROUTE}" --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
