From 78483991b780152a00f7f4032e0012638a1aba93 Mon Sep 17 00:00:00 2001 From: kezhenxu94 Date: Fri, 6 Sep 2024 22:01:22 +0800 Subject: [PATCH] Add self observability metrics for otel handler (#12598) --- docs/en/changes/changes.md | 1 + .../otel/otlp/OpenTelemetryLogHandler.java | 24 ++- .../OpenTelemetryMetricRequestProcessor.java | 75 +++++---- .../otel/otlp/OpenTelemetryTraceHandler.java | 68 ++++++--- .../src/main/resources/otel-rules/oap.yaml | 35 +++-- .../so11y_oap/so11y-instance.json | 142 ++++++++++++++++++ 6 files changed, 278 insertions(+), 67 deletions(-) diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 6fa5460bb2..db4d10bbf1 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -58,6 +58,7 @@ * Fix the compatibility with Grafana 11 when using label_values query variables. * Nacos as config server and cluster coordinator supports configuration contextPath. * Update the endpoint name format to `:` in eBPF Access Log Receiver. +* Add self-observability metrics for OpenTelemetry receiver. #### UI diff --git a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryLogHandler.java b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryLogHandler.java index 3d6c6bc9db..35c75b7ace 100644 --- a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryLogHandler.java +++ b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryLogHandler.java @@ -27,6 +27,7 @@ import io.opentelemetry.proto.common.v1.KeyValue; import io.opentelemetry.proto.logs.v1.LogRecord; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import lombok.Getter; import org.apache.skywalking.apm.network.common.v3.KeyStringValuePair; import org.apache.skywalking.apm.network.logging.v3.LogData; import org.apache.skywalking.apm.network.logging.v3.LogDataBody; @@ -39,6 +40,10 @@ import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.library.module.ModuleStartException; import org.apache.skywalking.oap.server.receiver.otel.Handler; import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; +import org.apache.skywalking.oap.server.telemetry.TelemetryModule; +import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; +import org.apache.skywalking.oap.server.telemetry.api.MetricsCreator; +import org.apache.skywalking.oap.server.telemetry.api.MetricsTag; import java.util.Map; import java.util.stream.Collectors; @@ -54,6 +59,17 @@ public class OpenTelemetryLogHandler private ILogAnalyzerService logAnalyzerService; + @Getter(lazy = true) + private final MetricsCreator metricsCreator = manager.find(TelemetryModule.NAME).provider().getService(MetricsCreator.class); + + @Getter(lazy = true) + private final HistogramMetrics processHistogram = getMetricsCreator().createHistogramMetric( + "otel_logs_latency", + "The latency to process the logs request", + MetricsTag.EMPTY_KEY, + MetricsTag.EMPTY_VALUE + ); + @Override public String type() { return "otlp-logs"; @@ -87,9 +103,11 @@ public class OpenTelemetryLogHandler .getScopeLogsList() .stream() .flatMap(it -> it.getLogRecordsList().stream()) - .forEach( - logRecord -> - doAnalysisQuietly(service, layer, serviceInstance, logRecord)); + .forEach(logRecord -> { + try (final var timer = getProcessHistogram().createTimer()) { + doAnalysisQuietly(service, layer, serviceInstance, logRecord); + } + }); responseObserver.onNext(ExportLogsServiceResponse.getDefaultInstance()); responseObserver.onCompleted(); }); diff --git a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryMetricRequestProcessor.java b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryMetricRequestProcessor.java index 0be3aaaf2c..7ef39c0a8f 100644 --- a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryMetricRequestProcessor.java +++ b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryMetricRequestProcessor.java @@ -25,6 +25,7 @@ import io.opentelemetry.proto.common.v1.KeyValue; import io.opentelemetry.proto.metrics.v1.Sum; import io.opentelemetry.proto.metrics.v1.SummaryDataPoint; import io.vavr.Function1; +import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.meter.analyzer.MetricConvert; @@ -42,6 +43,10 @@ import org.apache.skywalking.oap.server.library.util.prometheus.metrics.Histogra import org.apache.skywalking.oap.server.library.util.prometheus.metrics.Metric; import org.apache.skywalking.oap.server.library.util.prometheus.metrics.Summary; import org.apache.skywalking.oap.server.receiver.otel.OtelMetricReceiverConfig; +import org.apache.skywalking.oap.server.telemetry.TelemetryModule; +import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; +import org.apache.skywalking.oap.server.telemetry.api.MetricsCreator; +import org.apache.skywalking.oap.server.telemetry.api.MetricsTag; import java.io.IOException; import java.util.HashMap; @@ -73,39 +78,51 @@ public class OpenTelemetryMetricRequestProcessor implements Service { .build(); private List converters; + @Getter(lazy = true) + private final MetricsCreator metricsCreator = manager.find(TelemetryModule.NAME).provider().getService(MetricsCreator.class); + + @Getter(lazy = true) + private final HistogramMetrics processHistogram = getMetricsCreator().createHistogramMetric( + "otel_metrics_latency", + "The latency to process the metrics request", + MetricsTag.EMPTY_KEY, + MetricsTag.EMPTY_VALUE + ); + public void processMetricsRequest(final ExportMetricsServiceRequest requests) { - requests.getResourceMetricsList().forEach(request -> { - if (log.isDebugEnabled()) { - log.debug("Resource attributes: {}", request.getResource().getAttributesList()); - } + try (final var unused = getProcessHistogram().createTimer()) { + requests.getResourceMetricsList().forEach(request -> { + if (log.isDebugEnabled()) { + log.debug("Resource attributes: {}", request.getResource().getAttributesList()); + } - final Map nodeLabels = - request - .getResource() - .getAttributesList() - .stream() - .collect(toMap( - it -> LABEL_MAPPINGS - .getOrDefault(it.getKey(), it.getKey()) - .replaceAll("\\.", "_"), - it -> anyValueToString(it.getValue()), - (v1, v2) -> v1 - )); - - converters - .forEach(convert -> convert.toMeter( + final Map nodeLabels = request - .getScopeMetricsList().stream() - .flatMap(scopeMetrics -> scopeMetrics - .getMetricsList().stream() - .flatMap(metric -> adaptMetrics(nodeLabels, metric)) - .map(Function1.liftTry(Function.identity())) - .flatMap(tryIt -> MetricConvert.log( - tryIt, - "Convert OTEL metric to prometheus metric" - ))))); - }); + .getResource() + .getAttributesList() + .stream() + .collect(toMap( + it -> LABEL_MAPPINGS + .getOrDefault(it.getKey(), it.getKey()) + .replaceAll("\\.", "_"), + it -> anyValueToString(it.getValue()), + (v1, v2) -> v1 + )); + converters + .forEach(convert -> convert.toMeter( + request + .getScopeMetricsList().stream() + .flatMap(scopeMetrics -> scopeMetrics + .getMetricsList().stream() + .flatMap(metric -> adaptMetrics(nodeLabels, metric)) + .map(Function1.liftTry(Function.identity())) + .flatMap(tryIt -> MetricConvert.log( + tryIt, + "Convert OTEL metric to prometheus metric" + ))))); + }); + } } public void start() throws ModuleStartException { diff --git a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryTraceHandler.java b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryTraceHandler.java index 18e76bee34..f1b687ba50 100644 --- a/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryTraceHandler.java +++ b/oap-server/server-receiver-plugin/otel-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/otel/otlp/OpenTelemetryTraceHandler.java @@ -30,6 +30,7 @@ import io.opentelemetry.proto.common.v1.KeyValue; import io.opentelemetry.proto.resource.v1.Resource; import io.opentelemetry.proto.trace.v1.ScopeSpans; import io.opentelemetry.proto.trace.v1.Status; +import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; @@ -40,6 +41,10 @@ import org.apache.skywalking.oap.server.receiver.otel.Handler; import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; import org.apache.skywalking.oap.server.receiver.zipkin.SpanForwardService; import org.apache.skywalking.oap.server.receiver.zipkin.ZipkinReceiverModule; +import org.apache.skywalking.oap.server.telemetry.TelemetryModule; +import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; +import org.apache.skywalking.oap.server.telemetry.api.MetricsCreator; +import org.apache.skywalking.oap.server.telemetry.api.MetricsTag; import zipkin2.Endpoint; import zipkin2.Span; @@ -64,6 +69,17 @@ public class OpenTelemetryTraceHandler private final ModuleManager manager; private SpanForwardService forwardService; + @Getter(lazy = true) + private final MetricsCreator metricsCreator = manager.find(TelemetryModule.NAME).provider().getService(MetricsCreator.class); + + @Getter(lazy = true) + private final HistogramMetrics processHistogram = getMetricsCreator().createHistogramMetric( + "otel_spans_latency", + "The latency to process the span request", + MetricsTag.EMPTY_KEY, + MetricsTag.EMPTY_VALUE + ); + @Override public String type() { return "otlp-traces"; @@ -80,34 +96,36 @@ public class OpenTelemetryTraceHandler @Override public void export(ExportTraceServiceRequest request, StreamObserver responseObserver) { final ArrayList result = new ArrayList<>(); - request.getResourceSpansList().forEach(resourceSpans -> { - final Resource resource = resourceSpans.getResource(); - final List scopeSpansList = resourceSpans.getScopeSpansList(); - if (resource.getAttributesCount() == 0 && scopeSpansList.size() == 0) { - return; - } - final Map resourceTags = convertAttributeToMap(resource.getAttributesList()); - String serviceName = extractZipkinServiceName(resourceTags); - if (StringUtil.isEmpty(serviceName)) { - log.warn("No service name found in resource attributes, discarding the trace"); - return; - } - - try { - for (ScopeSpans scopeSpans : scopeSpansList) { - extractScopeTag(scopeSpans.getScope(), resourceTags); - for (io.opentelemetry.proto.trace.v1.Span span : scopeSpans.getSpansList()) { - Span zipkinSpan = convertSpan(span, serviceName, resourceTags); - result.add(zipkinSpan); - } + try (final var unused = getProcessHistogram().createTimer()) { + request.getResourceSpansList().forEach(resourceSpans -> { + final Resource resource = resourceSpans.getResource(); + final List scopeSpansList = resourceSpans.getScopeSpansList(); + if (resource.getAttributesCount() == 0 && scopeSpansList.size() == 0) { + return; + } + final Map resourceTags = convertAttributeToMap(resource.getAttributesList()); + String serviceName = extractZipkinServiceName(resourceTags); + if (StringUtil.isEmpty(serviceName)) { + log.warn("No service name found in resource attributes, discarding the trace"); + return; } - } catch (Exception e) { - log.warn("convert span error, discarding the span: {}", e.getMessage()); - } - }); - getForwardService().send(result); + try { + for (ScopeSpans scopeSpans : scopeSpansList) { + extractScopeTag(scopeSpans.getScope(), resourceTags); + for (io.opentelemetry.proto.trace.v1.Span span : scopeSpans.getSpansList()) { + Span zipkinSpan = convertSpan(span, serviceName, resourceTags); + result.add(zipkinSpan); + } + } + } catch (Exception e) { + log.warn("convert span error, discarding the span: {}", e.getMessage()); + } + }); + getForwardService().send(result); + } + responseObserver.onNext(ExportTraceServiceResponse.getDefaultInstance()); responseObserver.onCompleted(); } diff --git a/oap-server/server-starter/src/main/resources/otel-rules/oap.yaml b/oap-server/server-starter/src/main/resources/otel-rules/oap.yaml index 6de0e00b53..a8aa11ee2e 100644 --- a/oap-server/server-starter/src/main/resources/otel-rules/oap.yaml +++ b/oap-server/server-starter/src/main/resources/otel-rules/oap.yaml @@ -37,15 +37,17 @@ metricsRules: - name: instance_jvm_memory_bytes_used exp: jvm_memory_bytes_used.sum(['service', 'host_name']) - name: instance_jvm_gc_count - exp: "jvm_gc_collection_seconds_count.tagMatch('gc', 'PS Scavenge|Copy|ParNew|G1 Young Generation|PS MarkSweep|MarkSweepCompact|ConcurrentMarkSweep|G1 Old Generation') - .sum(['service', 'host_name', 'gc']).increase('PT1M') - .tag({tags -> if (tags['gc'] == 'PS Scavenge' || tags['gc'] == 'Copy' || tags['gc'] == 'ParNew' || tags['gc'] == 'G1 Young Generation') {tags.gc = 'young_gc_count'} }) - .tag({tags -> if (tags['gc'] == 'PS MarkSweep' || tags['gc'] == 'MarkSweepCompact' || tags['gc'] == 'ConcurrentMarkSweep' || tags['gc'] == 'G1 Old Generation') {tags.gc = 'old_gc_count'} })" + exp: > + jvm_gc_collection_seconds_count.tagMatch('gc', 'PS Scavenge|Copy|ParNew|G1 Young Generation|PS MarkSweep|MarkSweepCompact|ConcurrentMarkSweep|G1 Old Generation') + .sum(['service', 'host_name', 'gc']).increase('PT1M') + .tag({tags -> if (tags['gc'] == 'PS Scavenge' || tags['gc'] == 'Copy' || tags['gc'] == 'ParNew' || tags['gc'] == 'G1 Young Generation') {tags.gc = 'young_gc_count'} }) + .tag({tags -> if (tags['gc'] == 'PS MarkSweep' || tags['gc'] == 'MarkSweepCompact' || tags['gc'] == 'ConcurrentMarkSweep' || tags['gc'] == 'G1 Old Generation') {tags.gc = 'old_gc_count'} }) - name: instance_jvm_gc_time - exp: "(jvm_gc_collection_seconds_sum * 1000).tagMatch('gc', 'PS Scavenge|Copy|ParNew|G1 Young Generation|PS MarkSweep|MarkSweepCompact|ConcurrentMarkSweep|G1 Old Generation') - .sum(['service', 'host_name', 'gc']).increase('PT1M') - .tag({tags -> if (tags['gc'] == 'PS Scavenge' || tags['gc'] == 'Copy' || tags['gc'] == 'ParNew' || tags['gc'] == 'G1 Young Generation') {tags.gc = 'young_gc_time'} }) - .tag({tags -> if (tags['gc'] == 'PS MarkSweep' || tags['gc'] == 'MarkSweepCompact' || tags['gc'] == 'ConcurrentMarkSweep' || tags['gc'] == 'G1 Old Generation') {tags.gc = 'old_gc_time'} })" + exp: > + (jvm_gc_collection_seconds_sum * 1000).tagMatch('gc', 'PS Scavenge|Copy|ParNew|G1 Young Generation|PS MarkSweep|MarkSweepCompact|ConcurrentMarkSweep|G1 Old Generation') + .sum(['service', 'host_name', 'gc']).increase('PT1M') + .tag({tags -> if (tags['gc'] == 'PS Scavenge' || tags['gc'] == 'Copy' || tags['gc'] == 'ParNew' || tags['gc'] == 'G1 Young Generation') {tags.gc = 'young_gc_time'} }) + .tag({tags -> if (tags['gc'] == 'PS MarkSweep' || tags['gc'] == 'MarkSweepCompact' || tags['gc'] == 'ConcurrentMarkSweep' || tags['gc'] == 'G1 Old Generation') {tags.gc = 'old_gc_time'} }) - name: instance_trace_count exp: trace_in_latency_count.sum(['service', 'host_name']).increase('PT1M') - name: instance_trace_latency_percentile @@ -59,8 +61,9 @@ metricsRules: - name: instance_mesh_analysis_error_count exp: mesh_analysis_error_count.sum(['service', 'host_name']).increase('PT1M') - name: instance_metrics_aggregation - exp: "metrics_aggregation.tagEqual('dimensionality', 'minute').sum(['service', 'host_name', 'level']).increase('PT1M') - .tag({tags -> if (tags['level'] == '1') {tags.level = 'L1 aggregation'} }).tag({tags -> if (tags['level'] == '2') {tags.level = 'L2 aggregation'} })" + exp: > + metrics_aggregation.tagEqual('dimensionality', 'minute').sum(['service', 'host_name', 'level']).increase('PT1M') + .tag({tags -> if (tags['level'] == '1') {tags.level = 'L1 aggregation'} }).tag({tags -> if (tags['level'] == '2') {tags.level = 'L2 aggregation'} }) - name: instance_persistence_execute_percentile exp: persistence_timer_bulk_execute_latency.sum(['le', 'service', 'host_name']).increase('PT5M').histogram().histogram_percentile([50,70,90,99]) - name: instance_persistence_prepare_percentile @@ -99,3 +102,15 @@ metricsRules: exp: k8s_als_drop_count.sum(['service', 'host_name']).increase('PT1M') - name: instance_k8s_als_latency_percentile exp: k8s_als_in_latency.sum(['le', 'service', 'host_name']).increase('PT1M').histogram().histogram_percentile([50,70,90,99]) + - name: otel_metrics_received + exp: otel_metrics_latency_count.sum(['service', 'host_name']).increase('PT1M') + - name: otel_logs_received + exp: otel_logs_latency_count.sum(['service', 'host_name']).increase('PT1M') + - name: otel_spans_received + exp: otel_spans_latency_count.sum(['service', 'host_name']).increase('PT1M') + - name: otel_metrics_latency_percentile + exp: otel_metrics_latency.sum(['le', 'service', 'host_name']).increase('PT1M').histogram().histogram_percentile([50,70,90,99]) + - name: otel_logs_latency_percentile + exp: otel_logs_latency.sum(['le', 'service', 'host_name']).increase('PT1M').histogram().histogram_percentile([50,70,90,99]) + - name: otel_spans_latency_percentile + exp: otel_spans_latency.sum(['le', 'service', 'host_name']).increase('PT1M').histogram().histogram_percentile([50,70,90,99]) diff --git a/oap-server/server-starter/src/main/resources/ui-initialized-templates/so11y_oap/so11y-instance.json b/oap-server/server-starter/src/main/resources/ui-initialized-templates/so11y_oap/so11y-instance.json index 090d3de9e2..0752732cb5 100644 --- a/oap-server/server-starter/src/main/resources/ui-initialized-templates/so11y_oap/so11y-instance.json +++ b/oap-server/server-starter/src/main/resources/ui-initialized-templates/so11y_oap/so11y-instance.json @@ -467,6 +467,148 @@ "expressions": [ "relabels(meter_oap_instance_k8s_als_latency_percentile{p='50,75,90,95,99'},p='50,75,90,95,99',percentile='50,75,90,95,99')" ] + }, + { + "x": 12, + "y": 13, + "w": 6, + "h": 13, + "i": "19", + "type": "Widget", + "expressions": [ + "meter_oap_otel_metrics_received" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Metrics (Requests / Second)" + }, + "metricConfig": [ + { + "label": "Received Requests" + } + ] + }, + { + "x": 12, + "y": 26, + "w": 6, + "h": 13, + "i": "20", + "type": "Widget", + "expressions": [ + "meter_oap_otel_logs_received" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Logs (Requests / Second)" + } + }, + { + "x": 0, + "y": 26, + "w": 6, + "h": 13, + "i": "21", + "type": "Widget", + "expressions": [ + "meter_oap_otel_spans_received" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Spans (Requests / Second)" + }, + "metricConfig": [ + { + "unit": "Requests / Second" + } + ] + }, + { + "x": 18, + "y": 13, + "w": 6, + "h": 13, + "i": "22", + "type": "Widget", + "expressions": [ + "relabels(meter_oap_otel_metrics_latency_percentile{p='50,75,90,95,99'},p='50,75,90,95,99',percentile='50,75,90,95,99')" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Metrics Latency (ms / min)" + } + }, + { + "x": 6, + "y": 26, + "w": 6, + "h": 13, + "i": "23", + "type": "Widget", + "expressions": [ + "relabels(meter_oap_otel_spans_latency_percentile{p='50,75,90,95,99'},p='50,75,90,95,99',percentile='50,75,90,95,99')" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Spans Latency (ms / min)" + } + }, + { + "x": 18, + "y": 26, + "w": 6, + "h": 13, + "i": "24", + "type": "Widget", + "expressions": [ + "relabels(meter_oap_otel_logs_latency_percentile{p='50,75,90,95,99'},p='50,75,90,95,99',percentile='50,75,90,95,99')" + ], + "graph": { + "type": "Line", + "step": false, + "smooth": false, + "showSymbol": true, + "showXAxis": true, + "showYAxis": true + }, + "widget": { + "title": "OpenTelemetry Logs Latency (ms / min)" + } } ] }