From 41910f3f7edc473e190a935590364492dc254c61 Mon Sep 17 00:00:00 2001 From: Rajkaran Yadav Date: Sun, 2 Aug 2026 17:47:07 +0530 Subject: [PATCH 1/4] 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/4] 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/4] 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/4] 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();