Skip to content
Merged
48 changes: 22 additions & 26 deletions docs/sqlalchemy.md
Original file line number Diff line number Diff line change
Expand Up @@ -1365,60 +1365,56 @@ events = Table('events', metadata,

#### Querying JSON data

When querying JSON data, PyAthena automatically parses JSON strings into Python dictionaries:
PyAthena's default converters decode results of Athena's `json` type, and the `JSON` type returns those values unchanged.
A JSON object becomes a `dict`, an array a `list`, and a string scalar a `str`.
If the cursor returns the text of a `json` result instead, for example through a custom converter that does not decode it, the `JSON` type returns that text unchanged.

```python
from sqlalchemy import select, literal_column
from sqlalchemy.sql import type_coerce
from sqlalchemy.types import JSON

# Query with explicit type coercion
result = connection.execute(
select(
type_coerce(
literal_column('CAST(\'{"name": "test", "value": 123}\' AS JSON)'),
literal_column('json_parse(\'{"name": "test", "items": [1, 2, 3]}\')'),
JSON
).label("json_col")
)
).fetchone()

# Result is automatically parsed as a dictionary
print(result.json_col) # {"name": "test", "value": 123}
print(result.json_col) # {'items': [1, 2, 3], 'name': 'test'}
print(type(result.json_col)) # <class 'dict'>
```

#### Important limitations

Athena's JSON type support has specific limitations:

- **JSON objects and arrays are supported** - `CAST('...' AS JSON)` accepts an object or a top-level array such as `[1, 2, 3]`
- **Arrays within objects are supported** - JSON objects can contain arrays as property values
- **DML only** - JSON type is supported for SELECT queries but not in CREATE TABLE statements; compiling `CREATE TABLE` with a `JSON` column raises `CompileError`
`json_parse()` parses text into a JSON value, while `CAST('...' AS JSON)` returns the text as a JSON string scalar:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2 (claims, compatibility, operations) — FINDINGS (PR text only, corrected)

Scope: git diff 4e7b55ff700edfc92d51baa219e92add9783fcb4..13d2c0e0e25eda6c433736c37e85bf1f3a79d287; claims in the PR body, commit messages, the AthenaJSON docstrings, and docs/sqlalchemy.md.

Claims checked:

  • "converters decode the Athena json type": the default, pandas, Arrow, and Polars mappings use _to_json, and S3FS copies _DEFAULT_CONVERTERS. Measured through SQLAlchemy with all five drivers: description = json/varchar, values dict/list/str/None. pandas and Arrow were re-measured without the SQL NULL column because of the pre-existing NULL defect recorded in round 1. UNLOAD rejects json (NOT_SUPPORTED), so it never reaches this path.
  • Docs: json_parse() vs CAST('...' AS JSON) and varchar decoding: the examples were run live on rest. The first example's output comment uses Athena's key order (items before name).
  • Other docs: the CAST(ROW/MAP/ARRAY ... AS JSON) examples (docs/sqlalchemy.md:878, :1026, :1300; docs/usage.md:708) use text() or cursors without the JSON type and cast structured values, so they are unaffected.
  • Existing callers: SQLAlchemy >= 2.0 keeps the result_processor(dialect, coltype) signature. A bound value still serializes through the inherited bind_processor. The CAST(text AS JSON) → str change is release-noted, and the PR now has the 4.0.0 milestone and the bug label.

Corrections to the PR body:

  • "reflected json view columns" and "JSON view columns" were not measured. Narrowed to the reflected json → JSON mapping and the measured expressions.
  • Added the compatibility consequence for a custom converter that returns raw text for json (it now gets that text back).
  • Tested commit updated to 13d2c0e; the live runs used fc1ddb4 (it differs only in the offline helper).

Operations: the live JSON test now issues 1 Athena query instead of 4. No retry or quota behavior changed.

Limits: async dialect not tested live (it inherits colspecs); the full pyathena/sqla/sqla-async suites are left to CI after Ready. just docs build passed locally.


```python
# Supported: JSON object with nested array
result = connection.execute(
select(
type_coerce(
literal_column('CAST(\'{"items": [1, 2, 3]}\' AS JSON)'),
JSON
).label("json_col")
type_coerce(literal_column("json_parse('[1, 2, 3]')"), JSON).label("parsed"),
type_coerce(literal_column("CAST('[1, 2, 3]' AS JSON)"), JSON).label("cast"),
)
).fetchone()
print(result.json_col) # {"items": [1, 2, 3]}

# Supported: Top-level array
print(result.parsed) # [1, 2, 3]
print(result.cast) # '[1, 2, 3]'
```

The `JSON` type decodes JSON text in columns of other types, such as a `varchar` column, with `json.loads()` or the dialect's `json_deserializer`:

```python
result = connection.execute(
select(
type_coerce(
literal_column("CAST('[1, 2, 3]' AS JSON)"),
JSON
).label("json_col")
)
select(type_coerce(literal_column("'{\"a\": 1}'"), JSON).label("json_col"))
).fetchone()
print(result.json_col) # [1, 2, 3]

print(result.json_col) # {'a': 1}
```

#### Important limitations

- **DML only** - JSON type is supported for SELECT queries but not in CREATE TABLE statements; compiling `CREATE TABLE` with a `JSON` column raises `CompileError`

#### Best practices

1. **Use with SELECT queries** - JSON type works best for querying existing data
Expand Down
4 changes: 2 additions & 2 deletions pyathena/arrow/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,11 @@

from pyathena.converter import (
Converter,
_csv_to_json,
_to_binary,
_to_date,
_to_decimal,
_to_default,
_to_json,
_to_time,
)
from pyathena.util import override
Expand All @@ -26,7 +26,7 @@
"time": _to_time,
"decimal": _to_decimal,
"varbinary": _to_binary,
"json": _to_json,
"json": _csv_to_json,
}


Expand Down
22 changes: 14 additions & 8 deletions pyathena/arrow/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@

from pyathena import OperationalError
from pyathena.arrow.util import to_column_info
from pyathena.converter import Converter
from pyathena.converter import Converter, _json_text_converter, _to_default
from pyathena.error import ProgrammingError
from pyathena.model import AthenaQueryExecution
from pyathena.result_set import AthenaResultSet
Expand Down Expand Up @@ -146,8 +146,8 @@ def __init__(
import pyarrow as pa

self._table = pa.Table.from_pydict({})
# The fetch methods convert only the values read from a result file.
# GetQueryResults values are already converted.
# The fetch methods convert the values read from a result file. GetQueryResults
# values are already converted, except json values, which stay text.
self._convert_rows = bool(self.output_location)
self._batches = iter(self._table.to_batches(arraysize))

Expand Down Expand Up @@ -254,11 +254,16 @@ def _fetch(self) -> None:
return
else:
dict_rows = rows.to_pydict()
if self._convert_rows:
converters = self.converters
converters = (
self.converters if self._convert_rows else self._json_converters(self.converters)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review after the rebase onto master fe21250 (both perspectives): CLEAN

The branch conflicted with #1023 (#1021), which fixed the ArrowCursor fallback double conversion independently. I rebased onto fe21250a3e228e3d38e61898334c5baacc1285c2. Scope: git range-diff 16aef64d10dce92cbf8b8d957c178bef9bd494b6..4eb5b0b202d1df3000042a9c4a648ac7b8f47c5c fe21250a3e228e3d38e61898334c5baacc1285c2..d5eac4c78ef6fb7c2f009137d9dd537087598487, plus the upstream changes to the touched files (#1023 Arrow, #1012 and #1029 Polars).

Resolution:

Upstream checks: #1012 changes the chunked Polars UNLOAD metadata and #1029 changes PolarsDataFrameIterator.close(); neither touches the GetQueryResults path or _df_converters. Codex's earlier chunked-UNLOAD finding (#1007) is fixed by #1012.

Validation on d5eac4c: just lint passed; offline 148 passed; live Arrow/Polars/S3FS sync+aio suites 315 passed, 1 skipped; pandas subset 74 passed; test_json_type passed. With this PR's cursor changes reverted on fe21250, five test_fetch_all_rows cases fail (Arrow default and managed, pandas default, Polars managed, S3FS custom-converter managed); with them, all 10 pass. The PR body and #1028 now credit #1023 for the Arrow fallback change.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up after the rebase (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a10239-6234-7140-af28-4bdf226e9084. The snapshot was at d5eac4c78ef6fb7c2f009137d9dd537087598487, and the reviewer got the range-diff from 4eb5b0b and the final Arrow/Polars/result-set/converter diff against fe21250a3e228e3d38e61898334c5baacc1285c2. Static review only; the snapshot and the PR worktree were confirmed unchanged.

Covered: sync, async, and aio Arrow/Polars cursors across CSV, UNLOAD, and managed fallback paths, including Polars chunking; Arrow's fetch-time converter lookup (a converter.set() after execute()); managed JSON kept as text and decoded once while other fallback values pass through; heterogeneous JSON, string scalars, and NULL json; the rewritten Arrow test.

Verdict: CLEAN. The conflict resolution keeps both the #1023 intent and this change's intent.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review of df476e4 (both perspectives): CLEAN

Maintainer question: the S3FS test helper had the same name as the production pyathena.converter._json_text_converter(), but builds a different thing. Scope: git diff d5eac4c78ef6fb7c2f009137d9dd537087598487..df476e484c3c0c8208288c54d27761cecf6f2662 (test file only).

  • tests/pyathena/s3fs/test_cursor.py: _json_text_converter() → _s3fs_converter_with_json_text(), which builds a DefaultS3FSTypeConverter (the S3FS cursor's converter type) with json → _to_default. The production helper returns a DefaultTypeConverter for the GetQueryResults fallback; the test does not reuse it, so it does not depend on a private helper meant for another path.
  • DefaultS3FSTypeConverter.__init__ deep-copies _DEFAULT_CONVERTERS, so set() changes only that instance; there is no shared-dict mutation.
  • No source change. just lint passed; pytest -n 2 tests/pyathena/s3fs/test_cursor.py -k test_fetch_all_rows: 4 passed (default/managed, default/custom converter). PR body tested commit updated.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up on df476e4 (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a10269-23f8-78a2-8741-ad8297b5dee3. The snapshot was at df476e484c3c0c8208288c54d27761cecf6f2662, and the reviewer got the incremental diff from d5eac4c. Static review only; the snapshot and the PR worktree were confirmed unchanged.

Covered: the incremental diff, helper references, fixture forwarding, and the S3 CSV and managed API conversion paths.

Verdict: CLEAN. The rename is complete (definition and both call sites). Both cases still pass a custom DefaultS3FSTypeConverter that keeps JSON text, and the unchanged assertion checks it. The managed case keeps its skip when AWS_ATHENA_MANAGED_WORKGROUP is unset.

)
if converters:
column_names = dict_rows.keys()
processed_rows = [
tuple(converters[k](v) for k, v in zip(column_names, row, strict=False))
tuple(
converters.get(k, _to_default)(v)
for k, v in zip(column_names, row, strict=False)
)
for row in zip(*dict_rows.values(), strict=False)
]
else:
Expand Down Expand Up @@ -381,11 +386,12 @@ def _as_arrow_from_api(self, converter: Converter | None = None) -> Table:

Args:
converter: Type converter for result values. Defaults to
``DefaultTypeConverter`` if not specified.
``DefaultTypeConverter`` with json values kept as text, as in
the CSV result file.
"""
import pyarrow as pa

rows = self._fetch_all_rows(converter)
rows = self._fetch_all_rows(converter or _json_text_converter())
if not rows:
return pa.Table.from_pydict({})
description = self.description if self.description else []
Expand Down
26 changes: 26 additions & 0 deletions pyathena/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,21 @@ def _to_json(varchar_value: str | None) -> Any | None:
return json.loads(varchar_value)


def _csv_to_json(varchar_value: str | None) -> Any | None:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review of the expanded scope (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a101f7-f1dd-7353-ab60-dba323d2576c. The snapshot was at c0ebcaaebc83eeaa7b635b753ef063e6fd8e441f (merge-base 16aef64d10dce92cbf8b8d957c178bef9bd494b6), and the reviewer got the full diff with a neutral prompt (no PR text, commits, or prior findings). Static review only. Afterwards the snapshot and the PR worktree were confirmed unchanged.

Covered (reviewer): SQLAlchemy JSON adaptation, TypeDecorator/Variant, nested ARRAY/MAP/STRUCT, reflection, sync/aio dialects, shared converters, and pandas/Arrow/Polars/S3FS CSV, UNLOAD, managed-storage, chunk, and DataFrame paths including async callers; tests and changed docs.

Verdict: FINDINGS

  1. P2, introduced: pyathena/converter.py:131 (_to_json returning None for "") turns nested empty JSON strings into NULL: DefaultTypeConverter().convert("array", '[""]', type_hint="array(json)") → [None] instead of [""].
  2. P2, introduced: pyathena/arrow/result_set.py:155: with managed storage, a custom Arrow converter for a non-json type (e.g. varchar → str.upper) no longer applies ("abc" instead of "ABC").
  3. P2, pre-existing: pyathena/arrow/result_set.py:325: single-column SQL NULL rows are dropped from CSV results (ignore_empty_lines).
  4. P2, pre-existing: pyathena/pandas/result_set.py:849: with managed storage, DataFrame inference turns a JSON column of numbers + NULL into float64 (9007199254740993 → 9007199254740992.0).
  5. P2, pre-existing: pyathena/pandas/result_set.py:844: with managed storage, a configured pandas json converter is ignored.
  6. P2, pre-existing: pyathena/polars/result_set.py:367: chunked Polars UNLOAD uses the initial metadata names, so row values are missing.
    The reviewer judged the added tests to assert observable behavior and to expose the double decoding and the mixed-shape failures.

Author verification and disposition

Self-review of repair c5d78c7 (both perspectives): CLEAN. _to_json is byte-identical to master again, so the REST cursor, the typed parser, and Polars are unchanged from master. The pandas pyarrow-engine check uses self.converters, which tests mapping membership, not function identity (pyathena/pandas/result_set.py:405, :509), so the engine choice is unchanged. On the Arrow managed path, _csv_to_json receives None for NULL and the text "" for a JSON empty-string scalar, so values are unchanged. No docs reference _to_json (only build artifacts). test_fetch_all_rows for all four cursors: 10 passed on c5d78c7; just lint passed.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up on c5d78c7 (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a10209-9aa8-7873-8919-6782c2b4725f. The snapshot was at c5d78c7556c295640924cd4f8ad3bc3bb08d35f1, and the reviewer got the incremental diff from c0ebcaa. Static review only; the snapshot and the PR worktree were confirmed unchanged.

Covered: nested array(json) / map(varchar,json) conversion; the pandas and Arrow CSV paths across sync, async, and aio cursors, including pandas chunks and Arrow batches; the Arrow managed results; the added regression tests and directly affected callers.

Verdict: CLEAN. The nested empty-string regression is fixed, empty CSV fields stay SQL NULL, a JSON "" stays an empty string, and Arrow managed JSON text is decoded once with None kept as None. No new actionable defects.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review of refactor 4eb5b0b (both perspectives): CLEAN

Maintainer request: the result sets should only choose converters. Scope: git diff c5d78c7556c295640924cd4f8ad3bc3bb08d35f1..4eb5b0b202d1df3000042a9c4a648ac7b8f47c5c (4 files).

  • The json-text converter (a DefaultTypeConverter with json → _to_default) moves from the AthenaResultSet._json_text_converter() staticmethod to pyathena.converter._json_text_converter(). Arrow and Polars call it from _as_arrow_from_api() / _as_polars_from_api(), and each call still builds a fresh converter. AthenaResultSet._json_converters() stays in the result set because it picks columns from description.
  • Arrow's new _identity is removed; _table_converters.get(k, _to_default) uses _to_default, which returns its argument unchanged (pyathena/converter.py:477), so behavior is the same. The pre-existing Polars _identity is untouched.
  • There is no import cycle (pyathena.converter does not import result sets). Both helpers were added in this unreleased PR, so no caller depends on the old location.
  • Validation: just lint passed; offline test_converter.py, arrow/test_converter.py, and sqlalchemy/test_types.py (148 passed); live tests/pyathena/arrow tests/pyathena/aio/arrow tests/pyathena/polars tests/pyathena/aio/polars (212 passed). The PR body now names the new location.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up on 4eb5b0b (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a1022a-9e31-7610-a5a5-40bc2467440b. The snapshot was at 4eb5b0b202d1df3000042a9c4a648ac7b8f47c5c, and the reviewer got the incremental diff from c5d78c7. Static review only; the snapshot and the PR worktree were confirmed unchanged.

Covered: sync, async/Future, and aio Arrow and Polars cursors; CSV, UNLOAD/Parquet, the managed-result fallback, and Polars chunked iteration; JSON/NULL handling, explicit converter precedence, converter isolation, and Arrow identity conversion; helper references, import dependencies, and callable/return typing.

Verdict: CLEAN. The relocated helper keeps its implementation, both callers are updated, _to_default returns values unchanged and fits the converter mapping, and there is no new import cycle or typing problem.

"""Convert an Athena JSON value read from a CSV result file.

Args:
varchar_value: The value as JSON text, or None. An empty string, which
CSV results use for SQL NULL, is also treated as NULL.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2, expanded scope (claims, compatibility, operations): CLEAN

Scope: git diff 16aef64d10dce92cbf8b8d957c178bef9bd494b6..c0ebcaaebc83eeaa7b635b753ef063e6fd8e441f, plus the updated PR title/body and the commit messages of 5792a2e and c0ebcaa.

Claims checked:

  • "An empty string, which CSV results use for SQL NULL" (here): docs/s3fs.md:216-218 and the measurement agree (unquoted empty field = NULL). JSON text is never empty; a JSON empty-string scalar is "".
  • "PolarsCursor got the same fix in Keep the whole pandas and Polars results reusable #1005": 5a60dd1 sets _df_converters = {} for the GetQueryResults path.
  • PR table and the heterogeneous-shape failure: measured on master 16aef64 before the fix (ArrowInvalid: cannot mix list and non-list, non-null values; Polars TypeError: unexpected value while building Series of type Int64); after the fix, Arrow and Polars return the same rows on both paths as Cursor.
  • "Defaults to DefaultTypeConverter with json values kept as text, as in the CSV result file" (_as_arrow_from_api, _as_polars_from_api): the CSV readers type json as pa.string() / pl.String (docs/arrow.md:159, docs/polars.md:178), and the tests now assert a string column on both paths.
  • _fetch_all_rows docstring (subclass converters may not handle API strings) stays true; only S3FS, whose converter reads text like DefaultTypeConverter, now passes its own.

Compatibility (all release-noted in the PR body, 4.0.0 milestone):

  • With managed storage, the json columns of as_arrow() / as_polars() / Polars iter_chunks() change from inferred struct/list columns to string columns; code that used struct operations on them must decode instead. The fetch methods return the same values as before for shapes that worked.
  • With managed storage, an S3FS custom converter now applies; the default converter gives the same values as DefaultTypeConverter there.
  • _to_json("") returns None instead of raising; no previously working input changes.

Docs: docs/s3fs.md, docs/arrow.md, docs/polars.md, docs/pandas.md, docs/usage.md (managed storage), and docs/aio.md:223 contain no statement that this change makes false.

Operations: the cursor tests add columns and a UNION ALL row to existing queries. The new S3FS test adds 2 queries, and test_json_type saves 3.

Limits: the full pyathena / sqla / sqla-async suites run in CI after Ready; locally: the Arrow/Polars/S3FS sync+aio suites (432 passed) and the pandas subset (106 passed, before c0ebcaa, which does not touch pandas).


Returns:
The decoded value, or None for SQL NULL.
"""
if not varchar_value:
return None
return json.loads(varchar_value)


def _to_array(varchar_value: str | None) -> list[Any] | str | None:
"""Convert array data to Python list.

Expand Down Expand Up @@ -727,3 +742,14 @@ def _parse_type_hint(self, type_hint: str) -> TypeNode:
if normalized not in self._parsed_hints:
self._parsed_hints[normalized] = self._parser.parse(normalized)
return self._parsed_hints[normalized]


def _json_text_converter() -> DefaultTypeConverter:
"""Return a ``DefaultTypeConverter`` that keeps json values as text.

Returns:
The converter.
"""
converter = DefaultTypeConverter()
converter.set("json", _to_default)
return converter
4 changes: 2 additions & 2 deletions pyathena/pandas/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,11 @@

from pyathena.converter import (
Converter,
_csv_to_json,
_to_binary,
_to_boolean,
_to_decimal,
_to_default,
_to_json,
)
from pyathena.util import override

Expand All @@ -24,7 +24,7 @@
"boolean": _to_boolean,
"decimal": _to_decimal,
"varbinary": _to_binary,
"json": _to_json,
"json": _csv_to_json,
}


Expand Down
11 changes: 7 additions & 4 deletions pyathena/polars/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
)

from pyathena import OperationalError
from pyathena.converter import Converter
from pyathena.converter import Converter, _json_text_converter
from pyathena.error import ProgrammingError
from pyathena.model import AthenaQueryExecution
from pyathena.polars.util import to_column_info
Expand Down Expand Up @@ -262,7 +262,7 @@ def __init__(
# Note: _as_polars() and _create_dataframe_iterator() update _metadata for unload
# queries, so the converters and column names must be read AFTER them.
self._df: pl.DataFrame | None = None
# Converters for the rows of self._df. GetQueryResults values are already converted.
# Converters for the rows of self._df.
self._df_converters: dict[str, Callable[[str | None], Any | None]] = {}
if self.state == AthenaQueryExecution.STATE_SUCCEEDED and self.output_location:
if self._chunksize is None:
Expand All @@ -272,6 +272,8 @@ def __init__(
self._df_iter = self._create_dataframe_iterator()
elif self.state == AthenaQueryExecution.STATE_SUCCEEDED:
self._df = self._as_polars_from_api()
# GetQueryResults values are already converted, except json values kept as text.
self._df_converters = self._json_converters(self.converters)
else:
self._df = pl.DataFrame()
if self._df is not None:
Expand Down Expand Up @@ -539,11 +541,12 @@ def _as_polars_from_api(self, converter: Converter | None = None) -> pl.DataFram

Args:
converter: Type converter for result values. Defaults to
``DefaultTypeConverter`` if not specified.
``DefaultTypeConverter`` with json values kept as text, as in
the CSV result file.
"""
import polars as pl

rows = self._fetch_all_rows(converter)
rows = self._fetch_all_rows(converter or _json_text_converter())
if not rows:
return pl.DataFrame()
description = self.description if self.description else []
Expand Down
16 changes: 16 additions & 0 deletions pyathena/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
from pyathena.util import RetryConfig, override, parse_output_location, retry_api_call

if TYPE_CHECKING:
from collections.abc import Callable

from pyathena.connection import Connection

_logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -679,6 +681,20 @@ def _is_first_row_column_labels(self, rows: list[dict[str, Any]]) -> bool:
return False
return True

def _json_converters(
self, converters: dict[str, Callable[[str | None], Any | None]]
) -> dict[str, Callable[[str | None], Any | None]]:
"""Select the converters of the json columns.

Args:
converters: The converters keyed by column name.

Returns:
The converters of the columns whose Athena type is json.
"""
description = self.description if self.description else []
return {d[0]: converters[d[0]] for d in description if d[1] == "json"}

def _fetch_all_rows(
self,
converter: Converter | None = None,
Expand Down
5 changes: 3 additions & 2 deletions pyathena/s3fs/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -120,8 +120,9 @@ def __init__(
if self.state == AthenaQueryExecution.STATE_SUCCEEDED and self.output_location:
self._init_csv_reader()
elif self.state == AthenaQueryExecution.STATE_SUCCEEDED:
# Managed query result storage: no output_location, use API
rows = self._fetch_all_rows()
# Managed query result storage: no output_location, use API.
# The converter reads text values, as from the CSV result file.
rows = self._fetch_all_rows(self._converter)
self._rows.extend(rows)

# If CSV reader was not initialized (e.g., CTAS, DDL),
Expand Down
2 changes: 2 additions & 0 deletions pyathena/sqlalchemy/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
AthenaArray,
AthenaBinary,
AthenaDate,
AthenaJSON,
AthenaMap,
AthenaStruct,
AthenaTimestamp,
Expand Down Expand Up @@ -212,6 +213,7 @@ class AthenaDialect(DefaultDialect):
types.ARRAY: AthenaArray,
types.Date: AthenaDate,
types.DateTime: AthenaTimestamp,
types.JSON: AthenaJSON,
}

ischema_names: dict[str, type[Any]] = ischema_names
Expand Down
36 changes: 34 additions & 2 deletions pyathena/sqlalchemy/types.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@

from __future__ import annotations

from typing import TYPE_CHECKING
from typing import TYPE_CHECKING, Any

from sqlalchemy import types
from sqlalchemy.sql import sqltypes
Expand All @@ -19,7 +19,7 @@

if TYPE_CHECKING:
from sqlalchemy import Dialect
from sqlalchemy.sql.type_api import _LiteralProcessorType
from sqlalchemy.sql.type_api import _LiteralProcessorType, _ResultProcessorType

__all__ = [
"ARRAY",
Expand All @@ -29,6 +29,7 @@
"AthenaArray",
"AthenaBinary",
"AthenaDate",
"AthenaJSON",
"AthenaMap",
"AthenaStruct",
"AthenaTimestamp",
Expand All @@ -47,6 +48,37 @@ def process(value: bytes) -> str:
return process


class AthenaJSON(types.JSON):
"""SQLAlchemy JSON type that keeps the values PyAthena has decoded.

PyAthena's default converters decode results of the Athena ``json`` type,
so this type returns them unchanged, and a JSON string scalar stays a
``str``. If the cursor returns the text of a ``json`` result instead, for
example through a custom converter that does not decode it, this type
returns that text unchanged. Results of other Athena types, such as JSON
text in a ``varchar`` column, are decoded with the dialect's JSON
deserializer.
"""

@override
def result_processor(
self, dialect: Dialect, coltype: object
) -> _ResultProcessorType[Any] | None:
"""Return a processor decoding JSON text of Athena types other than ``json``.

Args:
dialect: The dialect fetching the value.
coltype: The Athena type name from the cursor description.

Returns:
The processor, or None for the Athena ``json`` type.
"""
if coltype == "json":

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1 (implementation behavior) — FINDINGS (1, repaired)

Scope: git diff 4e7b55ff700edfc92d51baa219e92add9783fcb4..fc1ddb4e4cc2f0ae82a47b6c0d16ddb171744de6 (all 5 files).

Covered:

  • Behavior: the coltype SQLAlchemy passes comes from cursor.description. Measured with all five drivers, it is json for JSON expressions and varchar for text, so json results pass through and text is decoded as before. None stays None on both paths.
  • Callers: AthenaAioDialect inherits colspecs (AthenaAioDialect.colspecs[types.JSON] is AthenaJSON). A TypeDecorator with impl = JSON resolves its impl through colspecs and gets the same processor (checked through dialect_impl(...).result_processor). STRUCT/MAP processors do not touch JSON, and the ARRAY path decodes JSON elements in _ArrayValueProcessor._decode without a result processor, so both are unchanged.
  • Bind side: bind_processor is inherited unchanged, so a bound dict is sent as JSON text (varchar) and decoded on fetch.
  • Framework contract: like every colspec, adapt_type now replaces a user subclass of types.JSON with AthenaJSON, which is the standard SQLAlchemy behavior (for example psycopg2 maps JSON to its own type); custom result handling belongs in a TypeDecorator, which keeps working.
  • Tests: the offline tests fail 13/24 and the live test_json_type fails with the SQLAlchemy JSON columns fail on JSON objects and arrays returned by Athena #934 TypeError when the pyathena/ changes are reverted.

Finding: tests/pyathena/sqlalchemy/test_types.py:23 called SQLAlchemy's private _cached_result_processor; the public dialect_impl(...).result_processor(...) exercises the same path for both JSON and TypeDecorator.
Repair: 13d2c0e switches to the public method; the tests still pass with the fix and fail 13/24 without it.

Out of scope (pre-existing, not introduced here): PandasCursor and ArrowCursor fail on SQL NULL of the json type (CAST(NULL AS JSON) → JSONDecodeError from _to_json(''); one Arrow run returned no row). The live test therefore runs the SQL NULL case on the default rest engine only.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed): Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max (from the user config), codex exec -s read-only --ephemeral, session 01a10177-84e0-7ec0-9890-dc85396aa8fc.
The reviewer worked from a detached snapshot at 13d2c0e0e25eda6c433736c37e85bf1f3a79d287 (merge-base 4e7b55ff700edfc92d51baa219e92add9783fcb4) and received only the literal diff and a neutral prompt; it did not see the PR number, description, commits, or prior findings. Static review only: no edits, builds, tests, or network access. Afterwards the snapshot and the PR worktree were confirmed unchanged.

Covered (reviewer): all five sync and five asyncio dialects; shared cursor conversion, CSV and managed-result paths; SQLAlchemy 2.0.46 colspecs, TypeDecorator, with_variant(), nested ARRAY/MAP/STRUCT processing and reflection; changed tests, docs, and docstrings.

Verdict: FINDINGS

  1. P2, introduced: raw custom converters bypass JSON deserialization (pyathena/sqlalchemy/types.py:73). With DefaultTypeConverter().remove("json") passed through connect_args, json_parse('{"a":1}') read through JSON used to return {"a": 1}; it now returns the raw string, and a configured json_deserializer is skipped. The reviewer also said the blanket converter claim needs qualification.
  2. P2, pre-existing: with managed result storage and no output_location, Arrow and Polars decode JSON a second time (pyathena/arrow/result_set.py:256, pyathena/polars/result_set.py:138). _fetch_all_rows() decodes json with DefaultTypeConverter, and the fetch then applies _to_json to the dict, raising TypeError before SQLAlchemy receives the value.
  3. P2, pre-existing: pandas CSV results raise on SQL NULL json (pyathena/pandas/converter.py:27, pyathena/pandas/result_set.py:605), because _to_json("") reaches json.loads("").

The reviewer also judged the tests to assert observable values and to fail for the original defect, found no other defect in adaptation, reflection, or the typed ARRAY path, and confirmed the json_parse() vs CAST and DDL statements as accurate.

Author verification and disposition

  • (1) Confirmed. The behavior is intentional and already listed as a release-note item in the PR body. For an Athena json value that is a str, the value alone cannot tell a decoded JSON string scalar from undecoded text, so returning string scalars as str (the expected behavior in SQLAlchemy JSON columns fail on JSON objects and arrays returned by Athena #934) requires relying on the converter. Repair 9906bd3: the docstring and docs/sqlalchemy.md now say "default converters" and state that a custom converter that does not decode json gets its text back. The implementation is unchanged.
  • (2) Confirmed statically: ArrowResultSet._fetch and AthenaPolarsResultSet.iterrows apply the cursor's json converter to values built by _fetch_all_rows() with DefaultTypeConverter (pyathena/result_set.py:682-731). This is a cursor-level defect independent of this diff, and the SQLAlchemy layer never receives the value. Not measured (no managed-results work group used here). Deferred to a separate issue.
  • (3) Confirmed and measured live (also recorded in self-review round 1): CAST(NULL AS JSON) raises OperationalError from JSONDecodeError with PandasCursor (and ArrowCursor CSV). Pre-existing; deferred to a separate issue.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review of repair 9906bd3 (both perspectives) — CLEAN

Scope: git range-diff 4e7b55ff700edfc92d51baa219e92add9783fcb4..13d2c0e0e25eda6c433736c37e85bf1f3a79d287 4e7b55ff700edfc92d51baa219e92add9783fcb4..9906bd32b01fdde2bbc360d735814ea32df7b464, which adds one commit and changes only the AthenaJSON docstring and docs/sqlalchemy.md prose.

  • Behavior: no code change. just lint, just docs lint, and the offline test_types.py (24 passed) were rerun on 9906bd3.
  • Claims: "PyAthena's default converters decode json" holds because the default, pandas, Arrow, and Polars mappings use _to_json and S3FS copies _DEFAULT_CONVERTERS. The pre-existing NULL and managed-result failures in (2) and (3) are raised in the cursors before decoding completes, so they do not contradict it. "A custom converter that does not decode json gets its text back" follows from coltype == "json" returning no processor. The PR body now also states that json_deserializer no longer applies to json results, and the tested commit was updated.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a101c4-5925-7ce0-a409-4d1b66be085b. The snapshot was at 9906bd32b01fdde2bbc360d735814ea32df7b464, and the reviewer received the range-diff from 13d2c0e plus the incremental diff. Static review only. Covered: all ten SQLAlchemy drivers (sync + aio), converter selection, row fetching, managed-result fallback, UNLOAD variants, AthenaJSON processing.

Verdict: FINDINGS

  • P2, docs/sqlalchemy.md:1370 and pyathena/sqlalchemy/types.py:56: "custom converter → text" overstates the guarantee. With awsathena+s3fs and managed result storage, pyathena/s3fs/result_set.py:124 calls _fetch_all_rows() without the cursor's converter, so DefaultTypeConverter (pyathena/result_set.py:713) decodes json to a dict even when a custom converter is configured. The reviewer suggested describing the actual boundary instead: when the cursor returns raw JSON text, JSON passes it through.

Author verification and repair f2cda29: confirmed at pyathena/s3fs/result_set.py:124. The docstring and docs now say: if the cursor returns the text of a json result, for example through a custom converter that does not decode it, JSON returns that text unchanged. The PR body uses the same boundary. No code change; just lint and just docs lint pass.
Self-review of the repair (both perspectives): the new wording claims only what coltype == "json" → no processor guarantees (the value is passed through), so it holds on every fetch path, including the managed-result fallbacks. The S3FS fallback ignoring a custom converter is pre-existing cursor behavior outside this diff; it is added to the deferred cursor-level issues with (2) and (3).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up 2 (relayed): Codex CLI 0.160.0, gpt-6-astra, codex exec -s read-only --ephemeral, session 01a101c9-e825-7820-bca0-f943a5910c3f. The snapshot was at f2cda2981ee7d5b2becfae2199eea945ee367b3c, and the reviewer received the incremental diff from 9906bd3. Static review only. Covered: all ten SQLAlchemy drivers, converter selection, cursor metadata, CSV/UNLOAD and managed-result paths, AthenaJSON, and "Querying JSON data".

Verdict: FINDINGS. The reviewer considers the previous custom-converter finding resolved.

  • P2, docs/sqlalchemy.md:1369 and pyathena/sqlalchemy/types.py:55: with managed result storage, Arrow and Polars (sync and aio) return CAST('[1, 2, 3]' AS JSON) as a list, and raise TypeError for objects and arrays, because _fetch_all_rows() decodes through DefaultTypeConverter (pyathena/result_set.py:713) and the fetch applies _to_json again (pyathena/arrow/result_set.py:256, pyathena/polars/result_set.py:138). The reviewer suggests qualifying the result-shape claims for these fallback paths.

Author disposition: not applied to this PR (pre-existing, deferred). This is pre-existing defect (2) from the first review: the cursor decodes json twice before SQLAlchemy sees the value, and AthenaJSON preserves whatever the cursor returns, as the reviewer notes. The same cursors return the wrong shape for these values without SQLAlchemy too. Following the project's practice, the docs describe the intended behavior rather than a known cursor defect, and the fix belongs in the Arrow/Polars managed-result fetch path. It is deferred to a cursor-level issue together with the pandas/Arrow SQL NULL json failure and the S3FS managed fallback ignoring a custom converter.

return None
processor: _ResultProcessorType[Any] | None = super().result_processor(dialect, coltype)
return processor


class Tinyint(sqltypes.Integer):
"""SQLAlchemy type for Athena TINYINT (8-bit signed integer).

Expand Down
9 changes: 8 additions & 1 deletion tests/pyathena/arrow/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -977,8 +977,15 @@ def test_fetch_all_rows(self, arrow_cursor):
,X'0102' AS col_varbinary
,json_parse('{"a": 1}') AS col_json
,CAST('{"a": 1}' AS JSON) AS col_json_string
,CAST(NULL AS JSON) AS col_json_null
UNION ALL
SELECT
2, CAST('12:34:56' AS TIME), X'0102', json_parse('[1, "x"]'), json_parse('"s"'), NULL
ORDER BY col
"""
)
assert arrow_cursor.fetchall() == [
(1, datetime(2017, 1, 1, 12, 34, 56).time(), b"\x01\x02", {"a": 1}, '{"a": 1}')
(1, datetime(2017, 1, 1, 12, 34, 56).time(), b"\x01\x02", {"a": 1}, '{"a": 1}', None),
(2, datetime(2017, 1, 1, 12, 34, 56).time(), b"\x01\x02", [1, "x"], "s", None),
]
assert arrow_cursor.as_arrow().schema.field("col_json").type == pa.string()
Loading
Loading