Skip to content

Commit 5c79afb

Browse files
Fix OTel timer metrics using Gauge instead of Histogram (apache#64207) (apache#66865)
* Fix OTel timer metrics using Gauge instead of Histogram * Use ExponentialBucketHistogramAggregation for timing metrics * Use public API import path for ExponentialBucketHistogramAggregation and fix histogram map isolation (cherry picked from commit b2dadd2) Co-authored-by: namratachaudhary <namratachaudhary@users.noreply.github.com>
1 parent ea3657f commit 5c79afb

3 files changed

Lines changed: 86 additions & 20 deletions

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
OTel timer and timing metrics now use Histogram instead of Gauge, preserving count, sum, and bucket distribution across recordings.

shared/observability/src/airflow_shared/observability/metrics/otel_logger.py

Lines changed: 52 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
ConsoleMetricExporter,
3131
PeriodicExportingMetricReader,
3232
)
33+
from opentelemetry.sdk.metrics.view import ExponentialBucketHistogramAggregation, View
3334
from opentelemetry.sdk.resources import SERVICE_NAME, Resource
3435

3536
from ..common import get_otel_data_exporter
@@ -146,7 +147,8 @@ class _OtelTimer(Timer):
146147
"""
147148
An implementation of Stats.Timer() which records the result in the OTel Metrics Map.
148149
149-
OpenTelemetry does not have a native timer, we will store the values as a Gauge.
150+
OpenTelemetry does not have a native timer; values are stored as a Histogram so that
151+
all observations (count, sum, bucket distribution) are preserved across multiple recordings.
150152
151153
:param name: The name of the timer.
152154
:param tags: Tags to append to the timer.
@@ -160,9 +162,9 @@ def __init__(self, otel_logger: SafeOtelLogger, name: str | None, tags: Attribut
160162

161163
def stop(self, send: bool = True) -> None:
162164
super().stop(send)
163-
if self.name and send and self.duration:
164-
self.otel_logger.metrics_map.set_gauge_value(
165-
full_name(prefix=self.otel_logger.prefix, name=self.name), self.duration, False, self.tags
165+
if self.name and send and self.duration is not None:
166+
self.otel_logger.metrics_map.record_histogram_value(
167+
full_name(prefix=self.otel_logger.prefix, name=self.name), self.duration, self.tags
166168
)
167169

168170

@@ -278,11 +280,11 @@ def timing(
278280
*,
279281
tags: Attributes = None,
280282
) -> None:
281-
"""OTel does not have a native timer, stored as a Gauge whose value is elapsed ms."""
283+
"""Record a timing observation as a Histogram to preserve distribution information."""
282284
if self.metrics_validator.test(stat) and name_is_otel_safe(self.prefix, stat):
283285
if isinstance(dt, datetime.timedelta):
284286
dt = dt.total_seconds() * 1000.0
285-
self.metrics_map.set_gauge_value(full_name(prefix=self.prefix, name=stat), float(dt), False, tags)
287+
self.metrics_map.record_histogram_value(full_name(prefix=self.prefix, name=stat), float(dt), tags)
286288

287289
def timer(
288290
self,
@@ -314,15 +316,29 @@ def set_value(self, new_value: int | float, delta: bool):
314316
self.gauge.set(new_value, attributes=self.attributes)
315317

316318

319+
class InternalHistogram:
320+
"""Stores a histogram instrument for timer/timing metrics."""
321+
322+
def __init__(self, meter, name: str):
323+
otel_safe_name = _get_otel_safe_name(name)
324+
self.histogram = meter.create_histogram(name=otel_safe_name, unit="ms")
325+
log.debug("Created %s as type: %s", otel_safe_name, _type_as_str(self.histogram))
326+
327+
def record(self, value: float, tags: Attributes) -> None:
328+
self.histogram.record(value, attributes=tags)
329+
330+
317331
class MetricsMap:
318332
"""Stores Otel Instruments."""
319333

320334
def __init__(self, meter):
321335
self.meter = meter
322336
self.map = {}
337+
self.histograms: dict[str, InternalHistogram] = {}
323338

324339
def clear(self) -> None:
325340
self.map.clear()
341+
self.histograms.clear()
326342

327343
def _create_counter(self, name):
328344
"""Create a new counter or up_down_counter for the provided name."""
@@ -376,6 +392,21 @@ def set_gauge_value(self, name: str, value: int | float, delta: bool, tags: Attr
376392

377393
self.map[key].set_value(value, delta)
378394

395+
def record_histogram_value(self, name: str, value: float, tags: Attributes) -> None:
396+
"""
397+
Record a timing observation in a Histogram instrument.
398+
399+
Unlike a Gauge, a Histogram accumulates all observations so that count, sum,
400+
and bucket distribution are preserved across multiple recordings.
401+
402+
:param name: The name of the histogram to record.
403+
:param value: The timing observation in milliseconds.
404+
:param tags: Attributes to attach to the observation.
405+
"""
406+
if name not in self.histograms:
407+
self.histograms[name] = InternalHistogram(meter=self.meter, name=name)
408+
self.histograms[name].record(value, tags)
409+
379410

380411
def flush_otel_metrics():
381412
provider = metrics.get_meter_provider()
@@ -400,6 +431,15 @@ def get_otel_logger(
400431
stat_name_handler: Callable[[str], str] | None = None,
401432
statsd_influxdb_enabled: bool = False,
402433
) -> SafeOtelLogger:
434+
"""
435+
Build and return a :class:`SafeOtelLogger` backed by a configured :class:`MeterProvider`.
436+
437+
Histogram instruments (used for ``timing()`` / ``timer()`` metrics) are aggregated with
438+
:class:`~opentelemetry.sdk.metrics.view.ExponentialBucketHistogramAggregation`
439+
so that bucket boundaries adapt automatically to the observed data range. This avoids
440+
the need to hand-tune explicit bucket boundaries for metrics that span very different
441+
scales (milliseconds to hours).
442+
"""
403443
otel_env_config = load_metrics_env_config()
404444

405445
effective_service_name: str = otel_env_config.service_name or service_name or "airflow"
@@ -453,6 +493,12 @@ def get_otel_logger(
453493
MeterProvider(
454494
resource=resource,
455495
metric_readers=readers,
496+
views=[
497+
View(
498+
instrument_type=metrics.Histogram,
499+
aggregation=ExponentialBucketHistogramAggregation(),
500+
)
501+
],
456502
shutdown_on_exit=False,
457503
),
458504
)

shared/observability/tests/observability/metrics/test_otel_logger.py

Lines changed: 33 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525

2626
import pytest
2727
from opentelemetry.metrics import MeterProvider
28+
from opentelemetry.sdk.metrics.view import ExponentialBucketHistogramAggregation, View
2829

2930
from airflow_shared.observability.common import get_otel_data_exporter
3031
from airflow_shared.observability.exceptions import InvalidStatsNameException
@@ -244,25 +245,28 @@ def test_timing_new_metric(self, name):
244245

245246
self.stats.timing(name, dt=datetime.timedelta(seconds=123))
246247

247-
self.meter.get_meter().create_gauge.assert_called_once_with(name=full_name(name))
248-
expected_value = 123000.0
249-
assert self.map[full_name(name)].value == expected_value
248+
self.meter.get_meter().create_histogram.assert_called_once_with(name=full_name(name), unit="ms")
249+
self.meter.get_meter().create_histogram.return_value.record.assert_called_once_with(
250+
123000.0, attributes=None
251+
)
250252

251253
def test_timing_new_metric_with_tags(self, name):
252254
tags = {"hello": "world"}
253-
key = _generate_key_name(full_name(name), tags)
254255

255256
self.stats.timing(name, dt=1, tags=tags)
256257

257-
self.meter.get_meter().create_gauge.assert_called_once_with(name=full_name(name))
258-
self.map[key].attributes == tags
258+
self.meter.get_meter().create_histogram.assert_called_once_with(name=full_name(name), unit="ms")
259+
self.meter.get_meter().create_histogram.return_value.record.assert_called_once_with(
260+
1.0, attributes=tags
261+
)
259262

260263
def test_timing_existing_metric(self, name):
261264
self.stats.timing(name, dt=1)
262265
self.stats.timing(name, dt=2)
263266

264-
self.meter.get_meter().create_gauge.assert_called_once_with(name=full_name(name))
265-
assert self.map[full_name(name)].value == 2
267+
# histogram created only once, but both observations are recorded
268+
self.meter.get_meter().create_histogram.assert_called_once_with(name=full_name(name), unit="ms")
269+
assert self.meter.get_meter().create_histogram.return_value.record.call_count == 2
266270

267271
# For the four test_timer_foo tests below:
268272
# time.perf_count() is called once to get the starting timestamp and again
@@ -277,7 +281,7 @@ def test_timer_with_name_returns_float_and_stores_value(self, mock_time, name):
277281
expected_duration = 3140.0
278282
assert timer.duration == expected_duration
279283
assert mock_time.call_count == 2
280-
self.meter.get_meter().create_gauge.assert_called_once_with(name=full_name(name))
284+
self.meter.get_meter().create_histogram.assert_called_once_with(name=full_name(name), unit="ms")
281285

282286
@mock.patch.object(time, "perf_counter", side_effect=[0.0, 3.14])
283287
def test_timer_no_name_returns_float_but_does_not_store_value(self, mock_time, name):
@@ -288,7 +292,7 @@ def test_timer_no_name_returns_float_but_does_not_store_value(self, mock_time, n
288292
expected_duration = 3140.0
289293
assert timer.duration == expected_duration
290294
assert mock_time.call_count == 2
291-
self.meter.get_meter().create_gauge.assert_not_called()
295+
self.meter.get_meter().create_histogram.assert_not_called()
292296

293297
@mock.patch.object(time, "perf_counter", side_effect=[0.0, 3.14])
294298
def test_timer_start_and_stop_manually_send_false(self, mock_time, name):
@@ -301,7 +305,7 @@ def test_timer_start_and_stop_manually_send_false(self, mock_time, name):
301305
expected_value = 3140.0
302306
assert timer.duration == expected_value
303307
assert mock_time.call_count == 2
304-
self.meter.get_meter().create_gauge.assert_not_called()
308+
self.meter.get_meter().create_histogram.assert_not_called()
305309

306310
@mock.patch.object(time, "perf_counter", side_effect=[0.0, 3.14])
307311
def test_timer_start_and_stop_manually_send_true(self, mock_time, name):
@@ -314,7 +318,7 @@ def test_timer_start_and_stop_manually_send_true(self, mock_time, name):
314318
expected_value = 3140.0
315319
assert timer.duration == expected_value
316320
assert mock_time.call_count == 2
317-
self.meter.get_meter().create_gauge.assert_called_once_with(name=full_name(name))
321+
self.meter.get_meter().create_histogram.assert_called_once_with(name=full_name(name), unit="ms")
318322

319323
@pytest.mark.parametrize(
320324
(
@@ -415,15 +419,30 @@ def test_config_priorities(
415419
== f"opentelemetry.exporter.otlp.proto.{expected_exporter_module}.metric_exporter"
416420
)
417421

422+
@mock.patch("airflow_shared.observability.metrics.otel_logger.metrics")
423+
@mock.patch("airflow_shared.observability.metrics.otel_logger.MeterProvider")
424+
def test_get_otel_logger_uses_exponential_histogram_view(self, mock_provider, mock_metrics):
425+
get_otel_logger(host="localhost", port=4318)
426+
427+
call_kwargs = mock_provider.call_args.kwargs
428+
views = call_kwargs["views"]
429+
assert len(views) == 1
430+
view = views[0]
431+
assert isinstance(view, View)
432+
assert isinstance(view._aggregation, ExponentialBucketHistogramAggregation)
433+
418434
def test_atexit_flush_on_process_exit(self):
419435
"""
420436
Run a process that initializes a logger, creates a stat and then exits.
421437
422438
The logger initialization registers an atexit hook.
423439
Test that the hook runs and flushes the created stat at shutdown.
424440
"""
425-
test_module_name = "tests.observability.metrics.test_otel_logger"
426-
function_call_str = f"import {test_module_name} as m; m.mock_service_run()"
441+
function_call_str = (
442+
"from airflow_shared.observability.metrics.otel_logger import get_otel_logger; "
443+
"logger = get_otel_logger(debug=True); "
444+
"logger.incr('my_test_stat')"
445+
)
427446

428447
proc = subprocess.run(
429448
[sys.executable, "-c", function_call_str],

0 commit comments

Comments
 (0)