Skip to content
1 change: 1 addition & 0 deletions .changelog/5472.fixed
Original file line number Diff line number Diff line change
@@ -0,0 +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`
Original file line number Diff line number Diff line change
Expand Up @@ -230,7 +230,6 @@ def on_emit(self, log_record: ReadWriteLogRecord):
set_value(_ON_EMIT_RECURSION_COUNT_KEY, cnt + 1), # pyright: ignore[reportOperatorIssue]
)
)
error: Exception | None = None
try:
if self._shutdown:
_logger.warning("Processor is already shutdown, ignoring call")
Expand All @@ -248,12 +247,12 @@ def on_emit(self, log_record: ReadWriteLogRecord):
instrumentation_scope=log_record.instrumentation_scope,
limits=log_record.limits,
)
# Record on submission to the exporter.
self._metrics.finish_items(1)
self._exporter.export((readable_log_record,))
except Exception as err: # pylint: disable=broad-exception-caught
error = err
except Exception: # pylint: disable=broad-exception-caught
_logger.exception("Exception while exporting logs.")
finally:
self._metrics.finish_items(1, error)
detach(token)

def shutdown(self):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -172,27 +172,20 @@ def _export(self, batch_strategy: BatchExportStrategy) -> None:
while self._should_export_batch(batch_strategy, iteration):
iteration += 1
token = attach(set_value(_SUPPRESS_INSTRUMENTATION_KEY, True))
error: Exception | None = None
count = 0
count = min(
Comment thread
DylanRussell marked this conversation as resolved.
self._max_export_batch_size,
len(self._queue),
)
# Oldest records are at the back, so pop from there.
batch = [self._queue.pop() for _ in range(count)]
# Record on submission to the exporter.
self._metrics.finish_items(count)
try:
count = min(
self._max_export_batch_size,
len(self._queue),
)
self._exporter.export(
[
# Oldest records are at the back, so pop from there.
self._queue.pop()
for _ in range(count)
]
)
except Exception as err: # pylint: disable=broad-exception-caught
error = err
self._exporter.export(batch)
except Exception: # pylint: disable=broad-exception-caught
_logger.exception(
"Exception while exporting %s.", self._exporting
)
finally:
self._metrics.finish_items(count, error)
detach(token)

def emit(self, data: Telemetry) -> None:
Expand Down
Comment thread
cijothomas marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ def register_queue_size(

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

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


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

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


Expand Down Expand Up @@ -115,15 +115,8 @@ def record_queue_size(
def drop_items(self, count: int) -> None:
self._processed.add(count, self._dropped_attrs)

def finish_items(self, count: int, error: Exception | None) -> None:
if not error:
self._processed.add(count, self._standard_attrs)
return
attrs = {
**self._standard_attrs,
ERROR_TYPE: type(error).__name__,
}
self._processed.add(count, attrs)
def finish_items(self, count: int) -> None:
self._processed.add(count, self._standard_attrs)


def create_processor_metrics(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -124,16 +124,15 @@ def on_end(self, span: ReadableSpan) -> None:
if not (span.context and span.context.trace_flags.sampled):
return
token = attach(set_value(_SUPPRESS_INSTRUMENTATION_KEY, True))
error: Exception | None = None
# Record on submission to the exporter.
self._metrics.finish_items(1)
try:
self.span_exporter.export((span,))
# pylint: disable=broad-exception-caught
except Exception as err:
error = err
except Exception:
logger.exception("Exception while exporting Span.")
finally:
self._metrics.finish_items(1, error)
detach(token)
detach(token)

def shutdown(self) -> None:
self.span_exporter.shutdown()
Expand Down
88 changes: 50 additions & 38 deletions opentelemetry-sdk/tests/logs/test_export.py
Original file line number Diff line number Diff line change
Expand Up @@ -433,13 +433,12 @@ def export_logs(_logs):
metrics = sorted(scope_metrics.metrics, key=lambda m: m.name)
self.assertEqual(len(metrics), 1)
self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed")
processed_data_points = sorted(
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 2)
processed_data_points = metrics[0].data.data_points
self.assertEqual(len(processed_data_points), 1)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 2)
# All 3 logs are counted as processed when submitted to the exporter,
# independent of the export outcome (the 3rd export fails).
self.assertEqual(processed_data_point0.value, 3)
self.assertEqual(
processed_data_point0.attributes["otel.component.type"],
"simple_log_processor",
Expand All @@ -450,20 +449,39 @@ def export_logs(_logs):
)
)
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
processed_data_point1 = processed_data_points[1]
self.assertEqual(processed_data_point1.value, 1)
self.assertEqual(
processed_data_point1.attributes["otel.component.type"],
"simple_log_processor",
)
self.assertTrue(
processed_data_point1.attributes["otel.component.name"].startswith(
"simple_log_processor/"
)

@patch.dict(
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
)
def test_metrics_not_counted_after_shutdown(self):
metric_reader = InMemoryMetricReader()
meter_provider = MeterProvider(metric_readers=[metric_reader])

exporter = mock.MagicMock()
exporter.export.return_value = LogRecordExportResult.SUCCESS
processor = SimpleLogRecordProcessor(
exporter, meter_provider=meter_provider
)
self.assertEqual(
processed_data_point1.attributes["error.type"], "RuntimeError"

processor.on_emit(EMPTY_LOG)

# Shut only the processor down; the record emitted afterwards hits the
# already-shutdown early return and must not be counted as processed.
processor.shutdown()
processor.on_emit(EMPTY_LOG)

metrics_data = metric_reader.get_metrics_data()
scope_metrics = metrics_data.resource_metrics[0].scope_metrics[0]
metrics = scope_metrics.metrics
self.assertEqual(len(metrics), 1)
self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed")
processed_data_points = metrics[0].data.data_points
self.assertEqual(len(processed_data_points), 1)
self.assertEqual(processed_data_points[0].value, 1)
self.assertIsNone(
processed_data_points[0].attributes.get("error.type")
)
self.assertEqual(exporter.export.call_count, 1)


# Many more test cases for the BatchLogRecordProcessor exist under
Expand Down Expand Up @@ -746,7 +764,9 @@ def export_logs(_logs):
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 1)
# "foo" is counted as processed when submitted to the exporter (before
# its export call blocks); "baz" is dropped due to a full queue.
self.assertEqual(len(processed_data_points), 2)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 1)
self.assertEqual(
Expand All @@ -758,8 +778,12 @@ def export_logs(_logs):
"batching_log_processor/"
)
)
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
processed_data_point_queue_full = processed_data_points[1]
self.assertEqual(processed_data_point_queue_full.value, 1)
self.assertEqual(
processed_data_point0.attributes.get("error.type"), "queue_full"
processed_data_point_queue_full.attributes.get("error.type"),
"queue_full",
)
self.assertEqual(
metrics[1].name, "otel.sdk.processor.log.queue.capacity"
Expand Down Expand Up @@ -806,9 +830,12 @@ def export_logs(_logs):
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 3)
# "foo", "bar" and "failed" are all counted as processed when submitted
# to the exporter, independent of the export outcome ("failed" raises).
# "baz" remains a queue_full drop.
self.assertEqual(len(processed_data_points), 2)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 2)
self.assertEqual(processed_data_point0.value, 3)
self.assertEqual(
processed_data_point0.attributes["otel.component.type"],
"batching_log_processor",
Expand All @@ -831,22 +858,7 @@ def export_logs(_logs):
)
)
self.assertEqual(
processed_data_point1.attributes.get("error.type"),
"BrokenPipeError",
)
processed_data_point2 = processed_data_points[2]
self.assertEqual(processed_data_point2.value, 1)
self.assertEqual(
processed_data_point2.attributes["otel.component.type"],
"batching_log_processor",
)
self.assertTrue(
processed_data_point2.attributes["otel.component.name"].startswith(
"batching_log_processor/"
)
)
self.assertEqual(
processed_data_point2.attributes.get("error.type"), "queue_full"
processed_data_point1.attributes.get("error.type"), "queue_full"
)
self.assertEqual(
metrics[1].name, "otel.sdk.processor.log.queue.capacity"
Expand Down
58 changes: 19 additions & 39 deletions opentelemetry-sdk/tests/trace/export/test_export.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,13 +172,12 @@ def export_spans(_spans):
metrics = sorted(scope_metrics.metrics, key=lambda m: m.name)
self.assertEqual(len(metrics), 1)
self.assertEqual(metrics[0].name, "otel.sdk.processor.span.processed")
processed_data_points = sorted(
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 2)
processed_data_points = metrics[0].data.data_points
self.assertEqual(len(processed_data_points), 1)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 2)
# All 3 spans are counted as processed when submitted to the exporter,
# independent of the export outcome (the 3rd export fails).
self.assertEqual(processed_data_point0.value, 3)
self.assertEqual(
processed_data_point0.attributes["otel.component.type"],
"simple_span_processor",
Expand All @@ -189,20 +188,6 @@ def export_spans(_spans):
)
)
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
processed_data_point1 = processed_data_points[1]
self.assertEqual(processed_data_point1.value, 1)
self.assertEqual(
processed_data_point1.attributes["otel.component.type"],
"simple_span_processor",
)
self.assertTrue(
processed_data_point1.attributes["otel.component.name"].startswith(
"simple_span_processor/"
)
)
self.assertEqual(
processed_data_point1.attributes["error.type"], "RuntimeError"
)


# Many more test cases for the BatchSpanProcessor exist under
Expand Down Expand Up @@ -447,7 +432,9 @@ def export_spans(_spans):
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 1)
# "foo" is counted as processed when submitted to the exporter (before
# its export call blocks); "baz" is dropped due to a full queue.
self.assertEqual(len(processed_data_points), 2)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 1)
self.assertEqual(
Expand All @@ -459,8 +446,12 @@ def export_spans(_spans):
"batching_span_processor/"
)
)
self.assertIsNone(processed_data_point0.attributes.get("error.type"))
processed_data_point_queue_full = processed_data_points[1]
self.assertEqual(processed_data_point_queue_full.value, 1)
self.assertEqual(
processed_data_point0.attributes.get("error.type"), "queue_full"
processed_data_point_queue_full.attributes.get("error.type"),
"queue_full",
)
self.assertEqual(
metrics[1].name, "otel.sdk.processor.span.queue.capacity"
Expand Down Expand Up @@ -508,9 +499,12 @@ def export_spans(_spans):
metrics[0].data.data_points,
key=lambda dp: dp.attributes.get("error.type", ""),
)
self.assertEqual(len(processed_data_points), 3)
# "foo", "bar" and "failed" are all counted as processed when submitted
# to the exporter, independent of the export outcome ("failed" raises).
# "baz" remains a queue_full drop.
self.assertEqual(len(processed_data_points), 2)
processed_data_point0 = processed_data_points[0]
self.assertEqual(processed_data_point0.value, 2)
self.assertEqual(processed_data_point0.value, 3)
self.assertEqual(
processed_data_point0.attributes["otel.component.type"],
"batching_span_processor",
Expand All @@ -533,21 +527,7 @@ def export_spans(_spans):
)
)
self.assertEqual(
processed_data_point1.attributes.get("error.type"), "ValueError"
)
processed_data_point2 = processed_data_points[2]
self.assertEqual(processed_data_point2.value, 1)
self.assertEqual(
processed_data_point2.attributes["otel.component.type"],
"batching_span_processor",
)
self.assertTrue(
processed_data_point2.attributes["otel.component.name"].startswith(
"batching_span_processor/"
)
)
self.assertEqual(
processed_data_point2.attributes.get("error.type"), "queue_full"
processed_data_point1.attributes.get("error.type"), "queue_full"
)
self.assertEqual(
metrics[1].name, "otel.sdk.processor.span.queue.capacity"
Expand Down
Loading