refactor: update OpenTelemetry initialization to avoid multiple calls (#4306)

* Ensure single initialization of OpenTelemetry and optimize meter provider setup

* Simplify condition by removing redundant class attribute
This commit is contained in:
Gabriel Luiz Freitas Almeida 2024-10-29 12:04:22 -03:00 • committed by GitHub
commit b39ee0ae2d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 23 additions and 12 deletions

Binary file not shown.

View file

@ -1,5 +1,4 @@
import threading import threading
import warnings
from collections.abc import Mapping from collections.abc import Mapping
from enum import Enum from enum import Enum
from typing import Any from typing import Any
@ -110,6 +109,8 @@ class OpenTelemetry(metaclass=ThreadSafeSingletonMetaUsingWeakref):
_metrics_registry: dict[str, Metric] = {} _metrics_registry: dict[str, Metric] = {}
_metrics: dict[str, Counter | ObservableGaugeWrapper | Histogram | UpDownCounter] = {} _metrics: dict[str, Counter | ObservableGaugeWrapper | Histogram | UpDownCounter] = {}
_meter_provider: MeterProvider | None = None _meter_provider: MeterProvider | None = None
_initialized: bool = False # Add initialization flag
prometheus_enabled: bool = True
def _add_metric( def _add_metric(
self, name: str, description: str, unit: str, metric_type: MetricType, labels: dict[str, bool] self, name: str, description: str, unit: str, metric_type: MetricType, labels: dict[str, bool]
@ -141,20 +142,29 @@ class OpenTelemetry(metaclass=ThreadSafeSingletonMetaUsingWeakref):
) )
def __init__(self, *, prometheus_enabled: bool = True): def __init__(self, *, prometheus_enabled: bool = True):
# Only initialize once
self.prometheus_enabled = prometheus_enabled
if OpenTelemetry._initialized:
return
if not self._metrics_registry: if not self._metrics_registry:
self._register_metric() self._register_metric()
if self._meter_provider is None: if self._meter_provider is None:
resource = Resource.create({"service.name": "langflow"}) # Get existing meter provider if any
metric_readers = [] existing_provider = metrics.get_meter_provider()
# configure prometheus exporter # Check if FastAPI instrumentation is already set up
self.prometheus_enabled = prometheus_enabled if hasattr(existing_provider, "get_meter") and existing_provider.get_meter("http.server"):
if prometheus_enabled: self._meter_provider = existing_provider
metric_readers.append(PrometheusMetricReader()) else:
resource = Resource.create({"service.name": "langflow"})
metric_readers = []
if self.prometheus_enabled:
metric_readers.append(PrometheusMetricReader())
self._meter_provider = MeterProvider(resource=resource, metric_readers=metric_readers) self._meter_provider = MeterProvider(resource=resource, metric_readers=metric_readers)
metrics.set_meter_provider(self._meter_provider) metrics.set_meter_provider(self._meter_provider)
self.meter = self._meter_provider.get_meter(langflow_meter_name) self.meter = self._meter_provider.get_meter(langflow_meter_name)
@ -163,11 +173,12 @@ class OpenTelemetry(metaclass=ThreadSafeSingletonMetaUsingWeakref):
msg = f"Key '{name}' does not match metric name '{metric.name}'" msg = f"Key '{name}' does not match metric name '{metric.name}'"
raise ValueError(msg) raise ValueError(msg)
if name not in self._metrics: if name not in self._metrics:
with warnings.catch_warnings(): self._metrics[metric.name] = self._create_metric(metric)
warnings.simplefilter("ignore")
self._metrics[metric.name] = self._create_metric(metric) OpenTelemetry._initialized = True
def _create_metric(self, metric): def _create_metric(self, metric):
# Remove _created_instruments check
if metric.name in self._metrics: if metric.name in self._metrics:
return self._metrics[metric.name] return self._metrics[metric.name]