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..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; @@ -46,6 +45,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 +72,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 +215,26 @@ 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; } 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(); + boolean anyFailed = false; + for (Collection currentBatch : batches) { + CompletableResultCode currentResult = exporter.export(currentBatch); + currentResult.join(exporterTimeoutNanos, TimeUnit.NANOSECONDS); + if (!currentResult.isSuccess()) { + anyFailed = true; + } + } + if (anyFailed) { + 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. */