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 14e3751ba18..ace3d0a83cf 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 @@ -19,6 +19,7 @@ import io.opentelemetry.sdk.metrics.data.MetricData; import java.util.Collection; import java.util.Iterator; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -46,6 +47,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 +74,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,7 +217,8 @@ private Scheduled() {} private CompletableResultCode exportMetrics(Collection metricData) { if (maxExportBatchSize == 0) { - return exporter.export(metricData); + CompletableResultCode result = exporter.export(metricData); + return applyTimeout(result); } Collection> batches = MetricExportBatcher.batchMetrics(metricData, maxExportBatchSize); @@ -227,14 +232,15 @@ public void run() { while (batchIterator.hasNext()) { Collection currentBatch = batchIterator.next(); CompletableResultCode currentResult = exporter.export(currentBatch); - if (currentResult.isDone()) { - if (!currentResult.isSuccess()) { + CompletableResultCode timeoutResult = applyTimeout(currentResult); + if (timeoutResult.isDone()) { + if (!timeoutResult.isSuccess()) { anyFailed.set(true); } } else { - currentResult.whenComplete( + timeoutResult.whenComplete( () -> { - if (!currentResult.isSuccess()) { + if (!timeoutResult.isSuccess()) { anyFailed.set(true); } this.run(); @@ -253,6 +259,45 @@ public void run() { return sequentialResult; } + private CompletableResultCode applyTimeout(CompletableResultCode result) { + if (exporterTimeoutNanos == Long.MAX_VALUE) { + return result; + } + + if (result.isDone()) { + return result; + } + + 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( + () -> { + timeoutFuture.cancel(false); + if (result.isSuccess()) { + timeoutResult.succeed(); + } else { + timeoutResult.fail(); + } + }); + + return timeoutResult; + } catch (RejectedExecutionException e) { + // Scheduler is shutting down, return original result without timeout enforcement + return result; + } + } + void setMeterProvider(MeterProvider meterProvider) { instrumentation = new MetricReaderInstrumentation(COMPONENT_ID, 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 0a2ef9ea900..4bcf0ba6892 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); + @Nullable private Long exporterTimeoutNanos; + @Nullable private ScheduledExecutorService executor; private int maxExportBatchSize; @@ -57,6 +60,26 @@ 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. + */ + 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. + */ + 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"); @@ -83,10 +106,17 @@ 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, intervalNanos, executor, maxExportBatchSize, internalTelemetryVersion); + metricExporter, + intervalNanos, + exporterTimeoutNanos != null + ? exporterTimeoutNanos + : Math.min(intervalNanos, TimeUnit.MILLISECONDS.toNanos(DEFAULT_EXPORT_TIMEOUT_MILLIS)), + executor, + maxExportBatchSize, + internalTelemetryVersion); } /** Sets the internal telemetry version used to control self-observability metrics. */ 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..40cd0adcbeb 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 @@ -10,6 +10,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -717,6 +718,246 @@ void stringRepresentation() { + "}"); } + @Test + void setExporterTimeout_configuresTimeout() { + PeriodicMetricReaderBuilder builder = + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ofSeconds(5)); + assertThat(builder).isNotNull(); + } + + @Test + @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(5, TimeUnit.SECONDS); + assertThat(builder).isNotNull(); + } + + @Test + void setExporterTimeout_zeroMeansNoTimeout() { + PeriodicMetricReaderBuilder builder = + PeriodicMetricReader.builder(metricExporter).setExporterTimeout(Duration.ZERO); + assertThat(builder).isNotNull(); + } + + @Test + void setExporterTimeout_negativeThrowsException() { + assertThatThrownBy( + () -> + PeriodicMetricReader.builder(metricExporter) + .setExporterTimeout(Duration.ofSeconds(-1))) + .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 + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(100); // 100ms delay + PeriodicMetricReader reader = + PeriodicMetricReader.builder(delayingExporter).setInterval(Duration.ofMillis(50)).build(); + + reader.register(collectionRegistration); + try { + // Wait for an export to complete - should succeed since 100ms < 30s + assertThat(delayingExporter.waitForExport()).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void explicitTimeout_exporterCompletesBeforeTimeout() throws Exception { + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(0); // 0ms = fast + PeriodicMetricReader reader = + PeriodicMetricReader.builder(delayingExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(Duration.ofSeconds(1)) // 1 second timeout + .build(); + + reader.register(collectionRegistration); + try { + assertThat(delayingExporter.waitForExport()).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + @SuppressWarnings({"rawtypes", "unchecked"}) + void export_alreadyCompleted_doesNotScheduleTimeout() { + ScheduledExecutorService mockScheduler = mock(ScheduledExecutorService.class); + ScheduledFuture mockFuture = mock(ScheduledFuture.class); + when(mockScheduler.scheduleAtFixedRate(any(), anyLong(), anyLong(), any())) + .thenReturn(mockFuture); + + MetricExporter fastExporter = mock(MetricExporter.class); + when(fastExporter.export(any())).thenReturn(CompletableResultCode.ofSuccess()); + when(fastExporter.flush()).thenReturn(CompletableResultCode.ofSuccess()); + when(fastExporter.shutdown()).thenReturn(CompletableResultCode.ofSuccess()); + when(fastExporter.getAggregationTemporality(any())) + .thenReturn(AggregationTemporality.CUMULATIVE); + + PeriodicMetricReader reader = + PeriodicMetricReader.builder(fastExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(Duration.ofSeconds(1)) + .setExecutor(mockScheduler) + .build(); + + reader.register(collectionRegistration); + + reader.forceFlush(); + + // Verify scheduleAtFixedRate was called (for periodic export) + verify(mockScheduler, times(1)).scheduleAtFixedRate(any(), anyLong(), anyLong(), any()); + // Verify schedule (for timeout) was never called because the export completed synchronously + verify(mockScheduler, never()).schedule(any(Runnable.class), anyLong(), any()); + } + + @Test + void explicitTimeout_withBatching_completesBeforeTimeout() throws Exception { + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(0); // 0ms = fast + PeriodicMetricReader reader = + PeriodicMetricReader.builder(delayingExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(Duration.ofSeconds(1)) + .setMaxExportBatchSize(2) + .build(); + + reader.register(collectionRegistration); + try { + assertThat(delayingExporter.waitForExport()).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void explicitTimeout_zeroMeansNoTimeout() throws Exception { + DelayingMetricExporter delayingExporter = new DelayingMetricExporter(100); + PeriodicMetricReader reader = + PeriodicMetricReader.builder(delayingExporter) + .setInterval(Duration.ofMillis(50)) + .setExporterTimeout(Duration.ZERO) // No timeout + .build(); + + reader.register(collectionRegistration); + try { + assertThat(delayingExporter.waitForExport()).isTrue(); + } finally { + reader.shutdown(); + } + } + + @Test + void timeoutEnforcement_failsSlowExporter() throws Exception { + // Exporter that completes asynchronously after a delay longer than timeout + DelayingMetricExporter slowExporter = + new DelayingMetricExporter(100, /* async= */ true); // 100ms async 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.waitForExport()).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 DelayingMetricExporter implements MetricExporter { + private final long delayMs; + private final AtomicInteger exportCount = new AtomicInteger(); + private final CountDownLatch exportLatch = new CountDownLatch(1); + private final boolean async; + + DelayingMetricExporter(long delayMs) { + this(delayMs, /* async= */ false); + } + + DelayingMetricExporter(long delayMs, boolean async) { + this.delayMs = delayMs; + this.async = async; + } + + @Override + public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { + return AggregationTemporality.CUMULATIVE; + } + + @Override + public CompletableResultCode export(Collection metrics) { + exportCount.incrementAndGet(); + 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 + public CompletableResultCode flush() { + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode shutdown() { + return CompletableResultCode.ofSuccess(); + } + + boolean waitForExport() throws InterruptedException { + return exportLatch.await(5, TimeUnit.SECONDS); + } + } + private static class WaitingMetricExporter implements MetricExporter { private final AtomicBoolean hasShutdown = new AtomicBoolean(false);