Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changelog/27.fixed
Original file line number Diff line number Diff line change
@@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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())
Expand Down Expand Up @@ -473,7 +477,10 @@ 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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -11,15 +13,56 @@
)
from opentelemetry.util._importlib_metadata import entry_points

# 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


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


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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,10 @@
Compression,
)
from opentelemetry.exporter.otlp.proto.http._common import (
_MAX_BACKOFF,
_is_retryable,
_load_session_from_envvar,
_parse_retry_after_header,
)
from opentelemetry.metrics import MeterProvider
from opentelemetry.sdk._logs import ReadableLogRecord
Expand Down Expand Up @@ -203,7 +205,10 @@ def export(
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())
Expand All @@ -218,6 +223,11 @@ def export(
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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,10 @@
Compression,
)
from opentelemetry.exporter.otlp.proto.http._common import (
_MAX_BACKOFF,
_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
Expand Down Expand Up @@ -274,7 +276,10 @@ 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())
Expand All @@ -289,6 +294,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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,10 @@
Compression,
)
from opentelemetry.exporter.otlp.proto.http._common import (
_MAX_BACKOFF,
_is_retryable,
_load_session_from_envvar,
_parse_retry_after_header,
)
from opentelemetry.metrics import MeterProvider
from opentelemetry.sdk.environment_variables import (
Expand Down Expand Up @@ -196,7 +198,10 @@ 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())
Expand All @@ -211,6 +216,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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1351,6 +1351,42 @@ def test_preferred_aggregation_override(self):
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"}
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -529,6 +529,42 @@ 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)
Expand Down
Loading
Loading