From ad55311956bbec7dc5d147915b4bb6f5c82b63c4 Mon Sep 17 00:00:00 2001 From: Anuraag Agrawal Date: Thu, 6 Aug 2026 12:04:16 +0900 Subject: [PATCH 1/3] Record processed logs before export complete and reject on shutdown --- .../logs/export/BatchLogRecordProcessor.java | 19 +-- ...gacyLogRecordProcessorInstrumentation.java | 14 +- .../LogRecordProcessorInstrumentation.java | 10 +- ...ConvLogRecordProcessorInstrumentation.java | 31 ++-- .../logs/export/SimpleLogRecordProcessor.java | 13 +- .../logs/SdkLoggerProviderMetricsTest.java | 134 ++++++++++++++++-- 6 files changed, 168 insertions(+), 53 deletions(-) diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java index 374d49e7abc..25406317071 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java @@ -159,6 +159,7 @@ private static final class Worker implements Runnable { private final AtomicInteger logsNeeded = new AtomicInteger(Integer.MAX_VALUE); private final BlockingQueue signal; private final AtomicReference flushRequested = new AtomicReference<>(); + private final AtomicBoolean isShutdown = new AtomicBoolean(false); private volatile boolean continueWork = true; private final ArrayList batch; private final long maxQueueSize; @@ -186,9 +187,13 @@ private Worker( } private void addLog(ReadWriteLogRecord logData) { + if (isShutdown.get()) { + logProcessorInstrumentation.dropLogsAlreadyShutdown(1); + return; + } logProcessorInstrumentation.buildQueueMetricsOnce(maxQueueSize, queue::size); if (!queue.offer(logData)) { - logProcessorInstrumentation.dropLogs(1); + logProcessorInstrumentation.dropLogsQueueFull(1); } else { if (queue.size() >= logsNeeded.get()) { signal.offer(true); @@ -251,6 +256,9 @@ private void updateNextExportTime() { } private CompletableResultCode shutdown() { + if (isShutdown.getAndSet(true)) { + return CompletableResultCode.ofSuccess(); + } CompletableResultCode result = new CompletableResultCode(); CompletableResultCode flushResult = forceFlush(); @@ -289,25 +297,18 @@ private void exportCurrentBatch() { return; } - String error = null; try { + logProcessorInstrumentation.finishLogs(batch.size()); CompletableResultCode result = logRecordExporter.export(Collections.unmodifiableList(batch)); result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); if (!result.isSuccess()) { logger.log(Level.FINE, "Exporter failed"); - if (result.getFailureThrowable() != null) { - error = result.getFailureThrowable().getClass().getName(); - } else { - error = "export_failed"; - } } } catch (Throwable t) { ThrowableUtil.propagateIfFatal(t); logger.log(Level.WARNING, "Exporter threw an Exception", t); - error = t.getClass().getName(); } finally { - logProcessorInstrumentation.finishLogs(batch.size(), error); batch.clear(); } } diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LegacyLogRecordProcessorInstrumentation.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LegacyLogRecordProcessorInstrumentation.java index e31c71430fd..33549cadc87 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LegacyLogRecordProcessorInstrumentation.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LegacyLogRecordProcessorInstrumentation.java @@ -43,16 +43,18 @@ final class LegacyLogRecordProcessorInstrumentation implements LogRecordProcesso } @Override - public void dropLogs(int count) { + public void dropLogsQueueFull(int count) { processedLogs().add(count, droppedAttrs); } @Override - public void finishLogs(int count, @Nullable String error) { - // Legacy metrics only record when no error. - if (error == null) { - processedLogs().add(count, standardAttrs); - } + public void dropLogsAlreadyShutdown(int count) { + // Legacy did not record this metric. + } + + @Override + public void finishLogs(int count) { + processedLogs().add(count, standardAttrs); } @Override diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LogRecordProcessorInstrumentation.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LogRecordProcessorInstrumentation.java index 57a9ae6cd33..cc3a17a6fa2 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LogRecordProcessorInstrumentation.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/LogRecordProcessorInstrumentation.java @@ -9,7 +9,6 @@ import io.opentelemetry.sdk.common.InternalTelemetryVersion; import io.opentelemetry.sdk.common.internal.ComponentId; import java.util.function.Supplier; -import javax.annotation.Nullable; /** Metrics exported by span processors. */ interface LogRecordProcessorInstrumentation { @@ -27,10 +26,13 @@ static LogRecordProcessorInstrumentation get( } /** Records metrics for logs dropped because a queue is full. */ - void dropLogs(int count); + void dropLogsQueueFull(int count); - /** Record metrics for logs processed, possibly with an error. */ - void finishLogs(int count, @Nullable String error); + /** Record metrics for logs dropped since processor is shutdown. */ + void dropLogsAlreadyShutdown(int count); + + /** Record metrics for logs processed successfully. */ + void finishLogs(int count); /** Registers metrics for processor queue capacity and size. */ void buildQueueMetricsOnce(long capacity, LongCallable getSize); diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SemConvLogRecordProcessorInstrumentation.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SemConvLogRecordProcessorInstrumentation.java index 618d1b3756a..8a0f1bb9301 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SemConvLogRecordProcessorInstrumentation.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SemConvLogRecordProcessorInstrumentation.java @@ -27,7 +27,8 @@ final class SemConvLogRecordProcessorInstrumentation implements LogRecordProcess private final Supplier meterProvider; private final Attributes standardAttrs; - private final Attributes droppedAttrs; + private final Attributes queueFullAttrs; + private final Attributes shutdownAttrs; @Nullable private Meter meter; @Nullable private volatile LongCounter processedLogs; @@ -42,7 +43,7 @@ final class SemConvLogRecordProcessorInstrumentation implements LogRecordProcess componentId.getTypeName(), SemConvAttributes.OTEL_COMPONENT_NAME, componentId.getComponentName()); - droppedAttrs = + queueFullAttrs = Attributes.of( SemConvAttributes.OTEL_COMPONENT_TYPE, componentId.getTypeName(), @@ -50,23 +51,29 @@ final class SemConvLogRecordProcessorInstrumentation implements LogRecordProcess componentId.getComponentName(), SemConvAttributes.ERROR_TYPE, "queue_full"); + shutdownAttrs = + Attributes.of( + SemConvAttributes.OTEL_COMPONENT_TYPE, + componentId.getTypeName(), + SemConvAttributes.OTEL_COMPONENT_NAME, + componentId.getComponentName(), + SemConvAttributes.ERROR_TYPE, + "already_shutdown"); } @Override - public void dropLogs(int count) { - processedLogs().add(count, droppedAttrs); + public void dropLogsQueueFull(int count) { + processedLogs().add(count, queueFullAttrs); } @Override - public void finishLogs(int count, @Nullable String error) { - if (error == null) { - processedLogs().add(count, standardAttrs); - return; - } + public void dropLogsAlreadyShutdown(int count) { + processedLogs().add(count, shutdownAttrs); + } - Attributes attributes = - standardAttrs.toBuilder().put(SemConvAttributes.ERROR_TYPE, error).build(); - processedLogs().add(count, attributes); + @Override + public void finishLogs(int count) { + processedLogs().add(count, standardAttrs); } @Override diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java index f3d44de8d4a..ee0cee3c709 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java @@ -85,11 +85,17 @@ public static SimpleLogRecordProcessorBuilder builder(LogRecordExporter exporter @Override public void onEmit(Context context, ReadWriteLogRecord logRecord) { + if (isShutdown.get()) { + logProcessorInstrumentation.dropLogsAlreadyShutdown(1); + return; + } + try { List logs = Collections.singletonList(logRecord.toLogRecordData()); CompletableResultCode result; synchronized (exporterLock) { + logProcessorInstrumentation.finishLogs(1); result = logRecordExporter.export(logs); } @@ -97,16 +103,9 @@ public void onEmit(Context context, ReadWriteLogRecord logRecord) { result.whenComplete( () -> { pendingExports.remove(result); - String error = null; if (!result.isSuccess()) { logger.log(Level.FINE, "Exporter failed"); - if (result.getFailureThrowable() != null) { - error = result.getFailureThrowable().getClass().getName(); - } else { - error = "export_failed"; - } } - logProcessorInstrumentation.finishLogs(1, error); }); } catch (RuntimeException e) { logger.log(Level.WARNING, "Exporter threw an Exception", e); diff --git a/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java b/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java index f7a07acab84..8941b378281 100644 --- a/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java +++ b/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java @@ -46,12 +46,11 @@ void simple() { SdkMeterProvider.builder().registerMetricReader(metricReader).build(); InMemoryLogRecordExporter exporter = InMemoryLogRecordExporter.create(); - LoggerProvider loggerProvider = + SimpleLogRecordProcessor processor = + SimpleLogRecordProcessor.builder(exporter).setMeterProvider(() -> meterProvider).build(); + SdkLoggerProvider loggerProvider = SdkLoggerProvider.builder() - .addLogRecordProcessor( - SimpleLogRecordProcessor.builder(exporter) - .setMeterProvider(() -> meterProvider) - .build()) + .addLogRecordProcessor(processor) .setMeterProvider(() -> meterProvider) .build(); @@ -102,6 +101,42 @@ void simple() { "simple_log_processor/0", OTEL_COMPONENT_TYPE, "simple_log_processor"))))); + + // Logs rejected after to call to shutdown, regardless of completion so no join. + processor.shutdown(); + logger.logRecordBuilder().emit(); + + assertThat(metricReader.collectAllMetrics()) + .satisfiesExactlyInAnyOrder( + m -> + assertThat(m) + .hasName("otel.sdk.log.created") + .hasLongSumSatisfying( + s -> s.hasPointsSatisfying(p -> p.hasValue(3).hasAttributes())), + m -> + assertThat(m) + .hasName("otel.sdk.processor.log.processed") + .hasLongSumSatisfying( + s -> + s.hasPointsSatisfying( + p -> + p.hasValue(2) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "simple_log_processor/0", + OTEL_COMPONENT_TYPE, + "simple_log_processor")), + p -> + p.hasValue(1) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "simple_log_processor/0", + OTEL_COMPONENT_TYPE, + "simple_log_processor", + ERROR_TYPE, + "already_shutdown"))))); } @Test @@ -134,11 +169,12 @@ void batch() throws Exception { // Will immediately be processed. logger.logRecordBuilder().emit(); Thread.sleep(500); // give time to start processing a batch of size 1 - // We haven't completed the export so this span is queued. + // We haven't completed the export so this log is queued. logger.logRecordBuilder().emit(); - // Queue is full, this span is dropped. + // Queue is full, this log is dropped. logger.logRecordBuilder().emit(); + // Export hasn't finished but processed metric is incremented assertThat(metricReader.collectAllMetrics()) .satisfiesExactlyInAnyOrder( m -> @@ -175,6 +211,14 @@ void batch() throws Exception { .hasLongSumSatisfying( s -> s.hasPointsSatisfying( + p -> + p.hasValue(1) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "batching_log_processor/0", + OTEL_COMPONENT_TYPE, + "batching_log_processor")), p -> p.hasValue(1) .hasAttributes( @@ -231,8 +275,73 @@ void batch() throws Exception { .hasLongSumSatisfying( s -> s.hasPointsSatisfying( + p -> + p.hasValue(2) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "batching_log_processor/0", + OTEL_COMPONENT_TYPE, + "batching_log_processor")), p -> p.hasValue(1) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "batching_log_processor/0", + OTEL_COMPONENT_TYPE, + "batching_log_processor", + ERROR_TYPE, + "queue_full")))), + m -> + assertThat(m) + .hasName("otel.sdk.log.created") + .hasLongSumSatisfying( + s -> s.hasPointsSatisfying(p -> p.hasValue(3).hasAttributes()))); + + lenient().when(mockExporter.shutdown()).thenReturn(CompletableResultCode.ofSuccess()); + // Logs rejected after to call to shutdown, regardless of completion so no join. + processor.shutdown(); + + logger.logRecordBuilder().emit(); + assertThat(metricReader.collectAllMetrics()) + .satisfiesExactlyInAnyOrder( + m -> + assertThat(m) + .hasName("otel.sdk.processor.log.queue.capacity") + .hasLongSumSatisfying( + s -> + s.hasPointsSatisfying( + p -> + p.hasValue(1) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "batching_log_processor/0", + OTEL_COMPONENT_TYPE, + "batching_log_processor")))), + m -> + assertThat(m) + .hasName("otel.sdk.processor.log.queue.size") + .hasLongSumSatisfying( + s -> + s.hasPointsSatisfying( + p -> + p.hasValue(0) + .hasAttributes( + Attributes.of( + OTEL_COMPONENT_NAME, + "batching_log_processor/0", + OTEL_COMPONENT_TYPE, + "batching_log_processor")))), + m -> + assertThat(m) + .hasName("otel.sdk.processor.log.processed") + .hasLongSumSatisfying( + s -> + s.hasPointsSatisfying( + p -> + p.hasValue(2) .hasAttributes( Attributes.of( OTEL_COMPONENT_NAME, @@ -248,7 +357,7 @@ void batch() throws Exception { OTEL_COMPONENT_TYPE, "batching_log_processor", ERROR_TYPE, - "export_failed")), + "already_shutdown")), p -> p.hasValue(1) .hasAttributes( @@ -263,10 +372,7 @@ void batch() throws Exception { assertThat(m) .hasName("otel.sdk.log.created") .hasLongSumSatisfying( - s -> s.hasPointsSatisfying(p -> p.hasValue(3).hasAttributes()))); - - lenient().when(mockExporter.shutdown()).thenReturn(CompletableResultCode.ofSuccess()); - processor.shutdown(); + s -> s.hasPointsSatisfying(p -> p.hasValue(4).hasAttributes()))); } @Test @@ -305,9 +411,7 @@ void simpleExportError() { OTEL_COMPONENT_NAME, "simple_log_processor/0", OTEL_COMPONENT_TYPE, - "simple_log_processor", - ERROR_TYPE, - "export_failed")))), + "simple_log_processor")))), m -> assertThat(m) .hasName("otel.sdk.log.created") From 4a54c280bc7e3e6c6ccf5318cffb3e166cdd64c1 Mon Sep 17 00:00:00 2001 From: Anuraag Agrawal Date: Thu, 6 Aug 2026 12:11:58 +0900 Subject: [PATCH 2/3] grammar --- .../opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java b/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java index 8941b378281..d094faef9f6 100644 --- a/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java +++ b/sdk/logs/src/test/java/io/opentelemetry/sdk/logs/SdkLoggerProviderMetricsTest.java @@ -102,7 +102,7 @@ void simple() { OTEL_COMPONENT_TYPE, "simple_log_processor"))))); - // Logs rejected after to call to shutdown, regardless of completion so no join. + // Logs rejected after the call to shutdown, regardless of completion so no join. processor.shutdown(); logger.logRecordBuilder().emit(); @@ -300,7 +300,7 @@ void batch() throws Exception { s -> s.hasPointsSatisfying(p -> p.hasValue(3).hasAttributes()))); lenient().when(mockExporter.shutdown()).thenReturn(CompletableResultCode.ofSuccess()); - // Logs rejected after to call to shutdown, regardless of completion so no join. + // Logs rejected after the call to shutdown, regardless of completion so no join. processor.shutdown(); logger.logRecordBuilder().emit(); From 85a7b3f54d6220af45e349930e1a015de0398f81 Mon Sep 17 00:00:00 2001 From: Anuraag Agrawal Date: Fri, 7 Aug 2026 11:03:37 +0900 Subject: [PATCH 3/3] cleanup --- .../sdk/logs/export/BatchLogRecordProcessor.java | 6 ++---- .../sdk/logs/export/SimpleLogRecordProcessor.java | 2 ++ 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java index 25406317071..306ca41f66d 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/BatchLogRecordProcessor.java @@ -49,7 +49,6 @@ public final class BatchLogRecordProcessor implements LogRecordProcessor { BatchLogRecordProcessor.class.getSimpleName() + "_WorkerThread"; private final Worker worker; - private final AtomicBoolean isShutdown = new AtomicBoolean(false); /** * Returns a new Builder for {@link BatchLogRecordProcessor}. @@ -94,9 +93,6 @@ public void onEmit(Context context, ReadWriteLogRecord logRecord) { @Override public CompletableResultCode shutdown() { - if (isShutdown.getAndSet(true)) { - return CompletableResultCode.ofSuccess(); - } return worker.shutdown(); } @@ -298,6 +294,8 @@ private void exportCurrentBatch() { } try { + // We always increment for every export invocation, so we increment before the export call + // to make sure thrown errors don't affect it. logProcessorInstrumentation.finishLogs(batch.size()); CompletableResultCode result = logRecordExporter.export(Collections.unmodifiableList(batch)); diff --git a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java index ee0cee3c709..1dcfe2d2276 100644 --- a/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java +++ b/sdk/logs/src/main/java/io/opentelemetry/sdk/logs/export/SimpleLogRecordProcessor.java @@ -95,6 +95,8 @@ public void onEmit(Context context, ReadWriteLogRecord logRecord) { CompletableResultCode result; synchronized (exporterLock) { + // We always increment for every export invocation, so we increment before the export call + // to make sure thrown errors don't affect it. logProcessorInstrumentation.finishLogs(1); result = logRecordExporter.export(logs); }