-
Notifications
You must be signed in to change notification settings - Fork 116
Return Athena JSON values from SQLAlchemy JSON columns and the pandas, Arrow, and S3FS cursors #1010
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Return Athena JSON values from SQLAlchemy JSON columns and the pandas, Arrow, and S3FS cursors #1010
Changes from all commits
5ea497c
428287b
b175da9
c5d9830
a96a81e
90537de
fe4d744
d5eac4c
df476e4
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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)) | ||
|
|
||
|
|
@@ -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) | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 Resolution:
Upstream checks: #1012 changes the chunked Polars UNLOAD metadata and #1029 changes Validation on d5eac4c:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up after the rebase (relayed): Codex CLI 0.160.0, 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 Verdict: CLEAN. The conflict resolution keeps both the #1023 intent and this change's intent.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up on df476e4 (relayed): Codex CLI 0.160.0, 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 |
||
| ) | ||
| 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: | ||
|
|
@@ -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 [] | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review of the expanded scope (relayed): Codex CLI 0.160.0, 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
Author verification and disposition
Self-review of repair c5d78c7 (both perspectives): CLEAN.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up on c5d78c7 (relayed): Codex CLI 0.160.0, Covered: nested Verdict: CLEAN. The nested empty-string regression is fixed, empty CSV fields stay SQL NULL, a JSON
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up on 4eb5b0b (relayed): Codex CLI 0.160.0, 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, |
||
| """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. | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round 2, expanded scope (claims, compatibility, operations): CLEAN Scope: Claims checked:
Compatibility (all release-noted in the PR body, 4.0.0 milestone):
Docs: Operations: the cursor tests add columns and a Limits: the full |
||
|
|
||
| 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. | ||
|
|
||
|
|
@@ -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 | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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", | ||
|
|
@@ -29,6 +29,7 @@ | |
| "AthenaArray", | ||
| "AthenaBinary", | ||
| "AthenaDate", | ||
| "AthenaJSON", | ||
| "AthenaMap", | ||
| "AthenaStruct", | ||
| "AthenaTimestamp", | ||
|
|
@@ -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": | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round 1 (implementation behavior) — FINDINGS (1, repaired) Scope: Covered:
Finding: Out of scope (pre-existing, not introduced here):
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed): Codex CLI 0.160.0, model Covered (reviewer): all five sync and five asyncio dialects; shared cursor conversion, CSV and managed-result paths; SQLAlchemy 2.0.46 Verdict: FINDINGS
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 Author verification and disposition
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review of repair 9906bd3 (both perspectives) — CLEAN Scope:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up (relayed): Codex CLI 0.160.0, Verdict: FINDINGS
Author verification and repair f2cda29: confirmed at
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up 2 (relayed): Codex CLI 0.160.0, Verdict: FINDINGS. The reviewer considers the previous custom-converter finding resolved.
Author disposition: not applied to this PR (pre-existing, deferred). This is pre-existing defect (2) from the first review: the cursor decodes |
||
| 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). | ||
|
|
||
|
|
||
There was a problem hiding this comment.
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, theAthenaJSONdocstrings, anddocs/sqlalchemy.md.Claims checked:
jsontype": 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, valuesdict/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 rejectsjson(NOT_SUPPORTED), so it never reaches this path.json_parse()vsCAST('...' AS JSON)and varchar decoding: the examples were run live on rest. The first example's output comment uses Athena's key order (itemsbeforename).CAST(ROW/MAP/ARRAY ... AS JSON)examples (docs/sqlalchemy.md:878,:1026,:1300;docs/usage.md:708) usetext()or cursors without theJSONtype and cast structured values, so they are unaffected.result_processor(dialect, coltype)signature. A bound value still serializes through the inheritedbind_processor. TheCAST(text AS JSON)→strchange is release-noted, and the PR now has the 4.0.0 milestone and thebuglabel.Corrections to the PR body:
jsonview columns" and "JSON view columns" were not measured. Narrowed to the reflectedjson→JSONmapping and the measured expressions.json(it now gets that text back).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 fullpyathena/sqla/sqla-asyncsuites are left to CI after Ready.just docs buildpassed locally.