Skip to content

Commit 72649ac

Browse files
Preserve Hive string types in metadata fallback reflection
1 parent f7b2369 commit 72649ac

4 files changed

Lines changed: 20 additions & 6 deletions

File tree

‎docs/sqlalchemy.md‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -103,7 +103,8 @@ For failed metadata requests, only recognized `EntityNotFoundException` response
103103
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.
104104
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.
105105
A `retry_config` in `cursor_kwargs` replaces that policy entirely, including its throttling retries, which then run before the fallback.
106-
Columns reflected this way carry the type names Athena reports there, such as `varchar` for `string`, and partition columns are marked from its `extra_info` column.
106+
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.
107+
Partition columns are marked from the `extra_info` column.
107108
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.
108109
Athena applies its metadata API rate limits per account, and they are not listed in Service Quotas.
109110
PyAthena's API retries use exponential backoff with uniform jitter; `RetryConfig` documents the default attempt count and waits.

‎pyathena/sqlalchemy/base.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -425,7 +425,8 @@ def _columns_from_information_schema(
425425
return [
426426
self._column(
427427
column_name,
428-
data_type,
428+
# Athena exposes Hive STRING as unbounded VARCHAR in information_schema.
429+
"string" if data_type == "varchar" else data_type,
429430
comment if isinstance(comment, str) else None,
430431
extra_info == "partition key" or None,
431432
)

‎tests/pyathena/sqlalchemy/test_base.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ def execute(operation, **kwargs):
7878
assert isinstance(columns[0]["type"], types.INTEGER)
7979
assert columns[0]["comment"] == "identifier"
8080
assert isinstance(columns[1]["type"], AthenaStruct)
81-
assert isinstance(columns[2]["type"], types.VARCHAR)
81+
assert type(columns[2]["type"]) is types.String
8282
assert columns[2]["comment"] is None
8383
assert [column["dialect_options"]["awsathena_partition"] for column in columns] == [
8484
None,

‎tests/sqlalchemy/test_suite.py‎

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -240,7 +240,10 @@ def test_preserves_table_metadata_until_clear_cache(self, connection, metadata):
240240
assert [column["name"] for column in inspector.get_columns(table.name)] == ["id", "added"]
241241

242242
@sa_testing.combinations((String, None), (VARCHAR, 52), (CHAR, 52), argnames="type_,length")
243-
def test_hive_string_length_reflection(self, connection, metadata, type_, length):
243+
@sa_testing.combinations(False, True, argnames="information_schema")
244+
def test_hive_string_length_reflection(
245+
self, connection, metadata, type_, length, information_schema
246+
):
244247
table = Table(
245248
"string_length",
246249
metadata,
@@ -249,8 +252,17 @@ def test_hive_string_length_reflection(self, connection, metadata, type_, length
249252
awsathena_file_format="PARQUET",
250253
)
251254
table.create(connection)
252-
reflected_type = inspect(connection).get_columns(table.name)[0]["type"]
253-
assert isinstance(reflected_type, type_)
255+
if information_schema:
256+
# Exercise the fallback with a real query without forcing an API failure.
257+
dialect = connection.dialect
258+
raw_connection = dialect._raw_connection(connection)
259+
columns = dialect._columns_from_information_schema(
260+
raw_connection, raw_connection.schema_name, table.name
261+
)
262+
else:
263+
columns = inspect(connection).get_columns(table.name)
264+
reflected_type = columns[0]["type"]
265+
assert type(reflected_type) is type_
254266
# Generic String compiles to Hive STRING without a length constraint.
255267
assert reflected_type.length == length
256268

0 commit comments

Comments
 (0)