Skip to content

Commit ab22674

Browse files
fix(opentelemetry-sdk): keep synchronous gauge values across cumulative collections (#5637)
* fix(opentelemetry-sdk): keep synchronous gauge values across cumulative collections The MetricReader spec requires that for synchronous instruments with cumulative aggregation temporality, Collect receives the data points exposed in previous collections regardless of whether new measurements have been recorded: https://opentelemetry.io/docs/specs/otel/metrics/sdk/#metricreader _LastValueAggregation dropped its value on every collection, so a synchronous gauge disappeared from the export as soon as one collection interval passed without a set() call. Fixes #4512 Fixes #3971 * chore(changelog): rename fragment to match PR number * chore(changelog): apply reviewer wording to 5637 fragment * test(opentelemetry-sdk): cover synchronous gauge persistence via InMemoryMetricReader Collecting twice through an InMemoryMetricReader without recording a new measurement in between must still export the gauge with its last value. --------- Co-authored-by: Emídio <9735060+emdneto@users.noreply.github.com>
1 parent 477ffd4 commit ab22674

5 files changed

Lines changed: 85 additions & 10 deletions

File tree

.changelog/5637.fixed

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
`opentelemetry-sdk`: retain values from synchronous instruments using last-value aggregation across cumulative collections

opentelemetry-sdk/src/opentelemetry/sdk/metrics/_internal/aggregation.py

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -369,8 +369,10 @@ def __init__(
369369
self,
370370
attributes: Attributes,
371371
reservoir_builder: ExemplarReservoirBuilder,
372+
instrument_is_synchronous: bool,
372373
):
373374
super().__init__(attributes, reservoir_builder)
375+
self._instrument_is_synchronous = instrument_is_synchronous
374376
self._value = None
375377

376378
def aggregate(self, measurement: Measurement, should_sample_exemplar: bool = True):
@@ -391,7 +393,11 @@ def collect(
391393
if self._value is None:
392394
return None
393395
value = self._value
394-
self._value = None
396+
if not (
397+
self._instrument_is_synchronous
398+
and collection_aggregation_temporality is AggregationTemporality.CUMULATIVE
399+
):
400+
self._value = None
395401

396402
exemplars = self._collect_exemplars()
397403

@@ -1171,12 +1177,14 @@ def _create_aggregation(
11711177
return _LastValueAggregation(
11721178
attributes,
11731179
reservoir_builder=reservoir_factory(_LastValueAggregation),
1180+
instrument_is_synchronous=False,
11741181
)
11751182

11761183
if isinstance(instrument, _Gauge):
11771184
return _LastValueAggregation(
11781185
attributes,
11791186
reservoir_builder=reservoir_factory(_LastValueAggregation),
1187+
instrument_is_synchronous=True,
11801188
)
11811189

11821190
# pylint: disable=broad-exception-raised
@@ -1329,6 +1337,7 @@ def _create_aggregation(
13291337
return _LastValueAggregation(
13301338
attributes,
13311339
reservoir_builder=reservoir_factory(_LastValueAggregation),
1340+
instrument_is_synchronous=isinstance(instrument, Synchronous),
13321341
)
13331342

13341343

opentelemetry-sdk/tests/metrics/test_aggregation.py

Lines changed: 48 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -210,7 +210,11 @@ def test_aggregate(self):
210210
temporality
211211
"""
212212

213-
last_value_aggregation = _LastValueAggregation(Mock(), _default_reservoir_factory(_LastValueAggregation))
213+
last_value_aggregation = _LastValueAggregation(
214+
Mock(),
215+
_default_reservoir_factory(_LastValueAggregation),
216+
instrument_is_synchronous=True,
217+
)
214218

215219
last_value_aggregation.aggregate(measurement(1))
216220
self.assertEqual(last_value_aggregation._value, 1)
@@ -226,7 +230,11 @@ def test_collect(self):
226230
`LastValueAggregation` collects number data points
227231
"""
228232

229-
last_value_aggregation = _LastValueAggregation(Mock(), _default_reservoir_factory(_LastValueAggregation))
233+
last_value_aggregation = _LastValueAggregation(
234+
Mock(),
235+
_default_reservoir_factory(_LastValueAggregation),
236+
instrument_is_synchronous=True,
237+
)
230238

231239
self.assertIsNone(last_value_aggregation.collect(AggregationTemporality.CUMULATIVE, 1))
232240

@@ -240,7 +248,7 @@ def test_collect(self):
240248

241249
self.assertIsNone(first_number_data_point.start_time_unix_nano)
242250

243-
last_value_aggregation.aggregate(measurement(1))
251+
last_value_aggregation.aggregate(measurement(2))
244252

245253
# CI fails the last assertion without this
246254
sleep(0.1)
@@ -249,7 +257,7 @@ def test_collect(self):
249257
# collection process starts.
250258
second_number_data_point = last_value_aggregation.collect(AggregationTemporality.CUMULATIVE, 2)
251259

252-
self.assertEqual(second_number_data_point.value, 1)
260+
self.assertEqual(second_number_data_point.value, 2)
253261

254262
self.assertIsNone(second_number_data_point.start_time_unix_nano)
255263

@@ -258,10 +266,44 @@ def test_collect(self):
258266
first_number_data_point.time_unix_nano,
259267
)
260268

261-
# 3 is used here directly to simulate the instant the second
269+
# 3 is used here directly to simulate the instant the third
262270
# collection process starts.
263271
third_number_data_point = last_value_aggregation.collect(AggregationTemporality.CUMULATIVE, 3)
264-
self.assertIsNone(third_number_data_point)
272+
self.assertIsInstance(third_number_data_point, NumberDataPoint)
273+
self.assertEqual(third_number_data_point.value, 2)
274+
self.assertIsNone(third_number_data_point.start_time_unix_nano)
275+
self.assertGreater(
276+
third_number_data_point.time_unix_nano,
277+
second_number_data_point.time_unix_nano,
278+
)
279+
280+
def test_collect_delta_resets_value(self):
281+
last_value_aggregation = _LastValueAggregation(
282+
Mock(),
283+
_default_reservoir_factory(_LastValueAggregation),
284+
instrument_is_synchronous=True,
285+
)
286+
287+
last_value_aggregation.aggregate(measurement(1))
288+
289+
number_data_point = last_value_aggregation.collect(AggregationTemporality.DELTA, 1)
290+
self.assertEqual(number_data_point.value, 1)
291+
292+
self.assertIsNone(last_value_aggregation.collect(AggregationTemporality.DELTA, 2))
293+
294+
def test_collect_asynchronous_resets_value(self):
295+
last_value_aggregation = _LastValueAggregation(
296+
Mock(),
297+
_default_reservoir_factory(_LastValueAggregation),
298+
instrument_is_synchronous=False,
299+
)
300+
301+
last_value_aggregation.aggregate(measurement(1))
302+
303+
number_data_point = last_value_aggregation.collect(AggregationTemporality.CUMULATIVE, 1)
304+
self.assertEqual(number_data_point.value, 1)
305+
306+
self.assertIsNone(last_value_aggregation.collect(AggregationTemporality.CUMULATIVE, 2))
265307

266308

267309
class TestExplicitBucketHistogramAggregation(TestCase):

opentelemetry-sdk/tests/metrics/test_in_memory_metric_reader.py

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,3 +142,26 @@ def test_cumulative_multiple_collect(self):
142142
number_data_point_1.time_unix_nano,
143143
number_data_point_0.time_unix_nano,
144144
)
145+
146+
def test_synchronous_gauge_cumulative_multiple_collect(self):
147+
reader = InMemoryMetricReader()
148+
meter = MeterProvider(metric_readers=[reader]).get_meter("test_meter")
149+
gauge = meter.create_gauge("gauge1")
150+
gauge.set(5, attributes={"key": "value"})
151+
152+
first_metrics_data = reader.get_metrics_data()
153+
second_metrics_data = reader.get_metrics_data()
154+
155+
self.assertIsNotNone(second_metrics_data)
156+
157+
first_gauge_data = first_metrics_data.resource_metrics[0].scope_metrics[0].metrics[0].data
158+
second_gauge_data = second_metrics_data.resource_metrics[0].scope_metrics[0].metrics[0].data
159+
160+
self.assertEqual(
161+
[(point.attributes, point.value) for point in first_gauge_data.data_points],
162+
[({"key": "value"}, 5)],
163+
)
164+
self.assertEqual(
165+
[(point.attributes, point.value) for point in second_gauge_data.data_points],
166+
[({"key": "value"}, 5)],
167+
)

opentelemetry-sdk/tests/metrics/test_metric_reader_storage.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -91,9 +91,9 @@ def test_creates_view_instrument_matches(self, MockViewInstrumentMatch: Mock):
9191

9292
@patch("opentelemetry.sdk.metrics._internal.metric_reader_storage._ViewInstrumentMatch")
9393
def test_forwards_calls_to_view_instrument_match(self, MockViewInstrumentMatch: Mock):
94-
view_instrument_match1 = Mock(_aggregation=_LastValueAggregation({}, Mock()))
95-
view_instrument_match2 = Mock(_aggregation=_LastValueAggregation({}, Mock()))
96-
view_instrument_match3 = Mock(_aggregation=_LastValueAggregation({}, Mock()))
94+
view_instrument_match1 = Mock(_aggregation=_LastValueAggregation({}, Mock(), instrument_is_synchronous=True))
95+
view_instrument_match2 = Mock(_aggregation=_LastValueAggregation({}, Mock(), instrument_is_synchronous=True))
96+
view_instrument_match3 = Mock(_aggregation=_LastValueAggregation({}, Mock(), instrument_is_synchronous=True))
9797
MockViewInstrumentMatch.side_effect = [
9898
view_instrument_match1,
9999
view_instrument_match2,

0 commit comments

Comments
 (0)