refactor: improve metric creation logic and silence ot warnings (#3207)
refactor(opentelemetry.py): enhance metric creation logic warnings and improved metric handling methods
This commit is contained in:
parent
8e733ce580
commit
7e7d5a1a1c
1 changed files with 43 additions and 42 deletions
|
|
@ -1,14 +1,15 @@
|
||||||
|
import threading
|
||||||
|
import warnings
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
from opentelemetry import metrics
|
|
||||||
from opentelemetry.exporter.prometheus import PrometheusMetricReader
|
|
||||||
from opentelemetry.metrics import Observation, CallbackOptions
|
|
||||||
from opentelemetry.metrics._internal.instrument import Counter, Histogram, UpDownCounter
|
|
||||||
from opentelemetry.sdk.metrics import MeterProvider
|
|
||||||
from opentelemetry.sdk.resources import Resource
|
|
||||||
from typing import Any, Dict, Mapping, Tuple, Union
|
from typing import Any, Dict, Mapping, Tuple, Union
|
||||||
from weakref import WeakValueDictionary
|
from weakref import WeakValueDictionary
|
||||||
|
|
||||||
import threading
|
from opentelemetry import metrics
|
||||||
|
from opentelemetry.exporter.prometheus import PrometheusMetricReader
|
||||||
|
from opentelemetry.metrics import CallbackOptions, Observation
|
||||||
|
from opentelemetry.metrics._internal.instrument import Counter, Histogram, UpDownCounter
|
||||||
|
from opentelemetry.sdk.metrics import MeterProvider
|
||||||
|
from opentelemetry.sdk.resources import Resource
|
||||||
|
|
||||||
# a default OpenTelelmetry meter name
|
# a default OpenTelelmetry meter name
|
||||||
langflow_meter_name = "langflow"
|
langflow_meter_name = "langflow"
|
||||||
|
|
@ -141,52 +142,52 @@ class OpenTelemetry(metaclass=ThreadSafeSingletonMetaUsingWeakref):
|
||||||
self._register_metric()
|
self._register_metric()
|
||||||
|
|
||||||
resource = Resource.create({"service.name": "langflow"})
|
resource = Resource.create({"service.name": "langflow"})
|
||||||
meter_provider = MeterProvider(resource=resource)
|
metric_readers = []
|
||||||
|
|
||||||
# configure prometheus exporter
|
# configure prometheus exporter
|
||||||
self.prometheus_enabled = prometheus_enabled
|
self.prometheus_enabled = prometheus_enabled
|
||||||
if prometheus_enabled:
|
if prometheus_enabled:
|
||||||
reader = PrometheusMetricReader()
|
metric_readers.append(PrometheusMetricReader())
|
||||||
meter_provider = MeterProvider(resource=resource, metric_readers=[reader])
|
|
||||||
|
|
||||||
|
meter_provider = MeterProvider(resource=resource, metric_readers=metric_readers)
|
||||||
metrics.set_meter_provider(meter_provider)
|
metrics.set_meter_provider(meter_provider)
|
||||||
self.meter = meter_provider.get_meter(langflow_meter_name)
|
self.meter = meter_provider.get_meter(langflow_meter_name)
|
||||||
|
|
||||||
for name, metric in self._metrics_registry.items():
|
for name, metric in self._metrics_registry.items():
|
||||||
# enforce the key in the mapping and metric's name are the same
|
|
||||||
# this error can get caught at unit test
|
|
||||||
if name != metric.name:
|
if name != metric.name:
|
||||||
raise ValueError(f"Key '{name}' does not match metric name '{metric.name}'")
|
raise ValueError(f"Key '{name}' does not match metric name '{metric.name}'")
|
||||||
if metric.type == MetricType.COUNTER:
|
with warnings.catch_warnings():
|
||||||
counter = self.meter.create_counter(
|
warnings.simplefilter("ignore")
|
||||||
name=metric.name,
|
|
||||||
unit=metric.unit,
|
self._metrics[metric.name] = self._create_metric(metric)
|
||||||
description=metric.description,
|
|
||||||
)
|
def _create_metric(self, metric):
|
||||||
self._metrics[metric.name] = counter
|
if metric.type == MetricType.COUNTER:
|
||||||
elif metric.type == MetricType.OBSERVABLE_GAUGE:
|
return self.meter.create_counter(
|
||||||
gauge = ObservableGaugeWrapper(
|
name=metric.name,
|
||||||
name=metric.name,
|
unit=metric.unit,
|
||||||
description=metric.description,
|
description=metric.description,
|
||||||
unit=metric.unit,
|
)
|
||||||
)
|
elif metric.type == MetricType.OBSERVABLE_GAUGE:
|
||||||
self._metrics[metric.name] = gauge
|
return ObservableGaugeWrapper(
|
||||||
elif metric.type == MetricType.UP_DOWN_COUNTER:
|
name=metric.name,
|
||||||
up_down_counter = self.meter.create_up_down_counter(
|
description=metric.description,
|
||||||
name=metric.name,
|
unit=metric.unit,
|
||||||
unit=metric.unit,
|
)
|
||||||
description=metric.description,
|
elif metric.type == MetricType.UP_DOWN_COUNTER:
|
||||||
)
|
return self.meter.create_up_down_counter(
|
||||||
self._metrics[metric.name] = up_down_counter
|
name=metric.name,
|
||||||
elif metric.type == MetricType.HISTOGRAM:
|
unit=metric.unit,
|
||||||
histogram = self.meter.create_histogram(
|
description=metric.description,
|
||||||
name=metric.name,
|
)
|
||||||
unit=metric.unit,
|
elif metric.type == MetricType.HISTOGRAM:
|
||||||
description=metric.description,
|
return self.meter.create_histogram(
|
||||||
)
|
name=metric.name,
|
||||||
self._metrics[metric.name] = histogram
|
unit=metric.unit,
|
||||||
else:
|
description=metric.description,
|
||||||
raise ValueError(f"Unknown metric type: {metric.type}")
|
)
|
||||||
|
else:
|
||||||
|
raise ValueError(f"Unknown metric type: {metric.type}")
|
||||||
|
|
||||||
def validate_labels(self, metric_name: str, labels: Mapping[str, str]):
|
def validate_labels(self, metric_name: str, labels: Mapping[str, str]):
|
||||||
reg = self._metrics_registry.get(metric_name)
|
reg = self._metrics_registry.get(metric_name)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue