Skip to content

Commit 91e49a5

Browse files
authored
Count otel.sdk.processor.{span,log}.processed at exporter-submit time (open-telemetry#5472)
* Count processor.processed at exporter-submit time, not export completion * Add changelog fragment * Trim comments and fix changelog wording * Clarify changelog covers simple and batch processors * Use single-line changelog fragment without redundant PR link Assisted-by: Claude Opus 4.8 * Add test that shutdown-dropped records are not counted as processed Assisted-by: Claude Opus 4.8 * Move SimpleSpanProcessor finish_items outside the export try/except Assisted-by: Claude Opus 4.8 * Simplify shutdown metrics test to call on_emit directly Assisted-by: Claude Opus 4.8
1 parent cd298d5 commit 91e49a5

7 files changed

Lines changed: 91 additions & 114 deletions

File tree

‎.changelog/5472.fixed‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
`opentelemetry-sdk`: for both the simple and batch span/log processors, count `otel.sdk.processor.{span,log}.processed` when the processor submits records to the exporter instead of after export completes, and stop stamping exporter failures onto this metric as `error.type`

‎opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py‎

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -230,7 +230,6 @@ def on_emit(self, log_record: ReadWriteLogRecord):
230230
set_value(_ON_EMIT_RECURSION_COUNT_KEY, cnt + 1), # pyright: ignore[reportOperatorIssue]
231231
)
232232
)
233-
error: Exception | None = None
234233
try:
235234
if self._shutdown:
236235
_logger.warning("Processor is already shutdown, ignoring call")
@@ -248,12 +247,12 @@ def on_emit(self, log_record: ReadWriteLogRecord):
248247
instrumentation_scope=log_record.instrumentation_scope,
249248
limits=log_record.limits,
250249
)
250+
# Record on submission to the exporter.
251+
self._metrics.finish_items(1)
251252
self._exporter.export((readable_log_record,))
252-
except Exception as err: # pylint: disable=broad-exception-caught
253-
error = err
253+
except Exception: # pylint: disable=broad-exception-caught
254254
_logger.exception("Exception while exporting logs.")
255255
finally:
256-
self._metrics.finish_items(1, error)
257256
detach(token)
258257

259258
def shutdown(self):

‎opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/__init__.py‎

Lines changed: 10 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -172,27 +172,20 @@ def _export(self, batch_strategy: BatchExportStrategy) -> None:
172172
while self._should_export_batch(batch_strategy, iteration):
173173
iteration += 1
174174
token = attach(set_value(_SUPPRESS_INSTRUMENTATION_KEY, True))
175-
error: Exception | None = None
176-
count = 0
175+
count = min(
176+
self._max_export_batch_size,
177+
len(self._queue),
178+
)
179+
# Oldest records are at the back, so pop from there.
180+
batch = [self._queue.pop() for _ in range(count)]
181+
# Record on submission to the exporter.
182+
self._metrics.finish_items(count)
177183
try:
178-
count = min(
179-
self._max_export_batch_size,
180-
len(self._queue),
181-
)
182-
self._exporter.export(
183-
[
184-
# Oldest records are at the back, so pop from there.
185-
self._queue.pop()
186-
for _ in range(count)
187-
]
188-
)
189-
except Exception as err: # pylint: disable=broad-exception-caught
190-
error = err
184+
self._exporter.export(batch)
185+
except Exception: # pylint: disable=broad-exception-caught
191186
_logger.exception(
192187
"Exception while exporting %s.", self._exporting
193188
)
194-
finally:
195-
self._metrics.finish_items(count, error)
196189
detach(token)
197190

198191
def emit(self, data: Telemetry) -> None:

‎opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py‎

Lines changed: 4 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ def register_queue_size(
3333

3434
def drop_items(self, count: int) -> None: ...
3535

36-
def finish_items(self, count: int, error: Exception | None) -> None: ...
36+
def finish_items(self, count: int) -> None: ...
3737

3838

3939
class NoOpProcessorMetrics:
@@ -43,7 +43,7 @@ def register_queue_size(self, get_queue_size: Callable[[], int]) -> None:
4343
def drop_items(self, count: int) -> None:
4444
pass
4545

46-
def finish_items(self, count: int, error: Exception | None) -> None:
46+
def finish_items(self, count: int) -> None:
4747
pass
4848

4949

@@ -115,15 +115,8 @@ def record_queue_size(
115115
def drop_items(self, count: int) -> None:
116116
self._processed.add(count, self._dropped_attrs)
117117

118-
def finish_items(self, count: int, error: Exception | None) -> None:
119-
if not error:
120-
self._processed.add(count, self._standard_attrs)
121-
return
122-
attrs = {
123-
**self._standard_attrs,
124-
ERROR_TYPE: type(error).__name__,
125-
}
126-
self._processed.add(count, attrs)
118+
def finish_items(self, count: int) -> None:
119+
self._processed.add(count, self._standard_attrs)
127120

128121

129122
def create_processor_metrics(

‎opentelemetry-sdk/src/opentelemetry/sdk/trace/export/__init__.py‎

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -124,16 +124,15 @@ def on_end(self, span: ReadableSpan) -> None:
124124
if not (span.context and span.context.trace_flags.sampled):
125125
return
126126
token = attach(set_value(_SUPPRESS_INSTRUMENTATION_KEY, True))
127-
error: Exception | None = None
127+
# Record on submission to the exporter.
128+
self._metrics.finish_items(1)
128129
try:
129130
self.span_exporter.export((span,))
130131
# pylint: disable=broad-exception-caught
131-
except Exception as err:
132-
error = err
132+
except Exception:
133133
logger.exception("Exception while exporting Span.")
134134
finally:
135-
self._metrics.finish_items(1, error)
136-
detach(token)
135+
detach(token)
137136

138137
def shutdown(self) -> None:
139138
self.span_exporter.shutdown()

‎opentelemetry-sdk/tests/logs/test_export.py‎

Lines changed: 50 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -433,13 +433,12 @@ def export_logs(_logs):
433433
metrics = sorted(scope_metrics.metrics, key=lambda m: m.name)
434434
self.assertEqual(len(metrics), 1)
435435
self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed")
436-
processed_data_points = sorted(
437-
metrics[0].data.data_points,
438-
key=lambda dp: dp.attributes.get("error.type", ""),
439-
)
440-
self.assertEqual(len(processed_data_points), 2)
436+
processed_data_points = metrics[0].data.data_points
437+
self.assertEqual(len(processed_data_points), 1)
441438
processed_data_point0 = processed_data_points[0]
442-
self.assertEqual(processed_data_point0.value, 2)
439+
# All 3 logs are counted as processed when submitted to the exporter,
440+
# independent of the export outcome (the 3rd export fails).
441+
self.assertEqual(processed_data_point0.value, 3)
443442
self.assertEqual(
444443
processed_data_point0.attributes["otel.component.type"],
445444
"simple_log_processor",
@@ -450,20 +449,39 @@ def export_logs(_logs):
450449
)
451450
)
452451
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
453-
processed_data_point1 = processed_data_points[1]
454-
self.assertEqual(processed_data_point1.value, 1)
455-
self.assertEqual(
456-
processed_data_point1.attributes["otel.component.type"],
457-
"simple_log_processor",
458-
)
459-
self.assertTrue(
460-
processed_data_point1.attributes["otel.component.name"].startswith(
461-
"simple_log_processor/"
462-
)
452+
453+
@patch.dict(
454+
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
455+
)
456+
def test_metrics_not_counted_after_shutdown(self):
457+
metric_reader = InMemoryMetricReader()
458+
meter_provider = MeterProvider(metric_readers=[metric_reader])
459+
460+
exporter = mock.MagicMock()
461+
exporter.export.return_value = LogRecordExportResult.SUCCESS
462+
processor = SimpleLogRecordProcessor(
463+
exporter, meter_provider=meter_provider
463464
)
464-
self.assertEqual(
465-
processed_data_point1.attributes["error.type"], "RuntimeError"
465+
466+
processor.on_emit(EMPTY_LOG)
467+
468+
# Shut only the processor down; the record emitted afterwards hits the
469+
# already-shutdown early return and must not be counted as processed.
470+
processor.shutdown()
471+
processor.on_emit(EMPTY_LOG)
472+
473+
metrics_data = metric_reader.get_metrics_data()
474+
scope_metrics = metrics_data.resource_metrics[0].scope_metrics[0]
475+
metrics = scope_metrics.metrics
476+
self.assertEqual(len(metrics), 1)
477+
self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed")
478+
processed_data_points = metrics[0].data.data_points
479+
self.assertEqual(len(processed_data_points), 1)
480+
self.assertEqual(processed_data_points[0].value, 1)
481+
self.assertIsNone(
482+
processed_data_points[0].attributes.get("error.type")
466483
)
484+
self.assertEqual(exporter.export.call_count, 1)
467485

468486

469487
# Many more test cases for the BatchLogRecordProcessor exist under
@@ -746,7 +764,9 @@ def export_logs(_logs):
746764
metrics[0].data.data_points,
747765
key=lambda dp: dp.attributes.get("error.type", ""),
748766
)
749-
self.assertEqual(len(processed_data_points), 1)
767+
# "foo" is counted as processed when submitted to the exporter (before
768+
# its export call blocks); "baz" is dropped due to a full queue.
769+
self.assertEqual(len(processed_data_points), 2)
750770
processed_data_point0 = processed_data_points[0]
751771
self.assertEqual(processed_data_point0.value, 1)
752772
self.assertEqual(
@@ -758,8 +778,12 @@ def export_logs(_logs):
758778
"batching_log_processor/"
759779
)
760780
)
781+
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
782+
processed_data_point_queue_full = processed_data_points[1]
783+
self.assertEqual(processed_data_point_queue_full.value, 1)
761784
self.assertEqual(
762-
processed_data_point0.attributes.get("error.type"), "queue_full"
785+
processed_data_point_queue_full.attributes.get("error.type"),
786+
"queue_full",
763787
)
764788
self.assertEqual(
765789
metrics[1].name, "otel.sdk.processor.log.queue.capacity"
@@ -806,9 +830,12 @@ def export_logs(_logs):
806830
metrics[0].data.data_points,
807831
key=lambda dp: dp.attributes.get("error.type", ""),
808832
)
809-
self.assertEqual(len(processed_data_points), 3)
833+
# "foo", "bar" and "failed" are all counted as processed when submitted
834+
# to the exporter, independent of the export outcome ("failed" raises).
835+
# "baz" remains a queue_full drop.
836+
self.assertEqual(len(processed_data_points), 2)
810837
processed_data_point0 = processed_data_points[0]
811-
self.assertEqual(processed_data_point0.value, 2)
838+
self.assertEqual(processed_data_point0.value, 3)
812839
self.assertEqual(
813840
processed_data_point0.attributes["otel.component.type"],
814841
"batching_log_processor",
@@ -831,22 +858,7 @@ def export_logs(_logs):
831858
)
832859
)
833860
self.assertEqual(
834-
processed_data_point1.attributes.get("error.type"),
835-
"BrokenPipeError",
836-
)
837-
processed_data_point2 = processed_data_points[2]
838-
self.assertEqual(processed_data_point2.value, 1)
839-
self.assertEqual(
840-
processed_data_point2.attributes["otel.component.type"],
841-
"batching_log_processor",
842-
)
843-
self.assertTrue(
844-
processed_data_point2.attributes["otel.component.name"].startswith(
845-
"batching_log_processor/"
846-
)
847-
)
848-
self.assertEqual(
849-
processed_data_point2.attributes.get("error.type"), "queue_full"
861+
processed_data_point1.attributes.get("error.type"), "queue_full"
850862
)
851863
self.assertEqual(
852864
metrics[1].name, "otel.sdk.processor.log.queue.capacity"

‎opentelemetry-sdk/tests/trace/export/test_export.py‎

Lines changed: 19 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -172,13 +172,12 @@ def export_spans(_spans):
172172
metrics = sorted(scope_metrics.metrics, key=lambda m: m.name)
173173
self.assertEqual(len(metrics), 1)
174174
self.assertEqual(metrics[0].name, "otel.sdk.processor.span.processed")
175-
processed_data_points = sorted(
176-
metrics[0].data.data_points,
177-
key=lambda dp: dp.attributes.get("error.type", ""),
178-
)
179-
self.assertEqual(len(processed_data_points), 2)
175+
processed_data_points = metrics[0].data.data_points
176+
self.assertEqual(len(processed_data_points), 1)
180177
processed_data_point0 = processed_data_points[0]
181-
self.assertEqual(processed_data_point0.value, 2)
178+
# All 3 spans are counted as processed when submitted to the exporter,
179+
# independent of the export outcome (the 3rd export fails).
180+
self.assertEqual(processed_data_point0.value, 3)
182181
self.assertEqual(
183182
processed_data_point0.attributes["otel.component.type"],
184183
"simple_span_processor",
@@ -189,20 +188,6 @@ def export_spans(_spans):
189188
)
190189
)
191190
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
192-
processed_data_point1 = processed_data_points[1]
193-
self.assertEqual(processed_data_point1.value, 1)
194-
self.assertEqual(
195-
processed_data_point1.attributes["otel.component.type"],
196-
"simple_span_processor",
197-
)
198-
self.assertTrue(
199-
processed_data_point1.attributes["otel.component.name"].startswith(
200-
"simple_span_processor/"
201-
)
202-
)
203-
self.assertEqual(
204-
processed_data_point1.attributes["error.type"], "RuntimeError"
205-
)
206191

207192

208193
# Many more test cases for the BatchSpanProcessor exist under
@@ -447,7 +432,9 @@ def export_spans(_spans):
447432
metrics[0].data.data_points,
448433
key=lambda dp: dp.attributes.get("error.type", ""),
449434
)
450-
self.assertEqual(len(processed_data_points), 1)
435+
# "foo" is counted as processed when submitted to the exporter (before
436+
# its export call blocks); "baz" is dropped due to a full queue.
437+
self.assertEqual(len(processed_data_points), 2)
451438
processed_data_point0 = processed_data_points[0]
452439
self.assertEqual(processed_data_point0.value, 1)
453440
self.assertEqual(
@@ -459,8 +446,12 @@ def export_spans(_spans):
459446
"batching_span_processor/"
460447
)
461448
)
449+
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
450+
processed_data_point_queue_full = processed_data_points[1]
451+
self.assertEqual(processed_data_point_queue_full.value, 1)
462452
self.assertEqual(
463-
processed_data_point0.attributes.get("error.type"), "queue_full"
453+
processed_data_point_queue_full.attributes.get("error.type"),
454+
"queue_full",
464455
)
465456
self.assertEqual(
466457
metrics[1].name, "otel.sdk.processor.span.queue.capacity"
@@ -508,9 +499,12 @@ def export_spans(_spans):
508499
metrics[0].data.data_points,
509500
key=lambda dp: dp.attributes.get("error.type", ""),
510501
)
511-
self.assertEqual(len(processed_data_points), 3)
502+
# "foo", "bar" and "failed" are all counted as processed when submitted
503+
# to the exporter, independent of the export outcome ("failed" raises).
504+
# "baz" remains a queue_full drop.
505+
self.assertEqual(len(processed_data_points), 2)
512506
processed_data_point0 = processed_data_points[0]
513-
self.assertEqual(processed_data_point0.value, 2)
507+
self.assertEqual(processed_data_point0.value, 3)
514508
self.assertEqual(
515509
processed_data_point0.attributes["otel.component.type"],
516510
"batching_span_processor",
@@ -533,21 +527,7 @@ def export_spans(_spans):
533527
)
534528
)
535529
self.assertEqual(
536-
processed_data_point1.attributes.get("error.type"), "ValueError"
537-
)
538-
processed_data_point2 = processed_data_points[2]
539-
self.assertEqual(processed_data_point2.value, 1)
540-
self.assertEqual(
541-
processed_data_point2.attributes["otel.component.type"],
542-
"batching_span_processor",
543-
)
544-
self.assertTrue(
545-
processed_data_point2.attributes["otel.component.name"].startswith(
546-
"batching_span_processor/"
547-
)
548-
)
549-
self.assertEqual(
550-
processed_data_point2.attributes.get("error.type"), "queue_full"
530+
processed_data_point1.attributes.get("error.type"), "queue_full"
551531
)
552532
self.assertEqual(
553533
metrics[1].name, "otel.sdk.processor.span.queue.capacity"

0 commit comments

Comments
 (0)