From d5c7d5b89fcd770e49d3d7c17b59ee3d3c6189f5 Mon Sep 17 00:00:00 2001 From: Diego Hurtado Date: Wed, 22 Jul 2026 01:08:19 -0500 Subject: [PATCH 1/2] otlp exporter: retry on 429, honor Retry-After, and cap backoff The OTLP HTTP exporters (trace, metric, logs) now treat HTTP 429 (Too Many Requests) as retryable and parse the Retry-After response header, honoring both the delay-seconds and HTTP-date forms as the delay before the next retry, as required by the OTLP specification. The exponential backoff in the HTTP exporters, as well as the gRPC exporter, is now clamped to a maximum interval so it cannot grow without bound. --- .changelog/0000.fixed | 5 + .../exporter/otlp/proto/grpc/exporter.py | 7 +- .../otlp/proto/http/_common/__init__.py | 43 ++++++ .../otlp/proto/http/_log_exporter/__init__.py | 10 +- .../proto/http/metric_exporter/__init__.py | 10 +- .../proto/http/trace_exporter/__init__.py | 10 +- .../metrics/test_otlp_metrics_exporter.py | 35 +++++ .../tests/test_proto_log_exporter.py | 35 +++++ .../tests/test_proto_span_exporter.py | 125 ++++++++++++++++++ 9 files changed, 276 insertions(+), 4 deletions(-) create mode 100644 .changelog/0000.fixed diff --git a/.changelog/0000.fixed b/.changelog/0000.fixed new file mode 100644 index 0000000000..16e2d987e7 --- /dev/null +++ b/.changelog/0000.fixed @@ -0,0 +1,5 @@ +OTLP HTTP exporters now retry on HTTP `429 Too Many Requests`, honor the +`Retry-After` response header (both delay-seconds and HTTP-date forms) when +choosing the delay before the next retry, and clamp the exponential backoff to +a maximum interval. The gRPC exporter backoff is now clamped to the same +maximum. diff --git a/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py b/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py index 4735091091..eab954f432 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-grpc/src/opentelemetry/exporter/otlp/proto/grpc/exporter.py @@ -119,6 +119,10 @@ ] ) _MAX_RETRYS = 6 +# Upper bound (in seconds) applied to the exponential backoff between retries so +# it cannot grow without limit. This mirrors the behavior of the Go and Java +# OTLP exporters, which both cap the backoff interval. +_MAX_BACKOFF = 32 logger = getLogger(__name__) # This prevents logs generated when a log fails to be written to generate another log which fails to be written etc. etc. logger.addFilter(DuplicateFilter()) @@ -434,7 +438,8 @@ def _export( "google.rpc.retryinfo-bin" # type: ignore [reportArgumentType] ) # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min(2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF) if retry_info_bin is not None: retry_info = RetryInfo() retry_info.ParseFromString(retry_info_bin) diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py index b5959300df..4f4e41bfbc 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_common/__init__.py @@ -1,6 +1,8 @@ # Copyright The OpenTelemetry Authors # SPDX-License-Identifier: Apache-2.0 +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime from os import environ from typing import Literal @@ -14,6 +16,11 @@ # 64 MiB, in bytes. _DEFAULT_MAX_REQUEST_SIZE = 64 * 1024 * 1024 +# Upper bound (in seconds) applied to the exponential backoff between retries so +# it cannot grow without limit. This mirrors the behavior of the Go and Java +# OTLP exporters, which both cap the backoff interval. +_MAX_BACKOFF = 32 + class RequestPayloadTooLargeError(Exception): """A serialized OTLP request exceeded the configured ``max_request_size``. @@ -26,6 +33,8 @@ class RequestPayloadTooLargeError(Exception): def _is_retryable(resp: requests.Response) -> bool: if resp.status_code == 408: return True + if resp.status_code == 429: + return True if resp.status_code >= 500 and resp.status_code <= 599: return True return False @@ -42,6 +51,40 @@ def _is_request_too_large(serialized_data: bytes, max_request_size: int) -> bool return max_request_size > 0 and len(serialized_data) > max_request_size +def _parse_retry_after_header(resp: requests.Response) -> float | None: + """Return the ``Retry-After`` delay in seconds, or ``None`` if absent/invalid. + + The ``Retry-After`` header can be either an integer number of seconds + (delay-seconds form) or an HTTP-date. Both forms are defined by RFC 9110 + and the OpenTelemetry OTLP specification requires that the exporter honor + the value when present on a retryable response (e.g. 429 or 503). + """ + retry_after = resp.headers.get("Retry-After") + if retry_after is None: + return None + retry_after = retry_after.strip() + if not retry_after: + return None + try: + # delay-seconds form, e.g. "Retry-After: 120" + return float(int(retry_after)) + except ValueError: + pass + try: + # HTTP-date form, e.g. "Retry-After: Wed, 21 Oct 2015 07:28:00 GMT" + retry_date = parsedate_to_datetime(retry_after) + except (TypeError, ValueError): + return None + if retry_date is None: + return None + if retry_date.tzinfo is None: + retry_date = retry_date.replace(tzinfo=timezone.utc) + delay = (retry_date - datetime.now(timezone.utc)).total_seconds() + if delay < 0: + return 0.0 + return delay + + def _load_session_from_envvar( cred_envvar: Literal[ "OTEL_PYTHON_EXPORTER_OTLP_HTTP_LOGS_CREDENTIAL_PROVIDER", diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py index c7b1decc96..9bd70b10aa 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/_log_exporter/__init__.py @@ -26,10 +26,12 @@ ) from opentelemetry.exporter.otlp.proto.http._common import ( _DEFAULT_MAX_REQUEST_SIZE, + _MAX_BACKOFF, RequestPayloadTooLargeError, _is_request_too_large, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.sdk._logs import ReadableLogRecord @@ -227,7 +229,8 @@ def export(self, batch: Sequence[ReadableLogRecord]) -> LogRecordExportResult: deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min(2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -242,6 +245,11 @@ def export(self, batch: Sequence[ReadableLogRecord]) -> LogRecordExportResult: reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py index 32a9326ae4..62dd9eb551 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/metric_exporter/__init__.py @@ -40,10 +40,12 @@ ) from opentelemetry.exporter.otlp.proto.http._common import ( _DEFAULT_MAX_REQUEST_SIZE, + _MAX_BACKOFF, RequestPayloadTooLargeError, _is_request_too_large, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.proto.collector.metrics.v1.metrics_service_pb2 import ( # noqa: F401 @@ -283,7 +285,8 @@ def _export_with_retries( deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min(2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -298,6 +301,11 @@ def _export_with_retries( reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py index 306a13340b..dfe71b334a 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/src/opentelemetry/exporter/otlp/proto/http/trace_exporter/__init__.py @@ -28,10 +28,12 @@ ) from opentelemetry.exporter.otlp.proto.http._common import ( _DEFAULT_MAX_REQUEST_SIZE, + _MAX_BACKOFF, RequestPayloadTooLargeError, _is_request_too_large, _is_retryable, _load_session_from_envvar, + _parse_retry_after_header, ) from opentelemetry.metrics import MeterProvider from opentelemetry.sdk.environment_variables import ( @@ -222,7 +224,8 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: deadline_sec = time() + self._timeout for retry_num in range(_MAX_RETRYS): # multiplying by a random number between .8 and 1.2 introduces a +/20% jitter to each backoff. - backoff_seconds = 2**retry_num * random.uniform(0.8, 1.2) + # The backoff is clamped to _MAX_BACKOFF so it cannot grow without bound. + backoff_seconds = min(2**retry_num * random.uniform(0.8, 1.2), _MAX_BACKOFF) export_error: Exception | None = None try: resp = self._export(serialized_data, deadline_sec - time()) @@ -237,6 +240,11 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: reason = resp.reason retryable = _is_retryable(resp) status_code = resp.status_code + # Honor a Retry-After header when present, overriding the + # computed backoff with the server-requested delay. + retry_after_seconds = _parse_retry_after_header(resp) + if retry_after_seconds is not None: + backoff_seconds = retry_after_seconds if not retryable: _logger.error( diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py index 93465bd6b3..8d52ccbb46 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/metrics/test_otlp_metrics_exporter.py @@ -1267,6 +1267,41 @@ def test_preferred_aggregation_override(self): self.assertEqual(exporter._preferred_aggregation[Histogram], histogram_aggregation) + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPMetricExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export(self.metrics["sum_int"]), + MetricExportResult.FAILURE, + ) + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while exporting metrics batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPMetricExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export(self.metrics["sum_int"]) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + @patch.dict("os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}) @patch.object(Session, "post") def test_retry_timeout(self, mock_post): diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py index a907a8e0c4..f5c9670140 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_log_exporter.py @@ -484,6 +484,41 @@ def test_retry_timeout(self, mock_post): 503, ) + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPLogExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export(self._get_sdk_log_data()), + LogRecordExportResult.FAILURE, + ) + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while exporting logs batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPLogExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export(self._get_sdk_log_data()) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + @patch.object(Session, "post") def test_export_no_collector_available_retryable(self, mock_post): exporter = OTLPLogExporter(timeout=1.5) diff --git a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py index b577f2b7f4..db6baf9437 100644 --- a/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py +++ b/exporter/opentelemetry-exporter-otlp-proto-http/tests/test_proto_span_exporter.py @@ -5,6 +5,7 @@ import threading import time import unittest +from datetime import datetime, timedelta, timezone from http.server import BaseHTTPRequestHandler, HTTPServer from logging import WARNING from unittest.mock import MagicMock, Mock, patch @@ -18,6 +19,11 @@ encode_spans, ) from opentelemetry.exporter.otlp.proto.http import Compression +from opentelemetry.exporter.otlp.proto.http._common import ( + _MAX_BACKOFF, + _is_retryable, + _parse_retry_after_header, +) from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( DEFAULT_COMPRESSION, DEFAULT_ENDPOINT, @@ -514,6 +520,83 @@ def test_oversized_payload_records_failure_metric(self, mock_post): "RequestPayloadTooLargeError", ) + @patch.object(Session, "post") + def test_429_is_retryable(self, mock_post): + exporter = OTLPSpanExporter(timeout=1.5) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + self.assertEqual( + exporter.export([BASIC_SPAN]), + SpanExportResult.FAILURE, + ) + # A 429 must be retried, so more than a single POST is expected. + self.assertGreater(mock_post.call_count, 1) + self.assertIn( + "Transient error TOO_MANY_REQUESTS encountered while exporting span batch, retrying in", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_seconds_is_honored(self, mock_post): + exporter = OTLPSpanExporter(timeout=10) + + resp = Response() + resp.status_code = 429 + resp.reason = "TOO_MANY_REQUESTS" + resp.headers["Retry-After"] = "2" + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + self.assertIn( + "retrying in 2.00s", + warning.records[0].message, + ) + + @patch.object(Session, "post") + def test_retry_after_header_http_date_is_honored(self, mock_post): + exporter = OTLPSpanExporter(timeout=100) + + resp = Response() + resp.status_code = 503 + resp.reason = "UNAVAILABLE" + retry_at = datetime.now(timezone.utc) + timedelta(seconds=3) + resp.headers["Retry-After"] = retry_at.strftime("%a, %d %b %Y %H:%M:%S GMT") + mock_post.return_value = resp + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + message = warning.records[0].message + self.assertIn("retrying in", message) + reported = float(message.split("retrying in ")[1].rstrip("s.").rstrip("s")) + # The HTTP-date is ~3 seconds in the future. + self.assertTrue(1.5 < reported <= 3.0) + + @patch("opentelemetry.exporter.otlp.proto.http.trace_exporter.random") + @patch.object(Session, "post") + def test_backoff_is_clamped_to_max(self, mock_post, mock_random): + # No jitter, so the raw backoff is a clean power of two. + mock_random.uniform.return_value = 1.0 + + resp = Response() + resp.status_code = 503 + resp.reason = "UNAVAILABLE" + mock_post.return_value = resp + exporter = OTLPSpanExporter(timeout=1000) + with self.assertLogs(level=WARNING) as warning: + exporter.export([BASIC_SPAN]) + reported_backoffs = [] + for record in warning.records: + if "retrying in" in record.message: + reported_backoffs.append(float(record.message.split("retrying in ")[1].rstrip("s.").rstrip("s"))) + # Uncapped, retry #5 would ask for 2**5 == 32 and #6 would be + # larger; every reported backoff must be clamped at _MAX_BACKOFF. + self.assertTrue(reported_backoffs) + for backoff in reported_backoffs: + self.assertLessEqual(backoff, _MAX_BACKOFF) + def test_end_to_end_export_sends_request_over_http(self): server, thread = _start_recording_server() port = server.server_address[1] @@ -549,3 +632,45 @@ def assert_standard_metric_attrs(self, attributes): self.assertTrue(attributes["otel.component.name"].startswith("otlp_http_span_exporter/")) self.assertEqual(attributes["server.address"], "localhost") self.assertEqual(attributes["server.port"], 4318) + + +class TestCommonRetryHelpers(unittest.TestCase): + def _response(self, status_code, headers=None): + resp = Response() + resp.status_code = status_code + if headers: + resp.headers.update(headers) + return resp + + def test_is_retryable_status_codes(self): + for code in (408, 429, 500, 502, 503, 504): + self.assertTrue(_is_retryable(self._response(code)), code) + for code in (200, 400, 401, 403, 404): + self.assertFalse(_is_retryable(self._response(code)), code) + + def test_parse_retry_after_missing(self): + self.assertIsNone(_parse_retry_after_header(self._response(429))) + + def test_parse_retry_after_seconds(self): + self.assertEqual( + _parse_retry_after_header(self._response(429, {"Retry-After": "120"})), + 120.0, + ) + + def test_parse_retry_after_http_date(self): + retry_at = datetime.now(timezone.utc) + timedelta(seconds=30) + header = retry_at.strftime("%a, %d %b %Y %H:%M:%S GMT") + delay = _parse_retry_after_header(self._response(503, {"Retry-After": header})) + self.assertIsNotNone(delay) + self.assertTrue(25 < delay <= 30) + + def test_parse_retry_after_http_date_in_past(self): + retry_at = datetime.now(timezone.utc) - timedelta(seconds=30) + header = retry_at.strftime("%a, %d %b %Y %H:%M:%S GMT") + self.assertEqual( + _parse_retry_after_header(self._response(503, {"Retry-After": header})), + 0.0, + ) + + def test_parse_retry_after_invalid(self): + self.assertIsNone(_parse_retry_after_header(self._response(429, {"Retry-After": "not-a-date"}))) From b31d98342f37bed89b9b1e0dc6643a922ede6942 Mon Sep 17 00:00:00 2001 From: Diego Hurtado Date: Wed, 22 Jul 2026 08:37:51 -0500 Subject: [PATCH 2/2] Rename changelog fragment to match PR number --- .changelog/{0000.fixed => 27.fixed} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename .changelog/{0000.fixed => 27.fixed} (100%) diff --git a/.changelog/0000.fixed b/.changelog/27.fixed similarity index 100% rename from .changelog/0000.fixed rename to .changelog/27.fixed