Skip to content

Commit 9d68b59

Browse files
Outlast metadata throttling with jittered retries and an existence fallback
Metadata API throttling episodes observed in CI and in a controlled probe last tens of seconds, while the default retry policy waited about 15 seconds in total and clients retried in lockstep. - Add uniform jitter of up to one multiplier to retry_api_call waits and raise the default RetryConfig attempts to seven, so the exponential waits sum to 63 seconds plus jitter. - Expose is_retryable_error and reuse it in the dialect: when a table-metadata request is still throttled after PyAthena's retries, has_table() determines existence with information_schema.tables and logs a warning. The fallback answers existence only; column, comment and table-option reflection still propagate the throttling error, and nothing is cached from the fallback. - Expose retry_config on the async connection adapter. - Cover the fallback for direct and wrapped throttling codes on existing and missing tables, keep the error-propagation regression for access denial and unrecognized messages, and add wait-bound and jitter unit tests. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1 parent ea4ce1f commit 9d68b59

6 files changed

Lines changed: 218 additions & 19 deletions

File tree

‎docs/sqlalchemy.md‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,13 +98,18 @@ A later listing preserves metadata already fetched for a table.
9898
`clear_cache()` also discards this metadata; an absent entry in a listing is not cached as proof that a table does not exist.
9999

100100
Table-metadata lookups propagate throttling and permission errors rather than reporting missing tables.
101-
`has_table()` propagates these failures, including access denied by Lake Formation, instead of returning or caching `False`.
101+
`has_table()` propagates permission failures, including access denied by Lake Formation, instead of returning or caching `False`.
102102
For failed metadata requests, only recognized `EntityNotFoundException` responses establish absence; unrecognized errors are propagated rather than guessed to mean a missing table.
103+
When a table-metadata request is still throttled after PyAthena's retries, `has_table()` determines existence with a query on `information_schema.tables` and logs a warning.
104+
This fallback answers only existence: it does not populate the metadata cache, and column, comment, and table-option reflection still propagate the throttling error.
105+
Athena applies its metadata API rate limits per account; they are not listed in Service Quotas, and throttling episodes can last tens of seconds.
106+
PyAthena's API retries use exponential backoff with uniform jitter; the default `RetryConfig` makes seven attempts whose waits sum to 63 seconds plus jitter.
103107
PyAthena recognizes Glue error codes in Athena's `MetadataException` service-error envelope and applies `RetryConfig.exceptions` to those codes.
104108
`RetryConfig` accepts one exception-name string or an iterable and captures the names as a tuple at construction.
105109
Changes to the original input list or iterator no longer change the stored policy; construct a new `RetryConfig` when changing the retry policy.
106110
SDK retries and PyAthena retries are separate layers, so increasing both attempt limits can multiply requests and waiting time.
107111
Adaptive SDK retries regulate individual clients, not the aggregate traffic from independent CI runners.
112+
For highly concurrent reflection or `checkfirst` DDL, bound the concurrency of metadata requests and consider `botocore.config.Config(retries={"mode": "adaptive"})` on the connection.
108113

109114
Use SQLAlchemy's identifier quoting for reserved words or names beginning with an underscore.
110115
The dialect uses backticks for table DDL and double quotes for DML.

‎pyathena/aio/sqlalchemy/base.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
ProgrammingError,
2323
)
2424
from pyathena.sqlalchemy.base import AthenaDialect
25+
from pyathena.util import RetryConfig
2526

2627
if TYPE_CHECKING:
2728
from types import ModuleType
@@ -142,6 +143,10 @@ def schema_name(self) -> str | None:
142143
def cursor_kwargs(self) -> dict[str, Any]:
143144
return self._connection.cursor_kwargs # type: ignore[no-any-return]
144145

146+
@property
147+
def retry_config(self) -> RetryConfig:
148+
return self._connection.retry_config # type: ignore[no-any-return]
149+
145150
def cursor(self) -> AsyncAdapt_pyathena_cursor:
146151
raw_cursor = self._connection.cursor()
147152
return AsyncAdapt_pyathena_cursor(raw_cursor)

‎pyathena/sqlalchemy/base.py‎

Lines changed: 38 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
from __future__ import annotations
22

33
import contextlib
4+
import logging
45
import re
56
from collections.abc import Mapping, MutableMapping
67
from re import Pattern
@@ -37,7 +38,7 @@
3738
get_double_type,
3839
)
3940
from pyathena.sqlalchemy.util import _HashableDict
40-
from pyathena.util import _get_error_code, strtobool
41+
from pyathena.util import _get_error_code, is_retryable_error, strtobool
4142

4243
if TYPE_CHECKING:
4344
from types import ModuleType
@@ -55,6 +56,8 @@
5556
)
5657
from sqlalchemy.sql.schema import SchemaItem
5758

59+
_logger = logging.getLogger(__name__)
60+
5861

5962
ischema_names: dict[str, type[Any]] = {
6063
"boolean": types.BOOLEAN,
@@ -384,6 +387,40 @@ def has_table(self, connection: Connection, table_name: str, schema: str | None
384387
return bool(columns)
385388
except exc.NoSuchTableError:
386389
return False
390+
except pyathena.error.OperationalError as e:
391+
cause = e.__cause__
392+
raw_connection = self._raw_connection(connection)
393+
if not (
394+
isinstance(cause, botocore.exceptions.ClientError)
395+
and is_retryable_error(cause, raw_connection.retry_config) # type: ignore[union-attr]
396+
):
397+
raise
398+
# The metadata API is still throttled after PyAthena's retries.
399+
# Existence can be answered by a query; column, comment and option
400+
# reflection cannot, so only this check falls back.
401+
_logger.warning(
402+
f"Table metadata request for {table_name} was throttled; "
403+
"checking existence with information_schema."
404+
)
405+
return self._has_table_information_schema(connection, table_name, schema)
406+
407+
def _has_table_information_schema(
408+
self, connection: Connection, table_name: str, schema: str | None = None
409+
) -> bool:
410+
raw_connection = self._raw_connection(connection)
411+
schema = schema if schema else raw_connection.schema_name # type: ignore[union-attr]
412+
catalog = raw_connection.cursor_kwargs.get("catalog_name", raw_connection.catalog_name)
413+
name = (
414+
str(table_name).lower()
415+
if (catalog or "").lower() == "awsdatacatalog"
416+
else str(table_name)
417+
)
418+
query = text(
419+
"SELECT table_name FROM information_schema.tables "
420+
"WHERE table_schema = :schema AND table_name = :table_name"
421+
)
422+
rows = connection.execute(query, {"schema": schema, "table_name": name}).fetchall()
423+
return bool(rows)
387424

388425
@reflection.cache
389426
def get_view_definition(

‎pyathena/util.py‎

Lines changed: 41 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,13 @@
77
from typing import Any
88

99
import tenacity
10-
from tenacity import after_log, retry_if_exception, stop_after_attempt, wait_exponential
10+
from tenacity import (
11+
after_log,
12+
retry_if_exception,
13+
stop_after_attempt,
14+
wait_exponential,
15+
wait_random,
16+
)
1117

1218
from pyathena import DataError
1319

@@ -86,9 +92,11 @@ class RetryConfig:
8692
8793
Attributes:
8894
exceptions: Tuple of AWS exception names to retry on.
89-
attempt: Maximum number of retry attempts.
90-
multiplier: Base multiplier for exponential backoff.
91-
max_delay: Maximum delay between retries in seconds.
95+
attempt: Maximum number of attempts, including the first call.
96+
multiplier: Base multiplier for exponential backoff in seconds. Each
97+
wait also adds uniform random jitter of up to one multiplier.
98+
max_delay: Maximum exponential delay between retries in seconds,
99+
before jitter is added.
92100
exponential_base: Base for exponential backoff calculation.
93101
94102
Example:
@@ -116,6 +124,10 @@ class RetryConfig:
116124
original iterable do not change this configuration.
117125
Retries are applied to AWS API calls, not to SQL query execution.
118126
Query failures typically require manual intervention or query fixes.
127+
With the default settings, the exponential waits between attempts sum
128+
to 63 seconds plus jitter, which outlasts the metadata API throttling
129+
episodes observed under concurrent reflection. SDK retries configured
130+
on the boto3 client are a separate layer applied within each attempt.
119131
Recognized Glue error codes wrapped in Athena MetadataException are
120132
matched against exceptions in the same way as direct AWS error codes.
121133
"""
@@ -126,7 +138,7 @@ def __init__(
126138
"ThrottlingException",
127139
"TooManyRequestsException",
128140
),
129-
attempt: int = 5,
141+
attempt: int = 7,
130142
multiplier: int = 1,
131143
max_delay: int = 100,
132144
exponential_base: int = 2,
@@ -159,6 +171,23 @@ def _get_error_code(ex: BaseException, unwrap_metadata: bool = False) -> str | N
159171
return code if isinstance(code, str) else None
160172

161173

174+
def is_retryable_error(ex: BaseException, config: RetryConfig) -> bool:
175+
"""Return whether an exception matches the retry policy's AWS error codes.
176+
177+
Args:
178+
ex: The exception raised by an AWS API call.
179+
config: RetryConfig whose ``exceptions`` list the retryable error codes.
180+
181+
Returns:
182+
True if the direct error code, or a recognized Glue error code wrapped in
183+
an Athena ``MetadataException``, is listed in ``config.exceptions``.
184+
"""
185+
return any(
186+
code is not None and code in config.exceptions
187+
for code in (_get_error_code(ex), _get_error_code(ex, unwrap_metadata=True))
188+
)
189+
190+
162191
def retry_api_call(
163192
func: Callable[..., Any],
164193
config: RetryConfig,
@@ -169,8 +198,8 @@ def retry_api_call(
169198
"""Execute a function with automatic retry logic for AWS API calls.
170199
171200
This function wraps AWS API calls with retry behavior based on the provided
172-
configuration. It uses exponential backoff and only retries on specific
173-
AWS exceptions that indicate transient failures.
201+
configuration. It uses exponential backoff with uniform jitter and only
202+
retries on specific AWS exceptions that indicate transient failures.
174203
175204
Args:
176205
func: The AWS API function to call.
@@ -201,20 +230,17 @@ def retry_api_call(
201230
Other errors are propagated without retrying.
202231
"""
203232

204-
def _is_retryable(ex: BaseException) -> bool:
205-
return any(
206-
code is not None and code in config.exceptions
207-
for code in (_get_error_code(ex), _get_error_code(ex, unwrap_metadata=True))
208-
)
209-
210233
retry = tenacity.Retrying(
211-
retry=retry_if_exception(_is_retryable),
234+
retry=retry_if_exception(lambda ex: is_retryable_error(ex, config)),
212235
stop=stop_after_attempt(config.attempt),
236+
# Uniform jitter of up to one multiplier keeps concurrent clients from
237+
# retrying in lockstep after a shared throttling response.
213238
wait=wait_exponential(
214239
multiplier=config.multiplier,
215240
max=config.max_delay,
216241
exp_base=config.exponential_base,
217-
),
242+
)
243+
+ wait_random(0, config.multiplier),
218244
after=after_log(logger, logger.getEffectiveLevel()) if logger else None, # type: ignore[arg-type]
219245
reraise=True,
220246
)

‎tests/pyathena/test_util.py‎

Lines changed: 87 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,17 @@
11
from typing import Any
22

33
import pytest
4+
import tenacity.nap
45
from botocore.exceptions import ClientError
56

67
from pyathena import DataError
7-
from pyathena.util import RetryConfig, parse_output_location, retry_api_call, strtobool
8+
from pyathena.util import (
9+
RetryConfig,
10+
is_retryable_error,
11+
parse_output_location,
12+
retry_api_call,
13+
strtobool,
14+
)
815

916

1017
def test_parse_output_location():
@@ -160,3 +167,82 @@ def call():
160167

161168
assert retry_api_call(call, config) == "success"
162169
assert calls == 3
170+
171+
172+
def _throttling_error() -> ClientError:
173+
return ClientError(
174+
{"Error": {"Code": "ThrottlingException", "Message": "Rate exceeded"}},
175+
"GetTableMetadata",
176+
)
177+
178+
179+
def test_retry_config_default_attempts():
180+
assert RetryConfig().attempt == 7
181+
182+
183+
@pytest.mark.parametrize(
184+
("multiplier", "max_delay", "expected_bases"),
185+
[
186+
(1, 100, [1, 2, 4, 8, 16, 32]),
187+
(2, 10, [2, 4, 8, 10, 10, 10]),
188+
(0, 0, [0, 0, 0, 0, 0, 0]),
189+
],
190+
)
191+
def test_retry_api_call_waits_with_jitter(monkeypatch, multiplier, max_delay, expected_bases):
192+
sleeps = []
193+
monkeypatch.setattr(tenacity.nap.time, "sleep", sleeps.append)
194+
error = _throttling_error()
195+
calls = 0
196+
197+
def call():
198+
nonlocal calls
199+
calls += 1
200+
raise error
201+
202+
config = RetryConfig(attempt=7, multiplier=multiplier, max_delay=max_delay)
203+
with pytest.raises(ClientError):
204+
retry_api_call(call, config)
205+
assert calls == 7
206+
assert len(sleeps) == len(expected_bases)
207+
for waited, base in zip(sleeps, expected_bases, strict=True):
208+
assert base <= waited <= base + multiplier
209+
210+
211+
def test_retry_api_call_jitter_varies(monkeypatch):
212+
sleeps = []
213+
monkeypatch.setattr(tenacity.nap.time, "sleep", sleeps.append)
214+
error = _throttling_error()
215+
216+
def call():
217+
raise error
218+
219+
config = RetryConfig(attempt=20, multiplier=1, max_delay=1)
220+
with pytest.raises(ClientError):
221+
retry_api_call(call, config)
222+
assert all(1 <= waited <= 2 for waited in sleeps)
223+
assert len(set(sleeps)) > 1
224+
225+
226+
@pytest.mark.parametrize(
227+
("code", "message", "expected"),
228+
[
229+
("ThrottlingException", "Rate exceeded", True),
230+
(
231+
"MetadataException",
232+
"Rate exceeded (Service: AmazonDataCatalog; Status Code: 400; "
233+
"Error Code: ThrottlingException; Request ID: example; Proxy: null)",
234+
True,
235+
),
236+
(
237+
"MetadataException",
238+
"Not authorized (Service: AmazonDataCatalog; Status Code: 400; "
239+
"Error Code: AccessDeniedException; Request ID: example; Proxy: null)",
240+
False,
241+
),
242+
("MetadataException", "Table not found", False),
243+
],
244+
)
245+
def test_is_retryable_error(code, message, expected):
246+
error = ClientError({"Error": {"Code": code, "Message": message}}, "GetTableMetadata")
247+
assert is_retryable_error(error, RetryConfig()) is expected
248+
assert is_retryable_error(ValueError("no response"), RetryConfig()) is False

‎tests/sqlalchemy/test_suite.py‎

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -264,7 +264,47 @@ def test_long_convention_name(self, type_, metadata, connection):
264264

265265

266266
class HasTableTest(_HasTableTest):
267-
@sa_testing.combinations("AccessDeniedException", "ThrottlingException", None, argnames="code")
267+
@sa_testing.combinations(True, False, argnames="wrapped")
268+
@sa_testing.combinations(True, False, argnames="exists")
269+
def test_throttled_existence_check_uses_information_schema(
270+
self, connection, metadata, monkeypatch, wrapped, exists
271+
):
272+
table = Table("throttled_existence", metadata, Column("id", Integer))
273+
table.create(connection)
274+
raw_connection = connection.connection.driver_connection
275+
if connection.dialect.is_async:
276+
raw_connection = raw_connection.driver_connection
277+
error = ClientError(
278+
{
279+
"Error": {
280+
"Code": "MetadataException" if wrapped else "ThrottlingException",
281+
"Message": "Rate exceeded (Service: AmazonDataCatalog; Status Code: 400; "
282+
"Error Code: ThrottlingException; Request ID: example; Proxy: null)"
283+
if wrapped
284+
else "Rate exceeded",
285+
}
286+
},
287+
"GetTableMetadata",
288+
)
289+
calls = []
290+
291+
def fail_metadata(**kwargs):
292+
calls.append(kwargs)
293+
raise error
294+
295+
monkeypatch.setattr(raw_connection.client, "get_table_metadata", fail_metadata)
296+
monkeypatch.setattr(raw_connection.retry_config, "attempt", 1)
297+
inspector = inspect(connection)
298+
name = table.name if exists else "throttled_missing"
299+
assert inspector.has_table(name) is exists
300+
assert inspector.has_table(name.upper(), schema=raw_connection.schema_name) is exists
301+
assert len(calls) == 2
302+
# Only existence falls back; reflection still reports the throttled request.
303+
with pytest.raises(OperationalError) as caught:
304+
inspector.get_columns(table.name)
305+
assert caught.value.__cause__ is error
306+
307+
@sa_testing.combinations("AccessDeniedException", None, argnames="code")
268308
def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, code):
269309
raw_connection = connection.connection.driver_connection
270310
if connection.dialect.is_async:

0 commit comments

Comments
 (0)