Skip to content
Open
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
2 changes: 2 additions & 0 deletions .changelog/17.added
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
Enforce a default metric aggregation cardinality limit of 2000, folding
attribute sets beyond the limit into a single `otel.metric.overflow` series
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,16 @@

_logger = getLogger(__name__)

# Default aggregation cardinality limit as mandated by the specification:
# https://github.com/open-telemetry/opentelemetry-specification/blob/main/specification/metrics/sdk.md#cardinality-limits
_DEFAULT_CARDINALITY_LIMIT = 2000

# Synthetic attribute set used to aggregate measurements whose own attribute set
# could not be independently aggregated because the cardinality limit was
# reached.
_OVERFLOW_ATTRIBUTES = {"otel.metric.overflow": True}
_OVERFLOW_ATTRIBUTES_KEY = frozenset(_OVERFLOW_ATTRIBUTES.items())


class _ViewInstrumentMatch:
def __init__(
Expand All @@ -35,6 +45,7 @@ def __init__(
self._attributes_aggregation: dict[frozenset, _Aggregation] = {}
self._lock = Lock()
self._instrument_class_aggregation = instrument_class_aggregation
self._cardinality_limit = _DEFAULT_CARDINALITY_LIMIT
self._name = self._view._name or self._instrument.name
self._description = (
self._view._description or self._instrument.description
Expand Down Expand Up @@ -103,6 +114,30 @@ def consume_measurement(
if aggr_key not in self._attributes_aggregation:
with self._lock:
if aggr_key not in self._attributes_aggregation:
# Enforce the cardinality limit. Once the number of tracked
# attribute sets would exceed the limit, any previously
# unseen attribute set is aggregated into a single overflow
# series instead of allocating a new one. One slot is
# reserved for the overflow series so that the total number
# of metric points never exceeds the limit, matching the
# reference implementations (Go, Java).
if (
len(self._attributes_aggregation)
>= self._cardinality_limit - 1
):
aggr_key = _OVERFLOW_ATTRIBUTES_KEY
attributes = _OVERFLOW_ATTRIBUTES

if aggr_key not in self._attributes_aggregation:
if aggr_key is _OVERFLOW_ATTRIBUTES_KEY:
_logger.warning(
"Metric cardinality limit (%s) reached for "
"instrument %s. Further attribute sets are being "
"aggregated under the overflow attribute set "
"{'otel.metric.overflow': True}.",
self._cardinality_limit,
self._name,
)
if not isinstance(
self._view._aggregation, DefaultAggregation
):
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
# Copyright The OpenTelemetry Authors
# SPDX-License-Identifier: Apache-2.0

from unittest import TestCase

from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics._internal._view_instrument_match import (
_DEFAULT_CARDINALITY_LIMIT,
_OVERFLOW_ATTRIBUTES,
)
from opentelemetry.sdk.metrics.export import InMemoryMetricReader


class TestCardinalityLimit(TestCase):
@staticmethod
def _record_distinct_attribute_sets(count):
reader = InMemoryMetricReader()
meter_provider = MeterProvider(metric_readers=[reader])
meter = meter_provider.get_meter("testmeter")
counter = meter.create_counter("testcounter")

for index in range(count):
counter.add(1, {"index": index})

metrics_data = reader.get_metrics_data()
return (
metrics_data.resource_metrics[0]
.scope_metrics[0]
.metrics[0]
.data.data_points
)

def test_no_overflow_below_limit(self):
# One slot is reserved for the overflow series, so up to
# ``limit - 1`` distinct attribute sets are aggregated independently
# without producing an overflow series.
data_points = self._record_distinct_attribute_sets(
_DEFAULT_CARDINALITY_LIMIT - 1
)

self.assertEqual(len(data_points), _DEFAULT_CARDINALITY_LIMIT - 1)
self.assertNotIn(
_OVERFLOW_ATTRIBUTES,
[dict(data_point.attributes) for data_point in data_points],
)

def test_overflow_at_limit(self):
# At and beyond the limit, the total number of metric points stays
# capped at the limit and the excess is folded into a single overflow
# series.
excess = 500
data_points = self._record_distinct_attribute_sets(
_DEFAULT_CARDINALITY_LIMIT + excess
)

self.assertEqual(len(data_points), _DEFAULT_CARDINALITY_LIMIT)

overflow_points = [
data_point
for data_point in data_points
if dict(data_point.attributes) == _OVERFLOW_ATTRIBUTES
]
self.assertEqual(len(overflow_points), 1)

def test_no_measurement_dropped_during_overflow(self):
# Every measurement must be reflected in exactly one aggregator, so the
# summed value across all series (including overflow) equals the total
# number of measurements recorded.
total_measurements = _DEFAULT_CARDINALITY_LIMIT + 500
data_points = self._record_distinct_attribute_sets(total_measurements)

self.assertEqual(
sum(data_point.value for data_point in data_points),
total_measurements,
)
Loading