From ff51b4e6979b7519b839648eaaa07896b030e014 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 21:30:00 +0900 Subject: [PATCH 01/10] Read the dialect's own queries through an API cursor The dialect parses the rows of two queries it issues itself, and both ran through whichever cursor class the connection was configured with. Those cursors disagree on NULL and blank values, so: - get_view_definition() raised TypeError on pandas and polars when a view's DDL contained a blank line, and silently dropped that line on arrow. - the information_schema fallback needed a NaN guard, and wrote a Parquet UNLOAD to S3 for a five-row metadata lookup under unload=true. Both now open a plain Cursor. Reflection no longer inherits a result format chosen for user queries, and the NaN guard goes away with it. Also widen the fallback to a MetadataException that survived unwrapping. A federated catalog reports a missing table in its connector's own words, with no Glue error envelope, so has_table() raised OperationalError instead of returning False and create_all(checkfirst=True) failed against one. Deciding absence from information_schema keeps throttling and connector outages distinct from a genuinely missing table. Co-Authored-By: Claude Opus 5 --- docs/sqlalchemy.md | 10 +- pyathena/sqlalchemy/base.py | 77 ++++++---- tests/pyathena/sqlalchemy/test_base.py | 197 ++++++++++++++++++++++--- 3 files changed, 230 insertions(+), 54 deletions(-) diff --git a/docs/sqlalchemy.md b/docs/sqlalchemy.md index 041906a0..b3ae7cef 100644 --- a/docs/sqlalchemy.md +++ b/docs/sqlalchemy.md @@ -100,12 +100,14 @@ A later listing preserves metadata already fetched for a table. Table-metadata lookups propagate throttling and permission errors rather than reporting missing tables. `has_table()` propagates permission failures, including access denied by Lake Formation, instead of returning or caching `False`. For failed metadata requests, only recognized `EntityNotFoundException` responses establish absence; unrecognized errors are propagated rather than guessed to mean a missing table. -Column reflection and `has_table()` do not retry a throttled table-metadata request; when Athena reports throttling, they read `information_schema.columns` instead, executed without query result reuse, and log a warning. -Other error codes listed in the connection's `RetryConfig.exceptions` are still retried on that path, except `MetadataException` itself, which carries wrapped throttling; list the wrapped Glue codes instead. -A `retry_config` in `cursor_kwargs` replaces that policy entirely, including its throttling retries, which then run before the fallback. +Column reflection and `has_table()` do not retry a table-metadata request that `information_schema` can answer; they read `information_schema.columns` instead, executed without query result reuse, and log a warning. +That covers a throttled request, and a `MetadataException` carrying no recognized Glue error envelope, which is how a federated catalog reports a missing table: absence is then decided by the query against the same catalog rather than by an unrecognized message. +Other error codes listed in the connection's `RetryConfig.exceptions` are still retried on that path, except those; list the wrapped Glue codes instead of `MetadataException`. +A `retry_config` in `cursor_kwargs` replaces that policy entirely, including those retries, which then run before the fallback. The fallback maps unbounded `varchar` to SQLAlchemy `String`, matching Hive `STRING` reflection from the metadata API, and preserves explicit `VARCHAR(n)` and `CHAR(n)` lengths. Partition columns are marked from the `extra_info` column. -This fallback does not populate the table-metadata cache, and table comments and table options still come from the metadata API with the configured retries, so they propagate the throttling error. +This fallback does not populate the table-metadata cache, and table comments and table options still come from the metadata API with the configured retries, so they propagate the error. +The dialect runs its own queries — this fallback and `get_view_definition()` — through the API cursor, whatever `cursor_class` or `unload` setting the connection carries, because it parses those result rows itself. Athena applies its metadata API rate limits per account, and they are not listed in Service Quotas. PyAthena's API retries use exponential backoff with uniform jitter; `RetryConfig` documents the default attempt count and waits. PyAthena recognizes Glue error codes in Athena's `MetadataException` service-error envelope and applies `RetryConfig.exceptions` to those codes. diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 6a3850ea..1abb122d 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -11,7 +11,7 @@ cast, ) -from sqlalchemy import exc, schema, text, types, util +from sqlalchemy import exc, schema, types, util from sqlalchemy.engine import Engine, reflection from sqlalchemy.engine.default import DefaultDialect from sqlalchemy.engine.interfaces import ExecutionContext @@ -23,6 +23,7 @@ ) import pyathena +from pyathena.cursor import Cursor from pyathena.sqlalchemy.compiler import ( AthenaDDLCompiler, AthenaStatementCompiler, @@ -195,6 +196,12 @@ class AthenaDialect(DefaultDialect): _connect_options: dict[str, Any] = {} # type: ignore[override] # noqa: RUF012 _pattern_column_type: Pattern[str] = re.compile(r"^([a-zA-Z]+)(?:$|[\(|<](.+)[\)|>]$)") + # Metadata failures that information_schema answers better than a retry. + # Throttling, because one query costs less than the retry ladder, and a + # MetadataException that survived unwrapping, because a federated catalog + # reports a missing table in its connector's words rather than in Glue's + # EntityNotFoundException envelope. + _FALLBACK_ERROR_CODES: tuple[str, ...] = (*THROTTLING_ERROR_CODES, "MetadataException") def __init__(self, json_deserializer=None, json_serializer=None, **kwargs): DefaultDialect.__init__(self, **kwargs) @@ -347,13 +354,13 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * metadata = info_cache.get(metadata_key) if metadata is not None: return self._columns_from_metadata(metadata) - # A throttled metadata request switches to information_schema at once - # instead of waiting out the retry policy; the query answers existence - # and columns, while table comments and options still need the API. - # Other retryable codes keep the connection's policy. Connection.cursor() - # applies cursor_kwargs last, so a retry_config given there still runs - # its own throttling retries before the fallback. - retry_config = self._without_throttling_retries( + # A metadata request the fallback can answer switches to + # information_schema at once instead of waiting out the retry policy; the + # query answers existence and columns, while table comments and options + # still need the API. Other retryable codes keep the connection's policy. + # Connection.cursor() applies cursor_kwargs last, so a retry_config given + # there still runs its own retries before the fallback. + retry_config = self._without_fallback_retries( raw_connection.retry_config # type: ignore[union-attr] ) with raw_connection.driver_connection.cursor( # type: ignore[union-attr] @@ -362,13 +369,11 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * try: metadata = self._lookup_table(cursor, schema, name, table_name) except pyathena.error.OperationalError as e: - if ( - _get_error_code(e.__cause__ or e, unwrap_metadata=True) - not in THROTTLING_ERROR_CODES - ): + code = _get_error_code(e.__cause__ or e, unwrap_metadata=True) + if code not in self._FALLBACK_ERROR_CODES: raise _logger.warning( - f"Table metadata request for {table_name} was throttled; " + f"Table metadata request for {table_name} failed with {code}; " "reflecting columns from information_schema." ) columns = self._columns_from_information_schema(raw_connection, schema, name) @@ -379,22 +384,35 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * info_cache[metadata_key] = metadata return self._columns_from_metadata(metadata) - @staticmethod - def _without_throttling_retries(retry_config: RetryConfig) -> RetryConfig: - """Copy a policy without the codes that carry throttling. + @classmethod + def _without_fallback_retries(cls, retry_config: RetryConfig) -> RetryConfig: + """Copy a policy without the codes the information_schema fallback answers. - Athena wraps Glue throttling in ``MetadataException``, so that code is - dropped as well; specific wrapped codes stay retryable. + Retrying those spends the policy's whole budget on a question one query + settles; specific wrapped Glue codes stay retryable. """ - excluded = (*THROTTLING_ERROR_CODES, "MetadataException") return RetryConfig( - exceptions=[c for c in retry_config.exceptions if c not in excluded], + exceptions=[c for c in retry_config.exceptions if c not in cls._FALLBACK_ERROR_CODES], attempt=retry_config.attempt, multiplier=retry_config.multiplier, max_delay=retry_config.max_delay, exponential_base=retry_config.exponential_base, ) + @staticmethod + def _internal_cursor(raw_connection: PoolProxiedConnection) -> Cursor: + """Open an API cursor for the queries this dialect parses itself. + + Reflection reads these rows directly, so they must not arrive in the + result format chosen for user queries: a DataFrame cursor reports a NULL + or blank value as NaN, as an empty string, or as a dropped row depending + on its backend and on UNLOAD. + """ + return cast( + Cursor, + raw_connection.driver_connection.cursor(Cursor), # type: ignore[union-attr] + ) + def _column(self, name: str | None, type_: str, comment: str | None, partition: bool | None): return { "name": name, @@ -420,7 +438,7 @@ def _columns_from_information_schema( # The answer must reflect the catalog now, so query result reuse is off. schema = str(schema).lower().replace("'", "''") table_name = table_name.lower().replace("'", "''") - with raw_connection.driver_connection.cursor() as cursor: # type: ignore[union-attr] + with self._internal_cursor(raw_connection) as cursor: cursor.execute( "SELECT ordinal_position, column_name, data_type, comment, extra_info " "FROM information_schema.columns " @@ -428,15 +446,13 @@ def _columns_from_information_schema( result_reuse_enable=False, ) rows = cursor.fetchall() - # Sort here: the query has no ORDER BY and UNLOAD-backed cursors do not - # preserve result order. A CSV-backed pandas cursor reads a missing - # comment as NaN, which _column() cannot recognize as empty. + # Sort here: the query has no ORDER BY, so its result order is Athena's. return [ self._column( column_name, # Athena exposes Hive STRING as unbounded VARCHAR in information_schema. "string" if data_type == "varchar" else data_type, - comment if isinstance(comment, str) else None, + comment, extra_info == "partition key" or None, ) for _, column_name, data_type, comment, extra_info in sorted( @@ -524,11 +540,14 @@ def get_view_definition( schema = schema if schema else self._cursor_option(raw_connection, "schema_name") query = f"""SHOW CREATE VIEW "{schema}"."{view_name}";""" try: - res = connection.scalars(text(query)) - except exc.OperationalError as e: + with self._internal_cursor(raw_connection) as cursor: + cursor.execute(query) + rows = cursor.fetchall() + except pyathena.error.OperationalError as e: raise exc.NoSuchTableError(f"{schema}.{view_name}") from e - else: - return "\n".join(res) + # Athena returns the definition one line per row and blank lines as + # empty values, which are part of the definition. + return "\n".join(row[0] or "" for row in rows) @reflection.cache def get_columns(self, connection: Connection, table_name: str, schema: str | None = None, **kw): diff --git a/tests/pyathena/sqlalchemy/test_base.py b/tests/pyathena/sqlalchemy/test_base.py index 2dcda8f8..fd80b5ca 100644 --- a/tests/pyathena/sqlalchemy/test_base.py +++ b/tests/pyathena/sqlalchemy/test_base.py @@ -19,6 +19,7 @@ from sqlalchemy.sql.schema import Column, MetaData, Table from sqlalchemy.sql.selectable import TextualSelect +from pyathena.cursor import Cursor from pyathena.error import OperationalError from pyathena.sqlalchemy.base import AthenaDialect from pyathena.sqlalchemy.types import ( @@ -54,28 +55,36 @@ def unique_s3tables_table_name(base: str) -> str: class TestAthenaDialect: def test_columns_from_information_schema(self): - # Rows arrive unordered, and a cursor may read a NULL comment as NaN - # (pandas), as an empty string (arrow), or as None (polars). + # Rows arrive unordered, and Athena reports a missing comment as NULL, + # which an API cursor hands over as None or as an empty string. rows = [ - ("4", "dt", "varchar", float("nan"), "partition key"), + ("4", "dt", "varchar", None, "partition key"), ("1", "id", "integer", "identifier", None), ("2", "payload", "row(a integer, b array(varchar))", None, None), ("3", "label", "varchar", "", ""), ] executed = [] + opened = [] def execute(operation, **kwargs): executed.append((operation, kwargs)) cursor = SimpleNamespace(execute=execute, fetchall=lambda: rows) - raw_connection = SimpleNamespace( - driver_connection=SimpleNamespace(cursor=lambda: contextlib.nullcontext(cursor)) - ) + + def open_cursor(cursor_class): + opened.append(cursor_class) + return contextlib.nullcontext(cursor) + + raw_connection = SimpleNamespace(driver_connection=SimpleNamespace(cursor=open_cursor)) columns = AthenaDialect()._columns_from_information_schema( raw_connection, "My_Schema", "O'Neil" ) + # The dialect parses these rows itself, so it asks for an API cursor + # rather than whatever result format the user configured. + assert opened == [Cursor] + assert [column["name"] for column in columns] == ["id", "payload", "label", "dt"] assert isinstance(columns[0]["type"], types.INTEGER) assert isinstance(columns[1]["type"], AthenaStruct) @@ -104,7 +113,7 @@ def test_empty_metadata_comment_is_no_comment(self): assert [column["comment"] for column in columns] == [None, None] - def test_without_throttling_retries_keeps_other_codes(self): + def test_without_fallback_retries_keeps_other_codes(self): policy = RetryConfig( exceptions=( "ThrottlingException", @@ -117,8 +126,8 @@ def test_without_throttling_retries_keeps_other_codes(self): max_delay=30, exponential_base=3, ) - derived = AthenaDialect._without_throttling_retries(policy) - # MetadataException carries wrapped throttling, so it is dropped as well. + derived = AthenaDialect._without_fallback_retries(policy) + # The fallback answers these, so retrying them only spends the budget. assert derived.exceptions == ("InternalServerException",) assert (derived.attempt, derived.multiplier, derived.max_delay) == (10, 2, 30) assert derived.exponential_base == 3 @@ -149,7 +158,7 @@ def get_table_metadata(table_name, **kwargs): schema_name=None, retry_config=RetryConfig(), driver_connection=SimpleNamespace( - cursor=lambda **kwargs: contextlib.nullcontext(cursor) + cursor=lambda *args, **kwargs: contextlib.nullcontext(cursor) ), ) connection = SimpleNamespace(connection=raw_connection) @@ -186,7 +195,7 @@ def get_table_metadata(table_name, **kwargs): schema_name="default", retry_config=RetryConfig(), driver_connection=SimpleNamespace( - cursor=lambda **kwargs: contextlib.nullcontext(cursor) + cursor=lambda *args, **kwargs: contextlib.nullcontext(cursor) ), ) connection = SimpleNamespace(connection=raw_connection) @@ -197,6 +206,113 @@ def get_table_metadata(table_name, **kwargs): # Absence is not cached as reflected columns. assert info_cache == {} + @pytest.mark.parametrize( + ("rows", "expected"), + [ + ([], None), + ([("1", "id", "integer", None, None)], ["id"]), + ], + ids=["absent", "present"], + ) + def test_unrecognized_metadata_error_asks_information_schema(self, rows, expected): + # A federated catalog reports a missing table in its connector's own + # words, with no Glue error envelope to unwrap, so the code stays + # MetadataException. Guessing "missing" from that is what turned + # throttling into false absence; ask information_schema instead. + error = ClientError( + { + "Error": { + "Code": "MetadataException", + "Message": ( + "Failed to invoke lambda function due to " + "com.amazonaws.services.lambda.invoke.LambdaFunctionException: " + "Requested resource not found " + "(Service: DynamoDb, Status Code: 400, Request ID: example)" + ), + } + }, + "GetTableMetadata", + ) + + def get_table_metadata(table_name, **kwargs): + raise OperationalError(*error.args) from error + + cursor = SimpleNamespace( + get_table_metadata=get_table_metadata, + execute=lambda operation, **kwargs: None, + fetchall=lambda: rows, + ) + raw_connection = SimpleNamespace( + cursor_kwargs={}, + catalog_name="federated_catalog", + schema_name="default", + retry_config=RetryConfig(), + driver_connection=SimpleNamespace( + cursor=lambda *args, **kwargs: contextlib.nullcontext(cursor) + ), + ) + connection = SimpleNamespace(connection=raw_connection) + + if expected is None: + with pytest.raises(NoSuchTableError): + AthenaDialect()._get_columns(connection, "events") + else: + columns = AthenaDialect()._get_columns(connection, "events") + assert [column["name"] for column in columns] == expected + + def test_get_view_definition_keeps_blank_lines(self): + # Athena returns the definition one row per line, blank lines included. + # Reading them through the user's cursor loses or corrupts those rows, + # so this path asks for an API cursor. + rows = [("CREATE VIEW v AS",), ("",), ("SELECT 1",)] + executed = [] + opened = [] + + def open_cursor(cursor_class): + opened.append(cursor_class) + return contextlib.nullcontext( + SimpleNamespace( + execute=lambda operation, **kwargs: executed.append(operation), + fetchall=lambda: rows, + ) + ) + + raw_connection = SimpleNamespace( + cursor_kwargs={}, + schema_name="default", + driver_connection=SimpleNamespace(cursor=open_cursor), + ) + connection = SimpleNamespace(connection=raw_connection) + + definition = AthenaDialect().get_view_definition(connection, "v") + + assert definition == "CREATE VIEW v AS\n\nSELECT 1" + assert opened == [Cursor] + assert executed == ['SHOW CREATE VIEW "default"."v";'] + + def test_get_view_definition_reports_a_missing_view(self): + error = ClientError( + {"Error": {"Code": "InvalidRequestException", "Message": "does not exist"}}, + "StartQueryExecution", + ) + + def execute(operation, **kwargs): + raise OperationalError(*error.args) from error + + raw_connection = SimpleNamespace( + cursor_kwargs={}, + schema_name="default", + driver_connection=SimpleNamespace( + cursor=lambda *args, **kwargs: contextlib.nullcontext( + SimpleNamespace(execute=execute, fetchall=list) + ) + ), + ) + connection = SimpleNamespace(connection=raw_connection) + + with pytest.raises(NoSuchTableError): + AthenaDialect().get_view_definition(connection, "v") + def test_get_table_matches_long_names_case_insensitively(self): # GetTableMetadata rejects names over 128 characters, so the lookup lists # with a lowercase filter and matches the catalog's own casing. @@ -661,8 +777,9 @@ def test_get_columns(self, engine): assert not actual["autoincrement"] assert actual["comment"] == "some comment" - # `unload` states what each case intends, independently of the URL the - # fixture builds, so the executed mode is compared against the intent. + # `unload` states what each case configures, independently of the URL the + # fixture builds, so the engine's setting can be checked before asserting + # that the dialect's own query ignores it. @pytest.mark.parametrize( ("engine", "unload"), [ @@ -692,6 +809,9 @@ def test_throttled_columns_across_cursor_types(self, engine, unload, monkeypatch ).create(bind=conn) raw_connection = conn.connection.driver_connection + # The engine really is configured the way this case says; the assertion + # below then means the dialect ignored it, not that it never arrived. + assert raw_connection.cursor_kwargs.get("unload", False) is unload error = ClientError( {"Error": {"Code": "ThrottlingException", "Message": "Rate exceeded"}}, "GetTableMetadata", @@ -717,19 +837,18 @@ def record_query(params, **kwargs): finally: raw_connection.client.meta.events.unregister(event, record_query) - # Reflection issued the fallback and nothing else, in the mode this case - # asked for. Without the mode check, an unload option that stopped - # reaching the cursor would leave every case green while three of them - # silently tested CSV twice. The count is not pinned: a throttled - # StartQueryExecution is retried through the same client hook. + # Reflection issued the fallback and nothing else, and it ran as a plain + # query even where the engine asked for UNLOAD: the dialect parses these + # rows itself, so it uses an API cursor. The count is not pinned, since a + # throttled StartQueryExecution is retried through the same client hook. assert queries for query in queries: assert "FROM information_schema.columns" in query - assert query.strip().startswith("UNLOAD (") is unload + assert not query.strip().startswith("UNLOAD (") # The fallback query has no ORDER BY, so this order comes from the - # client-side ordinal_position sort, and every missing-comment shape a - # cursor can report has to arrive as None. + # client-side ordinal_position sort, and a column with no comment has to + # arrive as None for every cursor class. assert [column["name"] for column in columns] == ["col_int", "col_string", "dt"] assert [column["comment"] for column in columns] == ["identifier", None, None] assert [column["dialect_options"]["awsathena_partition"] for column in columns] == [ @@ -745,6 +864,42 @@ def record_query(params, **kwargs): types.String, ] + @pytest.mark.parametrize( + "engine", + [ + {"driver": "rest"}, + {"driver": "pandas"}, + {"driver": "pandas", "unload": "true"}, + {"driver": "arrow"}, + {"driver": "polars"}, + ], + indirect=True, + ) + def test_get_view_definition_across_cursor_types(self, engine): + engine, conn = engine + # Athena formats this definition with a blank line. Read through a + # DataFrame cursor those rows arrive as NaN, as an empty string, or not + # at all, so the definition came back corrupted or raised TypeError. + view_name = f"test_view_definition_{uuid.uuid4().hex[:8]}" + conn.execute( + text( + f"CREATE OR REPLACE VIEW {ENV.schema}.{view_name} AS " + "WITH t AS (SELECT 1 AS a, 'x' AS b) " + "SELECT a, b FROM t UNION ALL SELECT 2, 'y'" + ) + ) + try: + definition = sqlalchemy.inspect(conn).get_view_definition(view_name, schema=ENV.schema) + finally: + conn.execute(text(f"DROP VIEW IF EXISTS {ENV.schema}.{view_name}")) + + lines = definition.splitlines() + assert lines[0].startswith("CREATE VIEW") + assert "UNION ALL" in definition + assert [i for i, line in enumerate(lines) if not line.strip()], ( + f"blank line lost from the definition: {definition!r}" + ) + def test_char_length(self, engine): engine, conn = engine one_row_complex = Table("one_row_complex", MetaData(schema=ENV.schema), autoload_with=conn) From da9a336f16fbc1bd3d74e757a55628f733325b18 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 21:55:43 +0900 Subject: [PATCH 02/10] Ask information_schema instead of propagating an unrecognized error The compliance suite pinned "an unwrapped MetadataException propagates" as a contract. That contract is what breaks has_table() on a federated catalog, which reports a missing table in its connector's words with no Glue envelope. Replace that case with one that asserts the fallback resolves it, for both a present and an absent table, and keep the recognized Glue codes propagating. Co-Authored-By: Claude Opus 5 --- tests/sqlalchemy/test_suite.py | 38 +++++++++++++++++++++++++++++----- 1 file changed, 33 insertions(+), 5 deletions(-) diff --git a/tests/sqlalchemy/test_suite.py b/tests/sqlalchemy/test_suite.py index 75c1fab3..fbfcc503 100644 --- a/tests/sqlalchemy/test_suite.py +++ b/tests/sqlalchemy/test_suite.py @@ -1038,9 +1038,7 @@ def record_query(params, **kwargs): inspector.get_table_options(name) assert caught.value.__cause__ is error - @sa_testing.combinations( - "AccessDeniedException", "InternalServerException", None, argnames="code" - ) + @sa_testing.combinations("AccessDeniedException", "InternalServerException", argnames="code") def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, code): raw_connection = _raw_connection(connection) retried = code == "InternalServerException" @@ -1054,8 +1052,6 @@ def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, message = ( "Catalog error (Service: AmazonDataCatalog; Status Code: 400; " f"Error Code: {code}; Request ID: example; Proxy: null)" - if code - else "Table not found" ) error = _metadata_error("MetadataException", message) calls = _fail_get_table_metadata( @@ -1068,6 +1064,38 @@ def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, assert caught.value.__cause__ is error assert len(calls) == (4 if retried else 2) + @sa_testing.combinations(True, False, argnames="exists") + def test_unrecognized_metadata_error_asks_information_schema( + self, connection, metadata, monkeypatch, exists + ): + # A federated catalog reports a missing table in its connector's own + # words, with no Glue error envelope to unwrap. Guessing absence from an + # unrecognized message is what turned throttling into false absence, so + # the question goes to information_schema in the same catalog instead. + name = "unrecognized_present" if exists else "unrecognized_absent" + if exists: + Table(name, metadata, Column("id", Integer)).create(connection) + raw_connection = _raw_connection(connection) + error = _metadata_error( + "MetadataException", + "Failed to invoke lambda function due to " + "com.amazonaws.services.lambda.invoke.LambdaFunctionException: " + "Requested resource not found " + "(Service: DynamoDb, Status Code: 400, Request ID: example)", + ) + # Retries are not shortened: an unrecognized code must not wait for them. + _fail_get_table_metadata(monkeypatch, raw_connection, error, attempt=None) + + inspector = inspect(connection) + assert inspector.has_table(name) is exists + if exists: + assert [column["name"] for column in inspector.get_columns(name)] == ["id"] + # Table options still need the metadata API and still report the failure. + monkeypatch.setattr(raw_connection.retry_config, "attempt", 1) + with pytest.raises(OperationalError) as caught: + inspector.get_table_options(name) + assert caught.value.__cause__ is error + @sa_testing.combinations((True, sa_testing.requires.schemas), False, argnames="use_schema") def test_has_table_cache_drop(self, connection, metadata, use_schema): schema = sa_testing.config.test_schema if use_schema else None From c61c3611b841c26b7434adac04bdfba503f24095 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 22:31:25 +0900 Subject: [PATCH 03/10] Let the async adapter accept a cursor class by name The shared dialect now asks for the API cursor class positionally, which the async adapted connection did not accept, and it could not drive a synchronous Cursor anyway. Accept the argument and map it to the async counterpart. Co-Authored-By: Claude Opus 5 --- pyathena/aio/sqlalchemy/base.py | 11 +++++++++-- pyathena/sqlalchemy/base.py | 11 ++++++----- 2 files changed, 15 insertions(+), 7 deletions(-) diff --git a/pyathena/aio/sqlalchemy/base.py b/pyathena/aio/sqlalchemy/base.py index ad2cc05a..e7b1d464 100644 --- a/pyathena/aio/sqlalchemy/base.py +++ b/pyathena/aio/sqlalchemy/base.py @@ -10,6 +10,8 @@ import pyathena from pyathena.aio.connection import AioConnection +from pyathena.aio.cursor import AioCursor +from pyathena.cursor import Cursor from pyathena.error import ( DatabaseError, DataError, @@ -30,6 +32,9 @@ from sqlalchemy import URL +_ASYNC_CURSOR_CLASSES: dict[Any, Any] = {Cursor: AioCursor} + + class AsyncAdapt_pyathena_cursor: """Wraps any async PyAthena cursor with a sync DBAPI interface. @@ -147,8 +152,10 @@ def cursor_kwargs(self) -> dict[str, Any]: def retry_config(self) -> RetryConfig: return self._connection.retry_config # type: ignore[no-any-return] - def cursor(self, **kwargs: Any) -> AsyncAdapt_pyathena_cursor: - raw_cursor = self._connection.cursor(**kwargs) + def cursor(self, cursor: Any = None, **kwargs: Any) -> AsyncAdapt_pyathena_cursor: + # The shared dialect names a cursor class in its synchronous form; this + # connection can only drive the async counterpart. + raw_cursor = self._connection.cursor(_ASYNC_CURSOR_CLASSES.get(cursor, cursor), **kwargs) return AsyncAdapt_pyathena_cursor(raw_cursor) def close(self) -> None: diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 1abb122d..784a4017 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -400,18 +400,19 @@ def _without_fallback_retries(cls, retry_config: RetryConfig) -> RetryConfig: ) @staticmethod - def _internal_cursor(raw_connection: PoolProxiedConnection) -> Cursor: + def _internal_cursor(raw_connection: PoolProxiedConnection) -> Any: """Open an API cursor for the queries this dialect parses itself. Reflection reads these rows directly, so they must not arrive in the result format chosen for user queries: a DataFrame cursor reports a NULL or blank value as NaN, as an empty string, or as a dropped row depending on its backend and on UNLOAD. + + The async connection adapter maps ``Cursor`` to its own counterpart and + returns its wrapper, so this is typed by the interface used here rather + than by the class requested. """ - return cast( - Cursor, - raw_connection.driver_connection.cursor(Cursor), # type: ignore[union-attr] - ) + return raw_connection.driver_connection.cursor(Cursor) # type: ignore[union-attr] def _column(self, name: str | None, type_: str, comment: str | None, partition: bool | None): return { From 82d56489189193348685376939c19e8943b625ea Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 23:09:45 +0900 Subject: [PATCH 04/10] Keep the Glue Data Catalog out of the unrecognized-error fallback Review found that an unrecognized MetadataException can also carry a permission failure: _get_error_code only unwraps the AmazonDataCatalog envelope, and information_schema filters by Lake Formation rather than erroring, so answering from it would report a table the caller cannot see as absent -- the same false absence this work exists to prevent. Glue states missing tables and permission failures in an envelope this client recognizes, so an unrecognized one there keeps propagating. Only a catalog that is not awsdatacatalog reaches the fallback, which is where the connector's own wording makes recognition impossible. Also pin the converter on the internal cursor, so a connection-level one chosen for a DataFrame cursor cannot reintroduce NaN comments, and cover a failed StartQueryExecution in get_view_definition: it raises DatabaseError, the parent of OperationalError, and must not be reported as a missing view. Co-Authored-By: Claude Opus 5 --- pyathena/sqlalchemy/base.py | 27 ++++++++++- tests/pyathena/sqlalchemy/test_base.py | 62 +++++++++++++++++++++----- 2 files changed, 77 insertions(+), 12 deletions(-) diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 784a4017..05682b4d 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -370,7 +370,7 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * metadata = self._lookup_table(cursor, schema, name, table_name) except pyathena.error.OperationalError as e: code = _get_error_code(e.__cause__ or e, unwrap_metadata=True) - if code not in self._FALLBACK_ERROR_CODES: + if not self._answerable_from_information_schema(code, catalog): raise _logger.warning( f"Table metadata request for {table_name} failed with {code}; " @@ -384,6 +384,24 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * info_cache[metadata_key] = metadata return self._columns_from_metadata(metadata) + @staticmethod + def _answerable_from_information_schema(code: str | None, catalog: str | None) -> bool: + """Whether a failed metadata request should be re-asked of the catalog. + + Throttling always: one query costs less than the retry ladder. + + A ``MetadataException`` that survived unwrapping only outside the Glue + Data Catalog. Glue states missing tables and permission failures in an + envelope this client recognizes, so an unrecognized one there has an + unknown cause, and answering it from ``information_schema`` would report + a table the caller merely cannot see as absent. A federated catalog has + no such envelope: it reports a missing table in its connector's own + words, which cannot be recognized at all. + """ + if code in THROTTLING_ERROR_CODES: + return True + return code == "MetadataException" and (catalog or "").lower() != "awsdatacatalog" + @classmethod def _without_fallback_retries(cls, retry_config: RetryConfig) -> RetryConfig: """Copy a policy without the codes the information_schema fallback answers. @@ -408,11 +426,16 @@ def _internal_cursor(raw_connection: PoolProxiedConnection) -> Any: or blank value as NaN, as an empty string, or as a dropped row depending on its backend and on UNLOAD. + The converter is pinned too: a connection-level one chosen for a + DataFrame cursor would otherwise be applied to this one. + The async connection adapter maps ``Cursor`` to its own counterpart and returns its wrapper, so this is typed by the interface used here rather than by the class requested. """ - return raw_connection.driver_connection.cursor(Cursor) # type: ignore[union-attr] + return raw_connection.driver_connection.cursor( # type: ignore[union-attr] + Cursor, converter=Cursor.get_default_converter() + ) def _column(self, name: str | None, type_: str, comment: str | None, partition: bool | None): return { diff --git a/tests/pyathena/sqlalchemy/test_base.py b/tests/pyathena/sqlalchemy/test_base.py index fd80b5ca..e4b4c76f 100644 --- a/tests/pyathena/sqlalchemy/test_base.py +++ b/tests/pyathena/sqlalchemy/test_base.py @@ -19,8 +19,9 @@ from sqlalchemy.sql.schema import Column, MetaData, Table from sqlalchemy.sql.selectable import TextualSelect +from pyathena.converter import DefaultTypeConverter from pyathena.cursor import Cursor -from pyathena.error import OperationalError +from pyathena.error import DatabaseError, OperationalError from pyathena.sqlalchemy.base import AthenaDialect from pyathena.sqlalchemy.types import ( TINYINT, @@ -71,8 +72,8 @@ def execute(operation, **kwargs): cursor = SimpleNamespace(execute=execute, fetchall=lambda: rows) - def open_cursor(cursor_class): - opened.append(cursor_class) + def open_cursor(cursor_class, converter=None): + opened.append((cursor_class, type(converter))) return contextlib.nullcontext(cursor) raw_connection = SimpleNamespace(driver_connection=SimpleNamespace(cursor=open_cursor)) @@ -83,7 +84,7 @@ def open_cursor(cursor_class): # The dialect parses these rows itself, so it asks for an API cursor # rather than whatever result format the user configured. - assert opened == [Cursor] + assert opened == [(Cursor, DefaultTypeConverter)] assert [column["name"] for column in columns] == ["id", "payload", "label", "dt"] assert isinstance(columns[0]["type"], types.INTEGER) @@ -253,6 +254,16 @@ def get_table_metadata(table_name, **kwargs): ) connection = SimpleNamespace(connection=raw_connection) + # In the Glue Data Catalog the same error propagates instead: Glue does + # state a missing table and a permission failure in a recognized + # envelope, so an unrecognized one there has an unknown cause, and + # information_schema hides a table the caller cannot see rather than + # erroring on it. + raw_connection.catalog_name = "AwsDataCatalog" + with pytest.raises(OperationalError): + AthenaDialect()._get_columns(connection, "events") + raw_connection.catalog_name = "federated_catalog" + if expected is None: with pytest.raises(NoSuchTableError): AthenaDialect()._get_columns(connection, "events") @@ -268,8 +279,8 @@ def test_get_view_definition_keeps_blank_lines(self): executed = [] opened = [] - def open_cursor(cursor_class): - opened.append(cursor_class) + def open_cursor(cursor_class, converter=None): + opened.append((cursor_class, type(converter))) return contextlib.nullcontext( SimpleNamespace( execute=lambda operation, **kwargs: executed.append(operation), @@ -287,10 +298,40 @@ def open_cursor(cursor_class): definition = AthenaDialect().get_view_definition(connection, "v") assert definition == "CREATE VIEW v AS\n\nSELECT 1" - assert opened == [Cursor] + assert opened == [(Cursor, DefaultTypeConverter)] assert executed == ['SHOW CREATE VIEW "default"."v";'] + def test_get_view_definition_propagates_a_failed_request(self): + # A request that never ran is not a missing view. BaseCursor raises + # DatabaseError for a failed StartQueryExecution, and DatabaseError is + # the parent of OperationalError, so it is not caught and not reported + # as absence. + error = ClientError( + {"Error": {"Code": "ThrottlingException", "Message": "Rate exceeded"}}, + "StartQueryExecution", + ) + + def execute(operation, **kwargs): + raise DatabaseError(*error.args) from error + + raw_connection = SimpleNamespace( + cursor_kwargs={}, + schema_name="default", + driver_connection=SimpleNamespace( + cursor=lambda *args, **kwargs: contextlib.nullcontext( + SimpleNamespace(execute=execute, fetchall=list) + ) + ), + ) + connection = SimpleNamespace(connection=raw_connection) + + with pytest.raises(DatabaseError): + AthenaDialect().get_view_definition(connection, "v") + def test_get_view_definition_reports_a_missing_view(self): + # Athena rejects SHOW CREATE VIEW for a view that does not exist, which + # reaches the dialect as OperationalError; verified live for the rest + # and pandas dialects. error = ClientError( {"Error": {"Code": "InvalidRequestException", "Message": "does not exist"}}, "StartQueryExecution", @@ -896,9 +937,10 @@ def test_get_view_definition_across_cursor_types(self, engine): lines = definition.splitlines() assert lines[0].startswith("CREATE VIEW") assert "UNION ALL" in definition - assert [i for i, line in enumerate(lines) if not line.strip()], ( - f"blank line lost from the definition: {definition!r}" - ) + if not any(not line.strip() for line in lines): + # The blank line is Athena's formatting, not the dialect's, so its + # absence means this case can no longer reach the defect. + pytest.skip(f"Athena formatted this view without a blank line: {definition!r}") def test_char_length(self, engine): engine, conn = engine From a3fbf4df1725237278d11f32009308e16dc14e6c Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 23:12:47 +0900 Subject: [PATCH 05/10] Restore the propagation case for an unrecognized Glue error This suite runs against the Glue Data Catalog, where an unrecognized MetadataException now keeps propagating, so the case removed earlier belongs here after all -- with a permission message, which is the reason it must not be re-asked of information_schema. The federated case it was replaced with has no catalog in this suite and is covered by the stub-cursor test instead. Co-Authored-By: Claude Opus 5 --- tests/sqlalchemy/test_suite.py | 46 ++++++++++------------------------ 1 file changed, 13 insertions(+), 33 deletions(-) diff --git a/tests/sqlalchemy/test_suite.py b/tests/sqlalchemy/test_suite.py index fbfcc503..2af77f67 100644 --- a/tests/sqlalchemy/test_suite.py +++ b/tests/sqlalchemy/test_suite.py @@ -1038,8 +1038,18 @@ def record_query(params, **kwargs): inspector.get_table_options(name) assert caught.value.__cause__ is error - @sa_testing.combinations("AccessDeniedException", "InternalServerException", argnames="code") + @sa_testing.combinations( + "AccessDeniedException", "InternalServerException", None, argnames="code" + ) def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, code): + # This suite runs against the Glue Data Catalog, which states missing + # tables and permission failures in an envelope this client recognizes. + # An unrecognized message there (code None) has an unknown cause, so it + # propagates rather than being re-asked of information_schema, which + # filters by Lake Formation instead of erroring and would report a table + # the caller cannot see as absent. Outside Glue the fallback does answer + # it; that case has no catalog here and is covered by + # TestAthenaDialect::test_unrecognized_metadata_error_asks_information_schema. raw_connection = _raw_connection(connection) retried = code == "InternalServerException" if retried: @@ -1052,6 +1062,8 @@ def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, message = ( "Catalog error (Service: AmazonDataCatalog; Status Code: 400; " f"Error Code: {code}; Request ID: example; Proxy: null)" + if code + else "is not authorized to perform: glue:GetTable" ) error = _metadata_error("MetadataException", message) calls = _fail_get_table_metadata( @@ -1064,38 +1076,6 @@ def test_metadata_errors_do_not_establish_absence(self, connection, monkeypatch, assert caught.value.__cause__ is error assert len(calls) == (4 if retried else 2) - @sa_testing.combinations(True, False, argnames="exists") - def test_unrecognized_metadata_error_asks_information_schema( - self, connection, metadata, monkeypatch, exists - ): - # A federated catalog reports a missing table in its connector's own - # words, with no Glue error envelope to unwrap. Guessing absence from an - # unrecognized message is what turned throttling into false absence, so - # the question goes to information_schema in the same catalog instead. - name = "unrecognized_present" if exists else "unrecognized_absent" - if exists: - Table(name, metadata, Column("id", Integer)).create(connection) - raw_connection = _raw_connection(connection) - error = _metadata_error( - "MetadataException", - "Failed to invoke lambda function due to " - "com.amazonaws.services.lambda.invoke.LambdaFunctionException: " - "Requested resource not found " - "(Service: DynamoDb, Status Code: 400, Request ID: example)", - ) - # Retries are not shortened: an unrecognized code must not wait for them. - _fail_get_table_metadata(monkeypatch, raw_connection, error, attempt=None) - - inspector = inspect(connection) - assert inspector.has_table(name) is exists - if exists: - assert [column["name"] for column in inspector.get_columns(name)] == ["id"] - # Table options still need the metadata API and still report the failure. - monkeypatch.setattr(raw_connection.retry_config, "attempt", 1) - with pytest.raises(OperationalError) as caught: - inspector.get_table_options(name) - assert caught.value.__cause__ is error - @sa_testing.combinations((True, sa_testing.requires.schemas), False, argnames="use_schema") def test_has_table_cache_drop(self, connection, metadata, use_schema): schema = sa_testing.config.test_schema if use_schema else None From 8ecd6bf00e137bcb569c74f6ec25746786fa414f Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 23:27:28 +0900 Subject: [PATCH 06/10] State the catalog boundary of the fallback in the docs Co-Authored-By: Claude Opus 5 --- docs/sqlalchemy.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/docs/sqlalchemy.md b/docs/sqlalchemy.md index b3ae7cef..c0b11797 100644 --- a/docs/sqlalchemy.md +++ b/docs/sqlalchemy.md @@ -101,7 +101,9 @@ Table-metadata lookups propagate throttling and permission errors rather than re `has_table()` propagates permission failures, including access denied by Lake Formation, instead of returning or caching `False`. For failed metadata requests, only recognized `EntityNotFoundException` responses establish absence; unrecognized errors are propagated rather than guessed to mean a missing table. Column reflection and `has_table()` do not retry a table-metadata request that `information_schema` can answer; they read `information_schema.columns` instead, executed without query result reuse, and log a warning. -That covers a throttled request, and a `MetadataException` carrying no recognized Glue error envelope, which is how a federated catalog reports a missing table: absence is then decided by the query against the same catalog rather than by an unrecognized message. +That covers a throttled request in any catalog. +It also covers a `MetadataException` carrying no recognized Glue error envelope, but only outside `AwsDataCatalog`: a federated catalog reports a missing table in its connector's own words, so absence is decided by the query against that catalog rather than by an unrecognized message. +In `AwsDataCatalog` an unrecognized `MetadataException` still propagates, because Glue does state missing tables and permission failures in a recognized envelope, and `information_schema` filters by Lake Formation instead of failing, so reading it there would report a table the caller cannot see as absent. Other error codes listed in the connection's `RetryConfig.exceptions` are still retried on that path, except those; list the wrapped Glue codes instead of `MetadataException`. A `retry_config` in `cursor_kwargs` replaces that policy entirely, including those retries, which then run before the fallback. The fallback maps unbounded `varchar` to SQLAlchemy `String`, matching Hive `STRING` reflection from the metadata API, and preserves explicit `VARCHAR(n)` and `CHAR(n)` lengths. From 205c9b14d4e152b39c810609455852f22d38512d Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 20 Sep 2026 23:56:24 +0900 Subject: [PATCH 07/10] Repair three findings from the independent review - get_view_definition() caught result-paging failures as a missing view. A GetQueryResults page that exhausts its retries raises OperationalError from the same block, so an existing view was reported absent. Only the rejected query establishes absence now; the fetch runs outside that handler. - Restore the comment guard in the information_schema fallback. Pinning the converter on the internal cursor is not sufficient: Connection.cursor() applies cursor_kwargs after the explicit arguments, so a converter supplied there still reaches it and can report a missing value as NaN. - The blank-line test could skip on the very defect it guards: losing every blank row still left CREATE VIEW and UNION ALL, so both assertions passed and the skip blamed Athena's formatting. Compare the whole definition against the rows Athena returned through an API cursor instead, and skip only when those rows carry no blank line. Co-Authored-By: Claude Opus 5 --- pyathena/sqlalchemy/base.py | 17 +++++++++++------ tests/pyathena/sqlalchemy/test_base.py | 22 ++++++++++++++-------- 2 files changed, 25 insertions(+), 14 deletions(-) diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 05682b4d..33a2191f 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -471,12 +471,15 @@ def _columns_from_information_schema( ) rows = cursor.fetchall() # Sort here: the query has no ORDER BY, so its result order is Athena's. + # The comment is still normalized at this boundary: a converter given in + # cursor_kwargs is applied after the one _internal_cursor() pins, and one + # written for a DataFrame cursor reports a missing value as NaN. return [ self._column( column_name, # Athena exposes Hive STRING as unbounded VARCHAR in information_schema. "string" if data_type == "varchar" else data_type, - comment, + comment if isinstance(comment, str) else None, extra_info == "partition key" or None, ) for _, column_name, data_type, comment, extra_info in sorted( @@ -563,12 +566,14 @@ def get_view_definition( raw_connection = self._raw_connection(connection) schema = schema if schema else self._cursor_option(raw_connection, "schema_name") query = f"""SHOW CREATE VIEW "{schema}"."{view_name}";""" - try: - with self._internal_cursor(raw_connection) as cursor: + with self._internal_cursor(raw_connection) as cursor: + try: cursor.execute(query) - rows = cursor.fetchall() - except pyathena.error.OperationalError as e: - raise exc.NoSuchTableError(f"{schema}.{view_name}") from e + except pyathena.error.OperationalError as e: + # Only a rejected query says the view is absent. A failure while + # paging the results is a failed read of a view that does exist. + raise exc.NoSuchTableError(f"{schema}.{view_name}") from e + rows = cursor.fetchall() # Athena returns the definition one line per row and blank lines as # empty values, which are part of the definition. return "\n".join(row[0] or "" for row in rows) diff --git a/tests/pyathena/sqlalchemy/test_base.py b/tests/pyathena/sqlalchemy/test_base.py index e4b4c76f..9df03a0b 100644 --- a/tests/pyathena/sqlalchemy/test_base.py +++ b/tests/pyathena/sqlalchemy/test_base.py @@ -56,10 +56,12 @@ def unique_s3tables_table_name(base: str) -> str: class TestAthenaDialect: def test_columns_from_information_schema(self): - # Rows arrive unordered, and Athena reports a missing comment as NULL, - # which an API cursor hands over as None or as an empty string. + # Rows arrive unordered, and Athena reports a missing comment as NULL. + # An API cursor hands that over as None or as an empty string; a + # converter supplied in cursor_kwargs is applied after the one this path + # pins, and one written for a DataFrame cursor reports it as NaN. rows = [ - ("4", "dt", "varchar", None, "partition key"), + ("4", "dt", "varchar", float("nan"), "partition key"), ("1", "id", "integer", "identifier", None), ("2", "payload", "row(a integer, b array(varchar))", None, None), ("3", "label", "varchar", "", ""), @@ -929,18 +931,22 @@ def test_get_view_definition_across_cursor_types(self, engine): "SELECT a, b FROM t UNION ALL SELECT 2, 'y'" ) ) + raw_connection = conn.connection.driver_connection try: definition = sqlalchemy.inspect(conn).get_view_definition(view_name, schema=ENV.schema) + # What Athena actually returned, row by row, independent of the + # dialect. Comparing against this catches a partial loss too. + with raw_connection.cursor(Cursor) as cursor: + cursor.execute(f'SHOW CREATE VIEW "{ENV.schema}"."{view_name}";') + rows = [row[0] for row in cursor.fetchall()] finally: conn.execute(text(f"DROP VIEW IF EXISTS {ENV.schema}.{view_name}")) - lines = definition.splitlines() - assert lines[0].startswith("CREATE VIEW") - assert "UNION ALL" in definition - if not any(not line.strip() for line in lines): + if not any(row is None or not row.strip() for row in rows): # The blank line is Athena's formatting, not the dialect's, so its # absence means this case can no longer reach the defect. - pytest.skip(f"Athena formatted this view without a blank line: {definition!r}") + pytest.skip(f"Athena formatted this view without a blank line: {rows!r}") + assert definition == "\n".join(row or "" for row in rows) def test_char_length(self, engine): engine, conn = engine From 185f58212aba330a384492057a07a89ac5f14acd Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Wed, 23 Sep 2026 12:14:37 +0900 Subject: [PATCH 08/10] Tell a missing view from a failed first result page Cursor.execute() prefetches the first result page, so narrowing the handler to execute() still caught a GetQueryResults failure as a missing view. Measured live, a missing view is a query Athena runs and fails ("View not found or not a valid presto view"), which the cursor raises as OperationalError with no underlying API error; an API failure carries its ClientError as the cause. Only the former is absence now. The stub for a missing view modelled a rejected StartQueryExecution, which never reaches this handler (that raises DatabaseError). Model the failed query instead, cover the failed first page, and restore the content checks in the live test, since its baseline goes through the same API cursor path. Co-Authored-By: Claude Opus 5.5 --- pyathena/sqlalchemy/base.py | 9 +++-- tests/pyathena/sqlalchemy/test_base.py | 50 ++++++++++++++------------ 2 files changed, 35 insertions(+), 24 deletions(-) diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 33a2191f..8add6b38 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -570,8 +570,13 @@ def get_view_definition( try: cursor.execute(query) except pyathena.error.OperationalError as e: - # Only a rejected query says the view is absent. A failure while - # paging the results is a failed read of a view that does exist. + # Athena runs SHOW CREATE VIEW for a missing view and fails the + # query, which the cursor reports without an underlying API + # error. execute() also fetches the first result page, and a + # failed API call there carries its error as the cause: that is + # a failed read of a view that exists, not a missing one. + if e.__cause__ is not None: + raise raise exc.NoSuchTableError(f"{schema}.{view_name}") from e rows = cursor.fetchall() # Athena returns the definition one line per row and blank lines as diff --git a/tests/pyathena/sqlalchemy/test_base.py b/tests/pyathena/sqlalchemy/test_base.py index 9df03a0b..1d5ddd53 100644 --- a/tests/pyathena/sqlalchemy/test_base.py +++ b/tests/pyathena/sqlalchemy/test_base.py @@ -316,6 +316,11 @@ def test_get_view_definition_propagates_a_failed_request(self): def execute(operation, **kwargs): raise DatabaseError(*error.args) from error + with pytest.raises(DatabaseError): + AthenaDialect().get_view_definition(self._view_connection(execute), "v") + + @staticmethod + def _view_connection(execute): raw_connection = SimpleNamespace( cursor_kwargs={}, schema_name="default", @@ -325,36 +330,33 @@ def execute(operation, **kwargs): ) ), ) - connection = SimpleNamespace(connection=raw_connection) - - with pytest.raises(DatabaseError): - AthenaDialect().get_view_definition(connection, "v") + return SimpleNamespace(connection=raw_connection) def test_get_view_definition_reports_a_missing_view(self): - # Athena rejects SHOW CREATE VIEW for a view that does not exist, which - # reaches the dialect as OperationalError; verified live for the rest - # and pandas dialects. + # Athena accepts SHOW CREATE VIEW for a view that does not exist and + # fails the query; the cursor reports the failure reason with no + # underlying API error. Measured live for the rest and pandas dialects. + def execute(operation, **kwargs): + raise OperationalError("View not found or not a valid presto view: v") + + with pytest.raises(NoSuchTableError): + AthenaDialect().get_view_definition(self._view_connection(execute), "v") + + def test_get_view_definition_propagates_a_failed_result_page(self): + # execute() also fetches the first result page. A GetQueryResults call + # that exhausts its retries there is a failed read of an existing view, + # and must not be reported as absence. error = ClientError( - {"Error": {"Code": "InvalidRequestException", "Message": "does not exist"}}, - "StartQueryExecution", + {"Error": {"Code": "ThrottlingException", "Message": "Rate exceeded"}}, + "GetQueryResults", ) def execute(operation, **kwargs): raise OperationalError(*error.args) from error - raw_connection = SimpleNamespace( - cursor_kwargs={}, - schema_name="default", - driver_connection=SimpleNamespace( - cursor=lambda *args, **kwargs: contextlib.nullcontext( - SimpleNamespace(execute=execute, fetchall=list) - ) - ), - ) - connection = SimpleNamespace(connection=raw_connection) - - with pytest.raises(NoSuchTableError): - AthenaDialect().get_view_definition(connection, "v") + with pytest.raises(OperationalError) as caught: + AthenaDialect().get_view_definition(self._view_connection(execute), "v") + assert caught.value.__cause__ is error def test_get_table_matches_long_names_case_insensitively(self): # GetTableMetadata rejects names over 128 characters, so the lookup lists @@ -942,6 +944,10 @@ def test_get_view_definition_across_cursor_types(self, engine): finally: conn.execute(text(f"DROP VIEW IF EXISTS {ENV.schema}.{view_name}")) + # Content the baseline cannot vouch for itself, since it is read through + # the same API cursor path the dialect now uses. + assert definition.startswith("CREATE VIEW") + assert "UNION ALL" in definition if not any(row is None or not row.strip() for row in rows): # The blank line is Athena's formatting, not the dialect's, so its # absence means this case can no longer reach the defect. From bde5075cab51718d02ea2310d66b65d0bc3ef4f5 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Wed, 23 Sep 2026 12:22:34 +0900 Subject: [PATCH 09/10] State what the view-absence handler still cannot tell apart A cancelled query, or one that failed for another reason, also raises OperationalError without a cause and is still reported as a missing view. That mapping predates this branch, which only stopped API failures from taking it. The cursor does not keep the failed query's state, and a missing view reports ErrorCategory 1 / ErrorType 1502, which does not single it out, so the comment says so rather than implying a precise test. Co-Authored-By: Claude Opus 5.5 --- pyathena/sqlalchemy/base.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 8add6b38..7a1b643d 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -574,7 +574,10 @@ def get_view_definition( # query, which the cursor reports without an underlying API # error. execute() also fetches the first result page, and a # failed API call there carries its error as the cause: that is - # a failed read of a view that exists, not a missing one. + # a failed read of a view that exists, not a missing one. Any + # query that ends without success is still read as absence, as + # before; its state is not kept and its error codes do not + # single out a missing view. if e.__cause__ is not None: raise raise exc.NoSuchTableError(f"{schema}.{view_name}") from e From 8e906772d57b0894ff0ff6159dc63e40ca7c3d2f Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Wed, 23 Sep 2026 12:48:19 +0900 Subject: [PATCH 10/10] Rename the fallback predicate to _is_fallback_error Match _FALLBACK_ERROR_CODES and _without_fallback_retries, and say in the docstring that the catalog, not only the code, decides the outcome. Co-Authored-By: Claude Opus 5.5 --- pyathena/sqlalchemy/base.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/pyathena/sqlalchemy/base.py b/pyathena/sqlalchemy/base.py index 7a1b643d..12e02511 100644 --- a/pyathena/sqlalchemy/base.py +++ b/pyathena/sqlalchemy/base.py @@ -370,7 +370,7 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * metadata = self._lookup_table(cursor, schema, name, table_name) except pyathena.error.OperationalError as e: code = _get_error_code(e.__cause__ or e, unwrap_metadata=True) - if not self._answerable_from_information_schema(code, catalog): + if not self._is_fallback_error(code, catalog): raise _logger.warning( f"Table metadata request for {table_name} failed with {code}; " @@ -385,8 +385,11 @@ def _get_columns(self, connection, table_name: str, schema: str | None = None, * return self._columns_from_metadata(metadata) @staticmethod - def _answerable_from_information_schema(code: str | None, catalog: str | None) -> bool: - """Whether a failed metadata request should be re-asked of the catalog. + def _is_fallback_error(code: str | None, catalog: str | None) -> bool: + """Whether the information_schema fallback answers this failed request. + + The codes are those in ``_FALLBACK_ERROR_CODES``; the catalog decides + whether an unrecognized ``MetadataException`` qualifies. Throttling always: one query costs less than the retry ladder.