Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions datareservoirio/_logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
ivan-nikalaichanka-4ss marked this conversation as resolved.


@lru_cache(maxsize=1)
def get_exceptions_logger() -> logging.Logger:
exceptions_logger = logging.getLogger(__name__ + "_exception_logger")
Expand Down
81 changes: 31 additions & 50 deletions datareservoirio/client.py
Original file line number Diff line number Diff line change
@@ -1,19 +1,17 @@
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

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,
Expand All @@ -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
Expand Down Expand Up @@ -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

Expand Down
Loading