From 64f3e8b6615a1fb5e297213690350850f3e0bc6a Mon Sep 17 00:00:00 2001 From: brunobat Date: Fri, 28 Aug 2026 11:00:17 +0100 Subject: [PATCH] Add timeout shutdown config to PeriodicMetricReader and its builder --- CHANGELOG.md | 8 ++ .../opentelemetry-sdk-metrics.txt | 2 + .../metrics/export/PeriodicMetricReader.java | 11 ++- .../export/PeriodicMetricReaderBuilder.java | 34 ++++++++- .../export/PeriodicMetricReaderTest.java | 73 +++++++++++++++++++ 5 files changed, 123 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0e1687e75c4..5c854746485 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,14 @@ ## Unreleased +### SDK + +#### Metrics + +* Add `PeriodicMetricReaderBuilder.setShutdownTimeout` to configure the previously hardcoded + shutdown wait time. Default is kept as 5s. + ([#8756](https://github.com/open-telemetry/opentelemetry-java/pull/8756)) + ## Version 1.65.0 (2026-08-07) **NOTE:** The `opentelemetry-exporter-zipkin` artifact has stopped being published. It was diff --git a/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt b/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt index ea360d6fad9..69eebddd0fe 100644 --- a/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt +++ b/docs/apidiffs/current_vs_latest/opentelemetry-sdk-metrics.txt @@ -2,3 +2,5 @@ Comparing source compatibility of opentelemetry-sdk-metrics-1.66.0-SNAPSHOT.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 setInternalTelemetryVersion(io.opentelemetry.sdk.common.InternalTelemetryVersion) + +++ NEW METHOD: PUBLIC(+) io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder setShutdownTimeout(long, java.util.concurrent.TimeUnit) + +++ NEW METHOD: PUBLIC(+) io.opentelemetry.sdk.metrics.export.PeriodicMetricReaderBuilder setShutdownTimeout(java.time.Duration) 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..6ddb7eba423 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 @@ -48,6 +48,7 @@ public final class PeriodicMetricReader implements MetricReader { private final long intervalNanos; private final ScheduledExecutorService scheduler; private final Scheduled scheduled; + private final long shutdownTimeoutMillis; private final Object lock = new Object(); private final InternalTelemetryVersion internalTelemetryVersion; @@ -74,13 +75,15 @@ public static PeriodicMetricReaderBuilder builder(MetricExporter exporter) { long intervalNanos, ScheduledExecutorService scheduler, int maxExportBatchSize, - InternalTelemetryVersion internalTelemetryVersion) { + InternalTelemetryVersion internalTelemetryVersion, + long shutdownTimeoutMillis) { this.exporter = exporter; this.intervalNanos = intervalNanos; this.scheduler = scheduler; this.maxExportBatchSize = maxExportBatchSize; this.scheduled = new Scheduled(); this.internalTelemetryVersion = internalTelemetryVersion; + this.shutdownTimeoutMillis = shutdownTimeoutMillis; } @Override @@ -133,12 +136,12 @@ public CompletableResultCode shutdown() { } scheduler.shutdown(); try { - scheduler.awaitTermination(5, TimeUnit.SECONDS); + scheduler.awaitTermination(shutdownTimeoutMillis, TimeUnit.MILLISECONDS); // Wait for any in-flight export to complete before performing the final collection. // Without this, doRun() sees exportAvailable=false and drops the final metrics. - scheduled.flushInProgress.join(5, TimeUnit.SECONDS); + scheduled.flushInProgress.join(shutdownTimeoutMillis, TimeUnit.MILLISECONDS); CompletableResultCode flushResult = scheduled.doRun(); - flushResult.join(5, TimeUnit.SECONDS); + flushResult.join(shutdownTimeoutMillis, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { // force a shutdown if the export hasn't finished. scheduler.shutdownNow(); 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 2529efb4430..eec0ed06ddb 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 long DEFAULT_SHUTDOWN_TIMEOUT_MILLIS = 5000; private final MetricExporter metricExporter; @@ -35,6 +36,8 @@ public final class PeriodicMetricReaderBuilder { private int maxExportBatchSize; + private long shutdownTimeoutMillis = DEFAULT_SHUTDOWN_TIMEOUT_MILLIS; + PeriodicMetricReaderBuilder(MetricExporter metricExporter) { this.metricExporter = metricExporter; } @@ -78,6 +81,30 @@ PeriodicMetricReaderBuilder setMaxExportBatchSize(int maxExportBatchSize) { return this; } + /** + * Sets the maximum time to wait for shutdown to complete. This timeout is applied independently + * to each phase of the shutdown sequence (awaiting executor termination, joining any in-flight + * export, and joining the final export), so shutdown may take up to three times this value. If + * unset, defaults to {@value DEFAULT_SHUTDOWN_TIMEOUT_MILLIS}ms per phase. + */ + public PeriodicMetricReaderBuilder setShutdownTimeout(long timeout, TimeUnit unit) { + requireNonNull(unit, "unit"); + checkArgument(timeout > 0, "timeout must be positive"); + shutdownTimeoutMillis = unit.toMillis(timeout); + return this; + } + + /** + * Sets the maximum time to wait for shutdown to complete. This timeout is applied independently + * to each phase of the shutdown sequence (awaiting executor termination, joining any in-flight + * export, and joining the final export), so shutdown may take up to three times this value. If + * unset, defaults to {@value DEFAULT_SHUTDOWN_TIMEOUT_MILLIS}ms per phase. + */ + public PeriodicMetricReaderBuilder setShutdownTimeout(Duration timeout) { + requireNonNull(timeout, "timeout"); + return setShutdownTimeout(timeout.toMillis(), TimeUnit.MILLISECONDS); + } + /** Build a {@link PeriodicMetricReader} with the configuration of this builder. */ public PeriodicMetricReader build() { ScheduledExecutorService executor = this.executor; @@ -86,7 +113,12 @@ public PeriodicMetricReader build() { Executors.newScheduledThreadPool(1, new DaemonThreadFactory("PeriodicMetricReader")); } return new PeriodicMetricReader( - metricExporter, intervalNanos, executor, maxExportBatchSize, internalTelemetryVersion); + metricExporter, + intervalNanos, + executor, + maxExportBatchSize, + internalTelemetryVersion, + shutdownTimeoutMillis); } /** 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..536cda88164 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 @@ -470,6 +470,79 @@ void invalidConfig() { assertThatThrownBy(() -> PeriodicMetricReader.builder(metricExporter).setExecutor(null)) .isInstanceOf(NullPointerException.class) .hasMessage("executor"); + assertThatThrownBy( + () -> PeriodicMetricReader.builder(metricExporter).setShutdownTimeout(1, null)) + .isInstanceOf(NullPointerException.class) + .hasMessage("unit"); + assertThatThrownBy( + () -> + PeriodicMetricReader.builder(metricExporter) + .setShutdownTimeout(-1, TimeUnit.MILLISECONDS)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("timeout must be positive"); + assertThatThrownBy(() -> PeriodicMetricReader.builder(metricExporter).setShutdownTimeout(null)) + .isInstanceOf(NullPointerException.class) + .hasMessage("timeout"); + } + + @Test + @Timeout(10) + @SuppressLogger(PeriodicMetricReader.class) + void shutdown_respectsConfiguredWaitTime() throws Exception { + CompletableResultCode neverCompletes = new CompletableResultCode(); + CountDownLatch exportStarted = new CountDownLatch(1); + + MetricExporter blockingExporter = + new MetricExporter() { + @Override + public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { + return AggregationTemporality.CUMULATIVE; + } + + @Override + public CompletableResultCode export(Collection metrics) { + exportStarted.countDown(); + return neverCompletes; + } + + @Override + public CompletableResultCode flush() { + return CompletableResultCode.ofSuccess(); + } + + @Override + public CompletableResultCode shutdown() { + return CompletableResultCode.ofSuccess(); + } + }; + + PeriodicMetricReader reader = + PeriodicMetricReader.builder(blockingExporter) + .setInterval(Duration.ofSeconds(Integer.MAX_VALUE)) + .setShutdownTimeout(Duration.ofMillis(100)) + .build(); + reader.register(collectionRegistration); + + // Start an export that never completes so shutdown() has to wait on it. + reader.forceFlush(); + assertThat(exportStarted.await(5, TimeUnit.SECONDS)).isTrue(); + + // shutdown() must give up after the configured 100ms per phase, well before the + // 5s-per-phase default would elapse. + CountDownLatch shutdownDone = new CountDownLatch(1); + Thread shutdownThread = + new Thread( + () -> { + reader.shutdown(); + shutdownDone.countDown(); + }); + shutdownThread.setDaemon(true); + shutdownThread.start(); + + assertThat(shutdownDone.await(3, TimeUnit.SECONDS)).isTrue(); + + // Release the hanging export so the test thread does not leak a pending result. + neverCompletes.succeed(); } @Test