From 41910f3f7edc473e190a935590364492dc254c61 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sun, 2 Aug 2026 17:47:07 +0530 Subject: [PATCH 1/9] feat: align PeriodicMetricReader export timeout semantics --- .../metrics/export/PeriodicMetricReader.java | 52 ++++++++----------- .../export/PeriodicMetricReaderBuilder.java | 34 +++++++++++- 2 files changed, 54 insertions(+), 32 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index 14e3751ba18..f74de25fd7e 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -46,6 +46,7 @@ public final class PeriodicMetricReader implements MetricReader { private final MetricExporter exporter; private final long intervalNanos; + private final long exporterTimeoutNanos; private final ScheduledExecutorService scheduler; private final Scheduled scheduled; private final Object lock = new Object(); @@ -72,11 +73,13 @@ public static PeriodicMetricReaderBuilder builder(MetricExporter exporter) { PeriodicMetricReader( MetricExporter exporter, long intervalNanos, + long exporterTimeoutNanos, ScheduledExecutorService scheduler, int maxExportBatchSize, InternalTelemetryVersion internalTelemetryVersion) { this.exporter = exporter; this.intervalNanos = intervalNanos; + this.exporterTimeoutNanos = exporterTimeoutNanos; this.scheduler = scheduler; this.maxExportBatchSize = maxExportBatchSize; this.scheduled = new Scheduled(); @@ -213,43 +216,30 @@ private Scheduled() {} private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { - return exporter.export(metricData); + CompletableResultCode result = exporter.export(metricData); + result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + return result.isSuccess() + ? CompletableResultCode.ofSuccess() + : CompletableResultCode.ofFailure(); } Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); CompletableResultCode sequentialResult = new CompletableResultCode(); AtomicBoolean anyFailed = new AtomicBoolean(false); Iterator> batchIterator = batches.iterator(); - Runnable exportNext = - new Runnable() { - @Override - public void run() { - while (batchIterator.hasNext()) { - Collection currentBatch = batchIterator.next(); - CompletableResultCode currentResult = exporter.export(currentBatch); - if (currentResult.isDone()) { - if (!currentResult.isSuccess()) { - anyFailed.set(true); - } - } else { - currentResult.whenComplete( - () -> { - if (!currentResult.isSuccess()) { - anyFailed.set(true); - } - this.run(); - }); - return; - } - } - if (anyFailed.get()) { - sequentialResult.fail(); - } else { - sequentialResult.succeed(); - } - } - }; - exportNext.run(); + while (batchIterator.hasNext()) { + Collection currentBatch = batchIterator.next(); + CompletableResultCode currentResult = exporter.export(currentBatch); + currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + if (!currentResult.isSuccess()) { + anyFailed.set(true); + } + } + if (anyFailed.get()) { + sequentialResult.fail(); + } else { + sequentialResult.succeed(); + } return sequentialResult; } diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java index 0a2ef9ea900..d147e461b61 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java @@ -24,6 +24,7 @@ public final class PeriodicMetricReaderBuilder { static final long DEFAULT_SCHEDULE_DELAY_MINUTES = 1; + static final int DEFAULT_EXPORT_TIMEOUT_MILLIS = 30_000; private final MetricExporter metricExporter; @@ -31,6 +32,8 @@ public final class PeriodicMetricReaderBuilder { private long intervalNanos = TimeUnit.MINUTES.toNanos(DEFAULT_SCHEDULE_DELAY_MINUTES); + private long exporterTimeoutNanos = TimeUnit.MILLISECONDS.toNanos(DEFAULT_EXPORT_TIMEOUT_MILLIS); + @Nullable private ScheduledExecutorService executor; private int maxExportBatchSize; @@ -57,6 +60,30 @@ public PeriodicMetricReaderBuilder setInterval(Duration interval) { return setInterval(interval.toNanos(), TimeUnit.NANOSECONDS); } + /** + * Sets the timeout for the underlying exporter. If unset, defaults to {@value + * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * + * @since 1.40.0 + */ + public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit unit) { + requireNonNull(unit, "unit"); + checkArgument(timeout >= 0, "timeout must be non-negative"); + exporterTimeoutNanos = timeout == 0 ? Long.MAX_VALUE : unit.toNanos(timeout); + return this; + } + + /** + * Sets the timeout for the underlying exporter. If unset, defaults to {@value + * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * + * @since 1.40.0 + */ + public PeriodicMetricReaderBuilder setExporterTimeout(Duration timeout) { + requireNonNull(timeout, "timeout"); + return setExporterTimeout(timeout.toNanos(), TimeUnit.NANOSECONDS); + } + /** Sets the {@link ScheduledExecutorService} to schedule reads on. */ public PeriodicMetricReaderBuilder setExecutor(ScheduledExecutorService executor) { requireNonNull(executor, "executor"); @@ -86,7 +113,12 @@ public PeriodicMetricReader build() { Executors.newScheduledThreadPool(1, new DaemonThreadFactory("PeriodicMetricReader")); } return new PeriodicMetricReader( - metricExporter, intervalNanos, executor, maxExportBatchSize, internalTelemetryVersion); + metricExporter, + intervalNanos, + exporterTimeoutNanos, + executor, + maxExportBatchSize, + internalTelemetryVersion); } /** Sets the internal telemetry version used to control self-observability metrics. */ From e371867d2bebb3fd87988aeedf8b4e3a4dd41be7 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sun, 2 Aug 2026 17:57:06 +0530 Subject: [PATCH 2/9] fix: resolve test failures by implementing export timeout asynchronously --- .../metrics/export/PeriodicMetricReader.java | 70 ++++++++++++++----- 1 file changed, 52 insertions(+), 18 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index f74de25fd7e..ce11533d1ad 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -214,32 +214,66 @@ private final class Scheduled implements Runnable { private Scheduled() {} + private CompletableResultCode withTimeout(CompletableResultCode result) { + if (result.isDone() || exporterTimeoutNanos == Long.MAX_VALUE) { + return result; + } + CompletableResultCode timeoutResult = new CompletableResultCode(); + ScheduledFuture timeoutFuture = + scheduler.schedule(timeoutResult::fail, exporterTimeoutNanos, TimeUnit.NANOSECONDS); + result.whenComplete( + () -> { + if (timeoutFuture != null) { + timeoutFuture.cancel(false); + } + if (result.isSuccess()) { + timeoutResult.succeed(); + } else { + timeoutResult.fail(); + } + }); + return timeoutResult; + } + private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { - CompletableResultCode result = exporter.export(metricData); - result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); - return result.isSuccess() - ? CompletableResultCode.ofSuccess() - : CompletableResultCode.ofFailure(); + return withTimeout(exporter.export(metricData)); } Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); CompletableResultCode sequentialResult = new CompletableResultCode(); AtomicBoolean anyFailed = new AtomicBoolean(false); Iterator> batchIterator = batches.iterator(); - while (batchIterator.hasNext()) { - Collection currentBatch = batchIterator.next(); - CompletableResultCode currentResult = exporter.export(currentBatch); - currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); - if (!currentResult.isSuccess()) { - anyFailed.set(true); - } - } - if (anyFailed.get()) { - sequentialResult.fail(); - } else { - sequentialResult.succeed(); - } + Runnable exportNext = + new Runnable() { + @Override + public void run() { + while (batchIterator.hasNext()) { + Collection currentBatch = batchIterator.next(); + CompletableResultCode currentResult = withTimeout(exporter.export(currentBatch)); + if (currentResult.isDone()) { + if (!currentResult.isSuccess()) { + anyFailed.set(true); + } + } else { + currentResult.whenComplete( + () -> { + if (!currentResult.isSuccess()) { + anyFailed.set(true); + } + this.run(); + }); + return; + } + } + if (anyFailed.get()) { + sequentialResult.fail(); + } else { + sequentialResult.succeed(); + } + } + }; + exportNext.run(); return sequentialResult; } From d3fc20b5cba18e0ace6c65c664845bf6e090be42 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sun, 2 Aug 2026 18:21:55 +0530 Subject: [PATCH 3/9] fix: use CompletableResultCode.join() for export timeout instead of scheduler Replace scheduler-based withTimeout() with CompletableResultCode.join(), matching the established pattern in BatchSpanProcessor and BatchLogRecordProcessor. The previous approach used scheduler.schedule() which throws RejectedExecutionException during shutdown because the scheduler is intentionally shut down before the final export flush. --- .../metrics/export/PeriodicMetricReader.java | 68 +++++-------------- 1 file changed, 16 insertions(+), 52 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index ce11533d1ad..dba819120aa 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -214,66 +214,30 @@ private final class Scheduled implements Runnable { private Scheduled() {} - private CompletableResultCode withTimeout(CompletableResultCode result) { - if (result.isDone() || exporterTimeoutNanos == Long.MAX_VALUE) { - return result; - } - CompletableResultCode timeoutResult = new CompletableResultCode(); - ScheduledFuture timeoutFuture = - scheduler.schedule(timeoutResult::fail, exporterTimeoutNanos, TimeUnit.NANOSECONDS); - result.whenComplete( - () -> { - if (timeoutFuture != null) { - timeoutFuture.cancel(false); - } - if (result.isSuccess()) { - timeoutResult.succeed(); - } else { - timeoutResult.fail(); - } - }); - return timeoutResult; - } - private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { - return withTimeout(exporter.export(metricData)); + CompletableResultCode result = exporter.export(metricData); + result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + return result; } Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); CompletableResultCode sequentialResult = new CompletableResultCode(); AtomicBoolean anyFailed = new AtomicBoolean(false); Iterator> batchIterator = batches.iterator(); - Runnable exportNext = - new Runnable() { - @Override - public void run() { - while (batchIterator.hasNext()) { - Collection currentBatch = batchIterator.next(); - CompletableResultCode currentResult = withTimeout(exporter.export(currentBatch)); - if (currentResult.isDone()) { - if (!currentResult.isSuccess()) { - anyFailed.set(true); - } - } else { - currentResult.whenComplete( - () -> { - if (!currentResult.isSuccess()) { - anyFailed.set(true); - } - this.run(); - }); - return; - } - } - if (anyFailed.get()) { - sequentialResult.fail(); - } else { - sequentialResult.succeed(); - } - } - }; - exportNext.run(); + while (batchIterator.hasNext()) { + Collection currentBatch = batchIterator.next(); + CompletableResultCode currentResult = exporter.export(currentBatch); + currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + if (!currentResult.isSuccess()) { + anyFailed.set(true); + } + } + if (anyFailed.get()) { + sequentialResult.fail(); + } else { + sequentialResult.succeed(); + } return sequentialResult; } From b78cae9ca2830c208e9d3f5035dc317703a67e11 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sun, 2 Aug 2026 18:41:10 +0530 Subject: [PATCH 4/9] fix: resolve errorprone warnings --- .../sdk/metrics/export/PeriodicMetricReader.java | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index dba819120aa..8b8d4aaabd2 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -18,7 +18,6 @@ import io.opentelemetry.sdk.metrics.data.AggregationTemporality; import io.opentelemetry.sdk.metrics.data.MetricData; import java.util.Collection; -import java.util.Iterator; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -223,17 +222,15 @@ private CompletableResultCode exportMetrics(Collection metricData) { Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); CompletableResultCode sequentialResult = new CompletableResultCode(); - AtomicBoolean anyFailed = new AtomicBoolean(false); - Iterator> batchIterator = batches.iterator(); - while (batchIterator.hasNext()) { - Collection currentBatch = batchIterator.next(); + boolean anyFailed = false; + for (Collection currentBatch : batches) { CompletableResultCode currentResult = exporter.export(currentBatch); currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); if (!currentResult.isSuccess()) { - anyFailed.set(true); + anyFailed = true; } } - if (anyFailed.get()) { + if (anyFailed) { sequentialResult.fail(); } else { sequentialResult.succeed(); From 1b69dec5cd1930b35a6b7d0c1826b529655f519f Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Fri, 21 Aug 2026 11:11:39 +0530 Subject: [PATCH 5/9] feat: align PeriodicMetricReader export timeout semantics with specification - Added setExporterTimeout() API to PeriodicMetricReaderBuilder with 30-second default per spec - Removed @since 1.40.0 annotations as requested by reviewer - Updated API-diff to reflect new public API - Added comprehensive tests for timeout behavior - Implemented conditional timeout enforcement (only when explicitly configured) - Addressed blocking concern by making timeout opt-in via setExporterTimeout() --- .../opentelemetry-sdk-metrics.txt | 2 + .../metrics/export/PeriodicMetricReader.java | 8 +- .../export/PeriodicMetricReaderBuilder.java | 4 - .../export/PeriodicMetricReaderTest.java | 192 ++++++++++++++++++ 4 files changed, 200 insertions(+), 6 deletions(-) diff --git a/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt b/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt index 9eef89292bc..fbcbf09c9ad 100644 --- a/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt +++ b/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt @@ -1,4 +1,6 @@ Comparing source compatibility of opentelemetry-sdk-metrics-1.65.0-SNAPSHOT.jar against opentelemetry-sdk-metrics-1.64.0.jar *** MODIFIED CLASS: PUBLIC FINAL io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder (not serializable) === CLASS FILE FORMAT VERSION: 52.0 <- 52.0 + +++ NEW METHOD: PUBLIC(+) io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder setExporterTimeout(long, java.util.concurrent.TimeUnit) + +++ NEW METHOD: PUBLIC(+) io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder setExporterTimeout(java.time.Duration) +++ NEW METHOD: PUBLIC(+) io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder setInternalTelemetryVersion(io.opentelemetry.sdk.common.InternalTelemetryVersion) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index 8b8d4aaabd2..7cb8b0f0e48 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -216,7 +216,9 @@ private Scheduled() {} private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { CompletableResultCode result = exporter.export(metricData); - result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + if (exporterTimeoutNanos != Long.MAX_VALUE) { + result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + } return result; } Collection> batches = @@ -225,7 +227,9 @@ private CompletableResultCode exportMetrics(Collection metricData) { boolean anyFailed = false; for (Collection currentBatch : batches) { CompletableResultCode currentResult = exporter.export(currentBatch); - currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + if (exporterTimeoutNanos != Long.MAX_VALUE) { + currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + } if (!currentResult.isSuccess()) { anyFailed = true; } diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java index d147e461b61..c328d207b48 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java @@ -63,8 +63,6 @@ public PeriodicMetricReaderBuilder setInterval(Duration interval) { /** * Sets the timeout for the underlying exporter. If unset, defaults to {@value * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. - * - * @since 1.40.0 */ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit unit) { requireNonNull(unit, "unit"); @@ -76,8 +74,6 @@ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit uni /** * Sets the timeout for the underlying exporter. If unset, defaults to {@value * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. - * - * @since 1.40.0 */ public PeriodicMetricReaderBuilder setExporterTimeout(Duration timeout) { requireNonNull(timeout, "timeout"); diff --git a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java index df00a7e14d4..eec6f4092f8 100644 --- a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java +++ b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java @@ -717,6 +717,198 @@ void stringRepresentation() { + "}"); } + @Test + void setExporterTimeout_configuresTimeout() { + PeriodicMetricReaderBuilder builder = + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(5, TimeUnit.SECONDS); + assertThat(builder).isNotNull(); + } + + @Test + void setExporterTimeout_withDuration_configuresTimeout() { + PeriodicMetricReaderBuilder builder = + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ofSeconds(5)); + assertThat(builder).isNotNull(); + } + + @Test + void setExporterTimeout_zeroMeansNoTimeout() { + PeriodicMetricReaderBuilder builder = + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(0, TimeUnit.SECONDS); + assertThat(builder).isNotNull(); + } + + @Test + void setExporterTimeout_negativeThrowsException() { + assertThatThrownBy( + () -> + PeriodicMetricReader.builder(metricExporter) + .setExporterTimeout(-1, TimeUnit.SECONDS)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("timeout must be non-negative"); + } + + @Test + void setExporterTimeout_nullUnitThrowsException() { + assertThatThrownBy( + () -> PeriodicMetricReader.builder(metricExporter).setExporterTimeout(5, null)) + .isInstanceOf(NullPointerException.class) + .hasMessage("unit"); + } + + @Test + void setExporterTimeout_nullDurationThrowsException() { + assertThatThrownBy( + () -> PeriodicMetricReader.builder(metricExporter).setExporterTimeout((Duration) null)) + .isInstanceOf(NullPointerException.class) + .hasMessage("timeout"); + } + + @Test + void defaultBehavior_applies30SecondTimeout() throws Exception { + // This test verifies that by default, a 30-second timeout is applied + // We test this by having an exporter that takes less than 30 seconds + // and ensuring it completes successfully + SlowMetricExporter slowExporter = new SlowMetricExporter(100); // 100ms delay + PeriodicMetricReader reader = + PeriodicMetricReader.builder(slowExporter) + .setInterval(Duration.ofMillis(50)) + .build(); + + reader.register(collectionRegistration); + try { + // Wait for an export to complete - should succeed since 100ms < 30s + assertThat(slowExporter.waitForExport(1)).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void explicitTimeout_exporterCompletesBeforeTimeout() throws Exception { + FastMetricExporter fastExporter = new FastMetricExporter(); + PeriodicMetricReader reader = + PeriodicMetricReader.builder(fastExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(1, TimeUnit.SECONDS) // 1 second timeout + .build(); + + reader.register(collectionRegistration); + try { + assertThat(fastExporter.waitForExport(1)).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { + FastMetricExporter fastExporter = new FastMetricExporter(); + PeriodicMetricReader reader = + PeriodicMetricReader.builder(fastExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(1, TimeUnit.SECONDS) + .setMaxExportBatchSize(2) + .build(); + + reader.register(collectionRegistration); + try { + assertThat(fastExporter.waitForExport(1)).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void explicitTimeout_zeroMeansNoTimeout() throws Exception { + SlowMetricExporter slowExporter = new SlowMetricExporter(100); + PeriodicMetricReader reader = + PeriodicMetricReader.builder(slowExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(0, TimeUnit.SECONDS) // No timeout + .build(); + + reader.register(collectionRegistration); + try { + assertThat(slowExporter.waitForExport(1)).isTrue(); + } finally { + reader.shutdown(); + } + } + + // Helper test classes for timeout testing + private static class SlowMetricExporter implements MetricExporter { + private final long delayMs; + private final AtomicInteger exportCount = new AtomicInteger(); + private final CountDownLatch exportLatch = new CountDownLatch(1); + + SlowMetricExporter(long delayMs) { + this.delayMs = delayMs; + } + + @Override + public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { + return AggregationTemporality.CUMULATIVE; + } + + @Override + public CompletableResultCode export(Collection metrics) { + try { + Thread.sleep(delayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + exportCount.incrementAndGet(); + exportLatch.countDown(); + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode flush() { + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode shutdown() { + return CompletableResultCode.ofSuccess(); + } + + boolean waitForExport(int count) throws InterruptedException { + return exportLatch.await(5, TimeUnit.SECONDS); + } + } + + private static class FastMetricExporter implements MetricExporter { + private final AtomicInteger exportCount = new AtomicInteger(); + private final CountDownLatch exportLatch = new CountDownLatch(1); + + @Override + public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { + return AggregationTemporality.CUMULATIVE; + } + + @Override + public CompletableResultCode export(Collection metrics) { + exportCount.incrementAndGet(); + exportLatch.countDown(); + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode flush() { + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode shutdown() { + return CompletableResultCode.ofSuccess(); + } + + boolean waitForExport(int count) throws InterruptedException { + return exportLatch.await(5, TimeUnit.SECONDS); + } + } + private static class WaitingMetricExporter implements MetricExporter { private final AtomicBoolean hasShutdown = new AtomicBoolean(false); From c3b8e7f9f0600cc3ea460ade14851ecda27c6b18 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Fri, 21 Aug 2026 11:23:58 +0530 Subject: [PATCH 6/9] fix: resolve JavaDoc formatting violations in PeriodicMetricReaderBuilder --- .../sdk/metrics/export/PeriodicMetricReaderBuilder.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java index c328d207b48..f3910d4dee7 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java @@ -61,8 +61,7 @@ public PeriodicMetricReaderBuilder setInterval(Duration interval) { } /** - * Sets the timeout for the underlying exporter. If unset, defaults to {@value - * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * Sets the timeout for the underlying exporter. If unset, defaults to {@value DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. */ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit unit) { requireNonNull(unit, "unit"); @@ -72,8 +71,7 @@ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit uni } /** - * Sets the timeout for the underlying exporter. If unset, defaults to {@value - * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * Sets the timeout for the underlying exporter. If unset, defaults to {@value DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. */ public PeriodicMetricReaderBuilder setExporterTimeout(Duration timeout) { requireNonNull(timeout, "timeout"); From 7d07ef888c539feb1b67f09d6f22b3dbe89d973b Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Fri, 21 Aug 2026 11:45:56 +0530 Subject: [PATCH 7/9] fix: apply Spotless formatting changes --- .../sdk/metrics/export/PeriodicMetricReaderBuilder.java | 6 ++++-- .../sdk/metrics/export/PeriodicMetricReaderTest.java | 4 +--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java index f3910d4dee7..c328d207b48 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java @@ -61,7 +61,8 @@ public PeriodicMetricReaderBuilder setInterval(Duration interval) { } /** - * Sets the timeout for the underlying exporter. If unset, defaults to {@value DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * Sets the timeout for the underlying exporter. If unset, defaults to {@value + * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. */ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit unit) { requireNonNull(unit, "unit"); @@ -71,7 +72,8 @@ public PeriodicMetricReaderBuilder setExporterTimeout(long timeout, TimeUnit uni } /** - * Sets the timeout for the underlying exporter. If unset, defaults to {@value DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. + * Sets the timeout for the underlying exporter. If unset, defaults to {@value + * DEFAULT_EXPORT_TIMEOUT_MILLIS}ms. */ public PeriodicMetricReaderBuilder setExporterTimeout(Duration timeout) { requireNonNull(timeout, "timeout"); diff --git a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java index eec6f4092f8..7c112b7c04a 100644 --- a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java +++ b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java @@ -771,9 +771,7 @@ void defaultBehavior_applies30SecondTimeout() throws Exception { // and ensuring it completes successfully SlowMetricExporter slowExporter = new SlowMetricExporter(100); // 100ms delay PeriodicMetricReader reader = - PeriodicMetricReader.builder(slowExporter) - .setInterval(Duration.ofMillis(50)) - .build(); + PeriodicMetricReader.builder(slowExporter).setInterval(Duration.ofMillis(50)).build(); reader.register(collectionRegistration); try { From bc2f19de3dd027bf681fe24aff9bce7cc1259f1d Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Fri, 21 Aug 2026 13:10:26 +0530 Subject: [PATCH 8/9] feat: implement async exporter timeout enforcement for PeriodicMetricReader Add asynchronous timeout enforcement for MetricExporter operations in PeriodicMetricReader using a dedicated timeout executor. This preserves the existing asynchronous scheduling and batching behavior while enforcing the spec-required 30-second default export timeout. Changes: - Add dedicated ScheduledExecutorService for timeout scheduling - Implement applyTimeout() method with asynchronous timeout enforcement - Preserve async batch processing with Iterator-based sequential execution - Fix Error Prone warnings (UnusedVariable, PreferJavaTimeOverload) - Add timeout enforcement test The timeout executor is shut down after final export completes to avoid RejectedExecutionException during shutdown. Timeout enforcement fails the result when timeout expires without blocking the periodic scheduler. Resolves CI compilation failures in PR #8684. --- .../metrics/export/PeriodicMetricReader.java | 101 ++++++++++++++---- .../export/PeriodicMetricReaderTest.java | 97 ++++++++++++++--- 2 files changed, 165 insertions(+), 33 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index 7cb8b0f0e48..79194c478fd 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -11,6 +11,7 @@ import io.opentelemetry.sdk.common.InternalTelemetryVersion; import io.opentelemetry.sdk.common.export.MemoryMode; import io.opentelemetry.sdk.common.internal.ComponentId; +import io.opentelemetry.sdk.common.internal.DaemonThreadFactory; import io.opentelemetry.sdk.metrics.Aggregation; import io.opentelemetry.sdk.metrics.InstrumentType; import io.opentelemetry.sdk.metrics.SdkMeterProvider; @@ -18,6 +19,8 @@ import io.opentelemetry.sdk.metrics.data.AggregationTemporality; import io.opentelemetry.sdk.metrics.data.MetricData; import java.util.Collection; +import java.util.Iterator; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -47,6 +50,7 @@ public final class PeriodicMetricReader implements MetricReader { private final long intervalNanos; private final long exporterTimeoutNanos; private final ScheduledExecutorService scheduler; + private final ScheduledExecutorService timeoutExecutor; private final Scheduled scheduled; private final Object lock = new Object(); private final InternalTelemetryVersion internalTelemetryVersion; @@ -80,6 +84,9 @@ public static PeriodicMetricReaderBuilder builder(MetricExporter exporter) { this.intervalNanos = intervalNanos; this.exporterTimeoutNanos = exporterTimeoutNanos; this.scheduler = scheduler; + this.timeoutExecutor = + Executors.newSingleThreadScheduledExecutor( + new DaemonThreadFactory("PeriodicMetricReader-timeout")); this.maxExportBatchSize = maxExportBatchSize; this.scheduled = new Scheduled(); this.internalTelemetryVersion = internalTelemetryVersion; @@ -147,6 +154,13 @@ public CompletableResultCode shutdown() { // reset the interrupted status Thread.currentThread().interrupt(); } finally { + timeoutExecutor.shutdown(); + try { + timeoutExecutor.awaitTermination(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + timeoutExecutor.shutdownNow(); + Thread.currentThread().interrupt(); + } CompletableResultCode shutdownResult = scheduled.shutdown(); shutdownResult.whenComplete( () -> { @@ -216,32 +230,81 @@ private Scheduled() {} private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { CompletableResultCode result = exporter.export(metricData); - if (exporterTimeoutNanos != Long.MAX_VALUE) { - result.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); - } - return result; + return applyTimeout(result); } Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); CompletableResultCode sequentialResult = new CompletableResultCode(); - boolean anyFailed = false; - for (Collection currentBatch : batches) { - CompletableResultCode currentResult = exporter.export(currentBatch); - if (exporterTimeoutNanos != Long.MAX_VALUE) { - currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); - } - if (!currentResult.isSuccess()) { - anyFailed = true; - } - } - if (anyFailed) { - sequentialResult.fail(); - } else { - sequentialResult.succeed(); - } + AtomicBoolean anyFailed = new AtomicBoolean(false); + Iterator> batchIterator = batches.iterator(); + Runnable exportNext = + new Runnable() { + @Override + public void run() { + while (batchIterator.hasNext()) { + Collection currentBatch = batchIterator.next(); + CompletableResultCode currentResult = exporter.export(currentBatch); + CompletableResultCode timeoutResult = applyTimeout(currentResult); + if (timeoutResult.isDone()) { + if (!timeoutResult.isSuccess()) { + anyFailed.set(true); + } + } else { + timeoutResult.whenComplete( + () -> { + if (!timeoutResult.isSuccess()) { + anyFailed.set(true); + } + this.run(); + }); + return; + } + } + if (anyFailed.get()) { + sequentialResult.fail(); + } else { + sequentialResult.succeed(); + } + } + }; + exportNext.run(); return sequentialResult; } + private CompletableResultCode applyTimeout(CompletableResultCode result) { + if (exporterTimeoutNanos == Long.MAX_VALUE) { + return result; + } + CompletableResultCode timeoutResult = new CompletableResultCode(); + AtomicBoolean timedOut = new AtomicBoolean(false); + ScheduledFuture timeoutFuture = + timeoutExecutor.schedule( + () -> { + if (!result.isDone()) { + timedOut.set(true); + logger.log( + Level.WARNING, "Export timed out after " + exporterTimeoutNanos + "ns"); + timeoutResult.fail(); + } + }, + exporterTimeoutNanos, + TimeUnit.NANOSECONDS); + result.whenComplete( + () -> { + if (timeoutFuture != null) { + timeoutFuture.cancel(false); + } + if (!timedOut.get()) { + if (result.isSuccess()) { + timeoutResult.succeed(); + } else { + timeoutResult.fail(); + } + } + }); + return timeoutResult; + } + void setMeterProvider(MeterProvider meterProvider) { instrumentation = new MetricReaderInstrumentation(COMPONENT_ID, meterProvider); } diff --git a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java index 7c112b7c04a..a7b2dd3e6d4 100644 --- a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java +++ b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java @@ -720,21 +720,23 @@ void stringRepresentation() { @Test void setExporterTimeout_configuresTimeout() { PeriodicMetricReaderBuilder builder = - PeriodicMetricReader.builder(metricExporter).setExporterTimeout(5, TimeUnit.SECONDS); + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ofSeconds(5)); assertThat(builder).isNotNull(); } @Test - void setExporterTimeout_withDuration_configuresTimeout() { + @SuppressWarnings("PreferJavaTimeOverload") + void setExporterTimeout_withTimeUnit_overload() { + // Keep at least one test for the long + TimeUnit public API to verify that overload works PeriodicMetricReaderBuilder builder = - PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ofSeconds(5)); + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(5, TimeUnit.SECONDS); assertThat(builder).isNotNull(); } @Test void setExporterTimeout_zeroMeansNoTimeout() { PeriodicMetricReaderBuilder builder = - PeriodicMetricReader.builder(metricExporter).setExporterTimeout(0, TimeUnit.SECONDS); + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ZERO); assertThat(builder).isNotNull(); } @@ -743,7 +745,7 @@ void setExporterTimeout_negativeThrowsException() { assertThatThrownBy( () -> PeriodicMetricReader.builder(metricExporter) - .setExporterTimeout(-1, TimeUnit.SECONDS)) + .setExporterTimeout(Duration.ofSeconds(-1))) .isInstanceOf(IllegalArgumentException.class) .hasMessage("timeout must be non-negative"); } @@ -776,7 +778,7 @@ void defaultBehavior_applies30SecondTimeout() throws Exception { reader.register(collectionRegistration); try { // Wait for an export to complete - should succeed since 100ms < 30s - assertThat(slowExporter.waitForExport(1)).isTrue(); + assertThat(slowExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -788,12 +790,12 @@ void explicitTimeout_exporterCompletesBeforeTimeout() throws Exception { PeriodicMetricReader reader = PeriodicMetricReader.builder(fastExporter) .setInterval(Duration.ofMillis(50)) - .setExporterTimeout(1, TimeUnit.SECONDS) // 1 second timeout + .setExporterTimeout(Duration.ofSeconds(1)) // 1 second timeout .build(); reader.register(collectionRegistration); try { - assertThat(fastExporter.waitForExport(1)).isTrue(); + assertThat(fastExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -805,13 +807,13 @@ void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { PeriodicMetricReader reader = PeriodicMetricReader.builder(fastExporter) .setInterval(Duration.ofMillis(50)) - .setExporterTimeout(1, TimeUnit.SECONDS) + .setExporterTimeout(Duration.ofSeconds(1)) .setMaxExportBatchSize(2) .build(); reader.register(collectionRegistration); try { - assertThat(fastExporter.waitForExport(1)).isTrue(); + assertThat(fastExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -823,18 +825,85 @@ void explicitTimeout_zeroMeansNoTimeout() throws Exception { PeriodicMetricReader reader = PeriodicMetricReader.builder(slowExporter) .setInterval(Duration.ofMillis(50)) - .setExporterTimeout(0, TimeUnit.SECONDS) // No timeout + .setExporterTimeout(Duration.ZERO) // No timeout .build(); reader.register(collectionRegistration); try { - assertThat(slowExporter.waitForExport(1)).isTrue(); + assertThat(slowExporter.waitForExport()).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void timeoutEnforcement_failsSlowExporter() throws Exception { + // Exporter that completes asynchronously after a delay longer than timeout + AsyncSlowMetricExporter slowExporter = new AsyncSlowMetricExporter(100); // 100ms delay + PeriodicMetricReader reader = + PeriodicMetricReader.builder(slowExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(Duration.ofMillis(10)) // 10ms timeout + .build(); + + reader.register(collectionRegistration); + try { + // Wait for export to be attempted + assertThat(slowExporter.exportStarted.await(5, TimeUnit.SECONDS)).isTrue(); + // Wait for timeout to occur + Thread.sleep(100); + // Export should have been attempted but timed out + assertThat(slowExporter.exportCount.get()).isGreaterThan(0); } finally { reader.shutdown(); } } // Helper test classes for timeout testing + private static class AsyncSlowMetricExporter implements MetricExporter { + private final long delayMs; + private final AtomicInteger exportCount = new AtomicInteger(); + private final CountDownLatch exportStarted = new CountDownLatch(1); + + AsyncSlowMetricExporter(long delayMs) { + this.delayMs = delayMs; + } + + @Override + public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { + return AggregationTemporality.CUMULATIVE; + } + + @Override + public CompletableResultCode export(Collection metrics) { + exportCount.incrementAndGet(); + exportStarted.countDown(); + CompletableResultCode result = new CompletableResultCode(); + new Thread( + () -> { + try { + Thread.sleep(delayMs); + result.succeed(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + result.fail(); + } + }) + .start(); + return result; + } + + @Override + public CompletableResultCode flush() { + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode shutdown() { + return CompletableResultCode.ofSuccess(); + } + } + private static class SlowMetricExporter implements MetricExporter { private final long delayMs; private final AtomicInteger exportCount = new AtomicInteger(); @@ -871,7 +940,7 @@ public CompletableResultCode shutdown() { return CompletableResultCode.ofSuccess(); } - boolean waitForExport(int count) throws InterruptedException { + boolean waitForExport() throws InterruptedException { return exportLatch.await(5, TimeUnit.SECONDS); } } @@ -902,7 +971,7 @@ public CompletableResultCode shutdown() { return CompletableResultCode.ofSuccess(); } - boolean waitForExport(int count) throws InterruptedException { + boolean waitForExport() throws InterruptedException { return exportLatch.await(5, TimeUnit.SECONDS); } } From b25a0bbb11c6f4795df05df5399d0aaa38ab4658 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sat, 22 Aug 2026 00:20:48 +0530 Subject: [PATCH 9/9] refactor: use 2-thread scheduler and simplify timeout enforcement per jack-berg review - Remove separate timeout executor and reuse existing scheduler - Increase scheduler from 1 to 2 threads to handle both periodic exports and timeout tasks - Simplify applyTimeout() by removing AtomicBoolean/timedOut state tracking - Combine SlowMetricExporter and FastMetricExporter into single DelayingMetricExporter - Add RejectedExecutionException handling for shutdown race condition - Fix BooleanParameter warnings in tests --- .../metrics/export/PeriodicMetricReader.java | 53 +++---- .../export/PeriodicMetricReaderBuilder.java | 2 +- .../export/PeriodicMetricReaderTest.java | 143 ++++++------------ 3 files changed, 68 insertions(+), 130 deletions(-) diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java index 79194c478fd..d44a95226e1 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReader.java @@ -11,7 +11,6 @@ import io.opentelemetry.sdk.common.InternalTelemetryVersion; import io.opentelemetry.sdk.common.export.MemoryMode; import io.opentelemetry.sdk.common.internal.ComponentId; -import io.opentelemetry.sdk.common.internal.DaemonThreadFactory; import io.opentelemetry.sdk.metrics.Aggregation; import io.opentelemetry.sdk.metrics.InstrumentType; import io.opentelemetry.sdk.metrics.SdkMeterProvider; @@ -20,7 +19,7 @@ import io.opentelemetry.sdk.metrics.data.MetricData; import java.util.Collection; import java.util.Iterator; -import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -50,7 +49,6 @@ public final class PeriodicMetricReader implements MetricReader { private final long intervalNanos; private final long exporterTimeoutNanos; private final ScheduledExecutorService scheduler; - private final ScheduledExecutorService timeoutExecutor; private final Scheduled scheduled; private final Object lock = new Object(); private final InternalTelemetryVersion internalTelemetryVersion; @@ -84,9 +82,6 @@ public static PeriodicMetricReaderBuilder builder(MetricExporter exporter) { this.intervalNanos = intervalNanos; this.exporterTimeoutNanos = exporterTimeoutNanos; this.scheduler = scheduler; - this.timeoutExecutor = - Executors.newSingleThreadScheduledExecutor( - new DaemonThreadFactory("PeriodicMetricReader-timeout")); this.maxExportBatchSize = maxExportBatchSize; this.scheduled = new Scheduled(); this.internalTelemetryVersion = internalTelemetryVersion; @@ -154,13 +149,6 @@ public CompletableResultCode shutdown() { // reset the interrupted status Thread.currentThread().interrupt(); } finally { - timeoutExecutor.shutdown(); - try { - timeoutExecutor.awaitTermination(5, TimeUnit.SECONDS); - } catch (InterruptedException e) { - timeoutExecutor.shutdownNow(); - Thread.currentThread().interrupt(); - } CompletableResultCode shutdownResult = scheduled.shutdown(); shutdownResult.whenComplete( () -> { @@ -275,34 +263,35 @@ private CompletableResultCode applyTimeout(CompletableResultCode result) { if (exporterTimeoutNanos == Long.MAX_VALUE) { return result; } - CompletableResultCode timeoutResult = new CompletableResultCode(); - AtomicBoolean timedOut = new AtomicBoolean(false); - ScheduledFuture timeoutFuture = - timeoutExecutor.schedule( - () -> { - if (!result.isDone()) { - timedOut.set(true); + + try { + CompletableResultCode timeoutResult = new CompletableResultCode(); + + ScheduledFuture timeoutFuture = + scheduler.schedule( + () -> { logger.log( Level.WARNING, "Export timed out after " + exporterTimeoutNanos + "ns"); timeoutResult.fail(); - } - }, - exporterTimeoutNanos, - TimeUnit.NANOSECONDS); - result.whenComplete( - () -> { - if (timeoutFuture != null) { + }, + exporterTimeoutNanos, + TimeUnit.NANOSECONDS); + + result.whenComplete( + () -> { timeoutFuture.cancel(false); - } - if (!timedOut.get()) { if (result.isSuccess()) { timeoutResult.succeed(); } else { timeoutResult.fail(); } - } - }); - return timeoutResult; + }); + + return timeoutResult; + } catch (RejectedExecutionException e) { + // Scheduler is shutting down, return original result without timeout enforcement + return result; + } } void setMeterProvider(MeterProvider meterProvider) { diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java index c328d207b48..5d19084f75f 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderBuilder.java @@ -106,7 +106,7 @@ public PeriodicMetricReader build() { ScheduledExecutorService executor = this.executor; if (executor == null) { executor = - Executors.newScheduledThreadPool(1, new DaemonThreadFactory("PeriodicMetricReader")); + Executors.newScheduledThreadPool(2, new DaemonThreadFactory("PeriodicMetricReader")); } return new PeriodicMetricReader( metricExporter, diff --git a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java index a7b2dd3e6d4..64797cbaf8e 100644 --- a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java +++ b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/export/PeriodicMetricReaderTest.java @@ -771,14 +771,14 @@ void defaultBehavior_applies30SecondTimeout() throws Exception { // This test verifies that by default, a 30-second timeout is applied // We test this by having an exporter that takes less than 30 seconds // and ensuring it completes successfully - SlowMetricExporter slowExporter = new SlowMetricExporter(100); // 100ms delay + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(100); // 100ms delay PeriodicMetricReader reader = - PeriodicMetricReader.builder(slowExporter).setInterval(Duration.ofMillis(50)).build(); + PeriodicMetricReader.builder(delayingExporter).setInterval(Duration.ofMillis(50)).build(); reader.register(collectionRegistration); try { // Wait for an export to complete - should succeed since 100ms < 30s - assertThat(slowExporter.waitForExport()).isTrue(); + assertThat(delayingExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -786,16 +786,16 @@ void defaultBehavior_applies30SecondTimeout() throws Exception { @Test void explicitTimeout_exporterCompletesBeforeTimeout() throws Exception { - FastMetricExporter fastExporter = new FastMetricExporter(); + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(0); // 0ms = fast PeriodicMetricReader reader = - PeriodicMetricReader.builder(fastExporter) + PeriodicMetricReader.builder(delayingExporter) .setInterval(Duration.ofMillis(50)) .setExporterTimeout(Duration.ofSeconds(1)) // 1 second timeout .build(); reader.register(collectionRegistration); try { - assertThat(fastExporter.waitForExport()).isTrue(); + assertThat(delayingExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -803,9 +803,9 @@ void explicitTimeout_exporterCompletesBeforeTimeout() throws Exception { @Test void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { - FastMetricExporter fastExporter = new FastMetricExporter(); + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(0); // 0ms = fast PeriodicMetricReader reader = - PeriodicMetricReader.builder(fastExporter) + PeriodicMetricReader.builder(delayingExporter) .setInterval(Duration.ofMillis(50)) .setExporterTimeout(Duration.ofSeconds(1)) .setMaxExportBatchSize(2) @@ -813,7 +813,7 @@ void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { reader.register(collectionRegistration); try { - assertThat(fastExporter.waitForExport()).isTrue(); + assertThat(delayingExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -821,16 +821,16 @@ void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { @Test void explicitTimeout_zeroMeansNoTimeout() throws Exception { - SlowMetricExporter slowExporter = new SlowMetricExporter(100); + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(100); PeriodicMetricReader reader = - PeriodicMetricReader.builder(slowExporter) + PeriodicMetricReader.builder(delayingExporter) .setInterval(Duration.ofMillis(50)) .setExporterTimeout(Duration.ZERO) // No timeout .build(); reader.register(collectionRegistration); try { - assertThat(slowExporter.waitForExport()).isTrue(); + assertThat(delayingExporter.waitForExport()).isTrue(); } finally { reader.shutdown(); } @@ -839,7 +839,8 @@ void explicitTimeout_zeroMeansNoTimeout() throws Exception { @Test void timeoutEnforcement_failsSlowExporter() throws Exception { // Exporter that completes asynchronously after a delay longer than timeout - AsyncSlowMetricExporter slowExporter = new AsyncSlowMetricExporter(100); // 100ms delay + DelayingMetricExporter slowExporter = + new DelayingMetricExporter(100, /* async= */ true); // 100ms async delay PeriodicMetricReader reader = PeriodicMetricReader.builder(slowExporter) .setInterval(Duration.ofMillis(50)) @@ -849,7 +850,7 @@ void timeoutEnforcement_failsSlowExporter() throws Exception { reader.register(collectionRegistration); try { // Wait for export to be attempted - assertThat(slowExporter.exportStarted.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(slowExporter.waitForExport()).isTrue(); // Wait for timeout to occur Thread.sleep(100); // Export should have been attempted but timed out @@ -860,94 +861,20 @@ void timeoutEnforcement_failsSlowExporter() throws Exception { } // Helper test classes for timeout testing - private static class AsyncSlowMetricExporter implements MetricExporter { - private final long delayMs; - private final AtomicInteger exportCount = new AtomicInteger(); - private final CountDownLatch exportStarted = new CountDownLatch(1); - - AsyncSlowMetricExporter(long delayMs) { - this.delayMs = delayMs; - } - - @Override - public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { - return AggregationTemporality.CUMULATIVE; - } - - @Override - public CompletableResultCode export(Collection metrics) { - exportCount.incrementAndGet(); - exportStarted.countDown(); - CompletableResultCode result = new CompletableResultCode(); - new Thread( - () -> { - try { - Thread.sleep(delayMs); - result.succeed(); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - result.fail(); - } - }) - .start(); - return result; - } - - @Override - public CompletableResultCode flush() { - return CompletableResultCode.ofSuccess(); - } - - @Override - public CompletableResultCode shutdown() { - return CompletableResultCode.ofSuccess(); - } - } - - private static class SlowMetricExporter implements MetricExporter { + private static class DelayingMetricExporter implements MetricExporter { private final long delayMs; private final AtomicInteger exportCount = new AtomicInteger(); private final CountDownLatch exportLatch = new CountDownLatch(1); + private final boolean async; - SlowMetricExporter(long delayMs) { - this.delayMs = delayMs; + DelayingMetricExporter(long delayMs) { + this(delayMs, /* async= */ false); } - @Override - public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { - return AggregationTemporality.CUMULATIVE; - } - - @Override - public CompletableResultCode export(Collection metrics) { - try { - Thread.sleep(delayMs); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - exportCount.incrementAndGet(); - exportLatch.countDown(); - return CompletableResultCode.ofSuccess(); - } - - @Override - public CompletableResultCode flush() { - return CompletableResultCode.ofSuccess(); - } - - @Override - public CompletableResultCode shutdown() { - return CompletableResultCode.ofSuccess(); - } - - boolean waitForExport() throws InterruptedException { - return exportLatch.await(5, TimeUnit.SECONDS); + DelayingMetricExporter(long delayMs, boolean async) { + this.delayMs = delayMs; + this.async = async; } - } - - private static class FastMetricExporter implements MetricExporter { - private final AtomicInteger exportCount = new AtomicInteger(); - private final CountDownLatch exportLatch = new CountDownLatch(1); @Override public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { @@ -957,8 +884,30 @@ public AggregationTemporality getAggregationTemporality(InstrumentType instrumen @Override public CompletableResultCode export(Collection metrics) { exportCount.incrementAndGet(); - exportLatch.countDown(); - return CompletableResultCode.ofSuccess(); + if (async) { + CompletableResultCode result = new CompletableResultCode(); + new Thread( + () -> { + try { + Thread.sleep(delayMs); + result.succeed(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + result.fail(); + } + }) + .start(); + exportLatch.countDown(); + return result; + } else { + try { + Thread.sleep(delayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + exportLatch.countDown(); + return CompletableResultCode.ofSuccess(); + } } @Override