diff --git a/datareservoirio/_logging.py b/datareservoirio/_logging.py index cc90fe2a..b9deb117 100644 --- a/datareservoirio/_logging.py +++ b/datareservoirio/_logging.py @@ -27,6 +27,24 @@ def _ensure_azure_monitor_configured(connection_string, logger_name): _configured_loggers[cache_key] = True +@lru_cache(maxsize=1) +def get_metric_logger() -> logging.Logger: + logger = logging.getLogger(__name__ + "_metric_appinsight") + if os.getenv(ENV_VAR_ENABLE_APP_INSIGHTS) is not None: + enable_app_insights = os.environ[ENV_VAR_ENABLE_APP_INSIGHTS].lower() + if enable_app_insights == "true" or enable_app_insights == "1": + # Prevent messages from being passed directly to ancestor logger handlers, + # since ancestor logger levels and filters are not considered during propagation. + # https://docs.python.org/3.12/library/logging.html#logging.Logger.propagate + logger.propagate = False + logger.setLevel(logging.DEBUG) + _ensure_azure_monitor_configured( + connection_string=environment._application_insight_connectionstring, + logger_name=__name__ + "_metric_appinsight", + ) + return logger + + @lru_cache(maxsize=1) def get_exceptions_logger() -> logging.Logger: exceptions_logger = logging.getLogger(__name__ + "_exception_logger") diff --git a/datareservoirio/client.py b/datareservoirio/client.py index 22c7ca4c..742be274 100644 --- a/datareservoirio/client.py +++ b/datareservoirio/client.py @@ -1,11 +1,9 @@ import logging -import os import time import warnings from collections import defaultdict from concurrent.futures import ThreadPoolExecutor -from datetime import datetime -from functools import lru_cache, wraps +from functools import wraps from operator import itemgetter from urllib.parse import urlencode from uuid import uuid4 @@ -13,7 +11,7 @@ import numpy as np import pandas as pd import requests -from azure.monitor.opentelemetry import configure_azure_monitor +from opentelemetry import trace from tenacity import ( retry, retry_if_exception_type, @@ -23,34 +21,13 @@ ) from tqdm.auto import tqdm -from datareservoirio._constants import ENV_VAR_ENABLE_APP_INSIGHTS - -from ._logging import _ensure_azure_monitor_configured, log_decorator +from ._logging import get_metric_logger, log_decorator from ._utils import function_translation, period_translation from .globalsettings import environment from .storage import Storage log = logging.getLogger(__name__) - -@lru_cache(maxsize=1) -def metric() -> logging.Logger: - logger = logging.getLogger(__name__ + "_metric_appinsight") - if os.getenv(ENV_VAR_ENABLE_APP_INSIGHTS) is not None: - enable_app_insights = os.environ[ENV_VAR_ENABLE_APP_INSIGHTS].lower() - if enable_app_insights == "true" or enable_app_insights == "1": - # Prevent messages from being passed directly to ancestor logger handlers, - # since ancestor logger levels and filters are not considered during propagation. - # https://docs.python.org/3.12/library/logging.html#logging.Logger.propagate - logger.propagate = False - logger.setLevel(logging.DEBUG) - _ensure_azure_monitor_configured( - connection_string=environment._application_insight_connectionstring, - logger_name=__name__ + "_metric_appinsight", - ) - return logger - - # Default values to push as start/end dates. (Limited by numpy.datetime64) _END_DEFAULT = 9214646400000000000 # 2262-01-01 _START_DEFAULT = 0 # 1970-01-01 @@ -328,30 +305,34 @@ def _timer(metric_name="Timer"): def decorator(func): @wraps(func) def wrapper(self, series_id, start=None, end=None, **kwargs): - start_time = time.perf_counter() - result = func(self, series_id, start=start, end=end, **kwargs) - end_time = time.perf_counter() - elapsed_time = end_time - start_time - start_date_as_str = None - end_date_as_str = None - if start: - start_date_as_str = pd.to_datetime( - start, dayfirst=True, unit="ns", utc=True - ).isoformat() - if end: - end_date_as_str = pd.to_datetime( - end, dayfirst=True, unit="ns", utc=True - ).isoformat() - number_of_samples = len(result) - properties = { - "series_id": series_id, - "start": start_date_as_str, - "end": end_date_as_str, - "elapsed": elapsed_time, - "number-of-samples": number_of_samples, - } - metric().info(metric_name, extra=properties) - return result + metric_logger = get_metric_logger() + tracer = trace.get_tracer(metric_logger.name) + with tracer.start_as_current_span(metric_name): + + start_time = time.perf_counter() + result = func(self, series_id, start=start, end=end, **kwargs) + end_time = time.perf_counter() + elapsed_time = end_time - start_time + start_date_as_str = None + end_date_as_str = None + if start: + start_date_as_str = pd.to_datetime( + start, dayfirst=True, unit="ns", utc=True + ).isoformat() + if end: + end_date_as_str = pd.to_datetime( + end, dayfirst=True, unit="ns", utc=True + ).isoformat() + number_of_samples = len(result) + properties = { + "series_id": series_id, + "start": start_date_as_str, + "end": end_date_as_str, + "elapsed": elapsed_time, + "number-of-samples": number_of_samples, + } + metric_logger.info(metric_name, extra=properties) + return result return wrapper