-
Notifications
You must be signed in to change notification settings - Fork 116
Read managed query results as a CSV result file in the pandas, Arrow, and Polars cursors #1061
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
Changes from all commits
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 |
|---|---|---|
|
|
@@ -5,9 +5,10 @@ | |
| import logging | ||
| from collections.abc import Callable | ||
| from copy import deepcopy | ||
| from typing import Any | ||
| from typing import TYPE_CHECKING, Any | ||
|
|
||
| from pyathena.converter import ( | ||
| _TIMESTAMP_TEXT_LENGTHS, | ||
| Converter, | ||
| _to_binary, | ||
| _to_date, | ||
|
|
@@ -20,6 +21,9 @@ | |
| ) | ||
| from pyathena.util import override | ||
|
|
||
| if TYPE_CHECKING: | ||
| from pyarrow import ChunkedArray, TimestampType | ||
|
|
||
| _logger = logging.getLogger(__name__) | ||
|
|
||
|
|
||
|
|
@@ -34,6 +38,28 @@ | |
| } | ||
|
|
||
|
|
||
| def _to_timestamp(column: ChunkedArray, type_: TimestampType) -> ChunkedArray: | ||
|
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 d376f51 (move to the converter modules), rounds 1 and 2: CLEAN The maintainer asked to keep conversions out of the result sets, as in #1010. Range: 1878a8c..d376f51 (base a89a4f7 unchanged).
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 d376f51 (relayed): CLEAN Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high,
|
||
| """Convert timestamp text to a timestamp type, truncating finer fractions. | ||
|
|
||
| Athena writes up to 12 fractional digits, which pyarrow does not parse into a | ||
| timestamp type whose unit holds fewer. | ||
|
|
||
| Args: | ||
| column: The timestamp text, with NULL as null or as an empty string. | ||
| type_: The timestamp type. | ||
|
|
||
| Returns: | ||
| The timestamps. | ||
| """ | ||
| import pyarrow as pa | ||
| import pyarrow.compute as pc | ||
|
|
||
| length = _TIMESTAMP_TEXT_LENGTHS[type_.unit] | ||
| if (pc.max(pc.utf8_length(column)).as_py() or 0) > length: | ||
| column = pc.utf8_slice_codeunits(column, 0, length) | ||
| return pc.if_else(pc.equal(column, ""), pa.scalar(None, pa.string()), column).cast(type_) | ||
|
|
||
|
|
||
| class DefaultArrowTypeConverter(Converter): | ||
| """Optimized type converter for Apache Arrow Table results. | ||
|
|
||
|
|
@@ -86,7 +112,7 @@ def _dtypes(self) -> dict[str, type[Any]]: | |
| "char": pa.string(), | ||
| "varchar": pa.string(), | ||
| "string": pa.string(), | ||
| "timestamp": pa.timestamp("ms"), | ||
| "timestamp": pa.timestamp("us"), | ||
| "date": pa.timestamp("ms"), | ||
| "time": pa.string(), | ||
| "time with time zone": pa.string(), | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -11,9 +11,9 @@ | |
| ) | ||
|
|
||
| from pyathena import OperationalError | ||
| from pyathena.arrow.converter import _to_timestamp | ||
| from pyathena.arrow.util import to_column_info | ||
| from pyathena.converter import _TEXT_VALUE_TYPES, Converter, _text_value_converter, _to_default | ||
| from pyathena.error import ProgrammingError | ||
| from pyathena.converter import Converter, _to_default | ||
| from pyathena.model import AthenaQueryExecution | ||
| from pyathena.result_set import AthenaResultSet | ||
| from pyathena.util import RetryConfig, override, parse_output_location | ||
|
|
@@ -141,15 +141,13 @@ def __init__( | |
| if self.state == AthenaQueryExecution.STATE_SUCCEEDED and self.output_location: | ||
| self._table = self._as_arrow() | ||
| elif self.state == AthenaQueryExecution.STATE_SUCCEEDED: | ||
| self._table = self._as_arrow_from_api() | ||
| # Without a result file, as with managed query result storage, the rows from | ||
| # GetQueryResults are read as a CSV result file. | ||
| self._table = self._read_csv() | ||
|
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): FINDINGS (1) Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort high, Covered: pagination and the header row; CSV quoting and newlines; NULL vs the empty string; empty and no-column results; DML and DDL; the pandas, Arrow, and Polars dtypes and converters; binary, JSON, and temporal values; duplicate and odd names; read options; sync, threaded async, and aio callers; exceptions, cursor state, fetch and close, resource ownership; tests and docs.
The reviewer found that the added tests assert observable bytes, schemas, values, and custom conversion, and can fail against the old implementation. Author verification: confirmed. Master's managed Arrow path converted with
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. Repair: 46bf916 (the maintainer chose to fix this in this PR rather than defer it)
Self-review of the repair
An independent follow-up review of the repair is pending.
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 review (relayed): FINDINGS (2) Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high,
Author verification: both reproduce offline on Polars 1.44.2. The same happens with Repair: 301b22e. Polars now follows the Arrow design: the first read is unchanged, and only when Self-review of the repair:
A second independent follow-up on 46bf916..301b22e is pending.
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 review #2 (relayed): CLEAN Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, Covered: retry eligibility and
Author note on the gaps:
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 after rebase (relayed): FINDINGS (1) Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, Covered: API pagination; CSV quoting and newlines; NULL vs empty values; empty and no-column results; DML and DDL; the pandas, Arrow, and Polars types, converters, and timestamp handling; duplicate labels; reader kwargs; sync, Future-based async, and aio callers; exceptions, cursor state, fetch and close, resource ownership; S3-path performance; test sensitivity to the old implementation.
Author verification: confirmed from the source. It is the S3-path defect filed as #1051 (pandas
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. Repair: rebased onto master a89a4f7 (#1066, fixing #1051) → fd80f98, plus test 1878a8c. The maintainer chose to merge #1051 first and rebase this PR onto it.
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 review (relayed): CLEAN Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, Covered: managed CSV rendering and pagination; pandas label resolution, header skipping, and the in-memory binary stream; Arrow positional types, timestamp conversion, binary NULL handling, and renaming;
|
||
| else: | ||
| import pyarrow as pa | ||
|
|
||
| self._table = pa.Table.from_pydict({}) | ||
| # The fetch methods convert the values read from a result file. GetQueryResults | ||
| # values are already converted, except json and time with time zone values, | ||
| # which stay text. | ||
| self._convert_rows = bool(self.output_location) | ||
| self._batches = iter(self._table.to_batches(arraysize)) | ||
|
|
||
| def _create_s3_file_system(self): | ||
|
|
@@ -258,12 +256,7 @@ def _fetch(self) -> None: | |
| # converters property keep one column per name. | ||
| columns = [column.to_pylist() for column in rows.columns] | ||
| description = self.description if self.description else [] | ||
| converters = [ | ||
| self._converter.get(d[1]) | ||
| if self._convert_rows or d[1] in _TEXT_VALUE_TYPES | ||
| else _to_default | ||
| for d in description | ||
| ] | ||
| converters = [self._converter.get(d[1]) for d in description] | ||
| if any(convert is not _to_default for convert in converters): | ||
| processed_rows = [ | ||
| tuple(convert(v) for convert, v in zip(converters, row, strict=False)) | ||
|
|
@@ -287,12 +280,18 @@ def fetchone( | |
| return self._rows.popleft() | ||
|
|
||
| def _read_csv(self) -> Table: | ||
| """Read the CSV result file, or the GetQueryResults rows as one without it. | ||
|
|
||
| Returns: | ||
| The Arrow Table of the results. | ||
|
|
||
| Raises: | ||
| OperationalError: If reading the results fails. | ||
| """ | ||
| import pyarrow as pa | ||
| from pyarrow import csv | ||
|
|
||
| if not self.output_location: | ||
| raise ProgrammingError("OutputLocation is none or empty.") | ||
| if not self.output_location.endswith((".csv", ".txt")): | ||
| if self.output_location and not self.output_location.endswith((".csv", ".txt")): | ||
| return pa.Table.from_pydict({}) | ||
| if self.substatement_type and self.substatement_type.upper() in ( | ||
| "UPDATE", | ||
|
|
@@ -301,7 +300,14 @@ def _read_csv(self) -> Table: | |
| "VACUUM_TABLE", | ||
| ): | ||
| return pa.Table.from_pydict({}) | ||
| length = self._get_content_length() | ||
| if self.output_location: | ||
| data = None | ||
| length = self._get_content_length() | ||
| location = "/".join(parse_output_location(self.output_location)) | ||
| else: | ||
| data = self._fetch_all_rows_as_csv() | ||
| length = len(data) | ||
| location = "the GetQueryResults rows" | ||
| description = self.description if self.description else [] | ||
| names = [d[0] for d in description] | ||
| # pyarrow types every column with a name by its column_types entry, so columns | ||
|
|
@@ -318,8 +324,16 @@ def _read_csv(self) -> Table: | |
| else: | ||
| column_names = names | ||
| column_types = self.column_types | ||
| # Timestamp columns are read as text: pyarrow does not parse fractions finer | ||
| # than the unit, and on some platforms its strptime fallbacks accept such a | ||
| # value without its fraction. | ||
| timestamp_types = { | ||
| i: dtype | ||
| for i, (name, d) in enumerate(zip(column_names, description, strict=True)) | ||
| if d[1] == "timestamp" and isinstance(dtype := column_types.get(name), pa.TimestampType) | ||
| } | ||
| binary_columns = {i for i, d in enumerate(description) if d[1] == "varbinary"} | ||
| if length and self.output_location.endswith(".txt"): | ||
| if length and self.output_location and self.output_location.endswith(".txt"): | ||
| read_opts = csv.ReadOptions( | ||
| skip_rows=0, | ||
| column_names=column_names, | ||
|
|
@@ -332,7 +346,7 @@ def _read_csv(self) -> Table: | |
| double_quote=False, | ||
| escape_char=False, | ||
| ) | ||
| elif length and self.output_location.endswith(".csv"): | ||
| elif length: | ||
| read_opts = csv.ReadOptions(skip_rows=0, block_size=self._block_size, use_threads=True) | ||
| if has_duplicate_names: | ||
| read_opts.column_names = column_names | ||
|
|
@@ -353,19 +367,27 @@ def _read_csv(self) -> Table: | |
| else: | ||
| return pa.Table.from_pydict({}) | ||
|
|
||
| bucket, key = parse_output_location(self.output_location) | ||
| try: | ||
| table = csv.read_csv( | ||
| self._fs.open_input_stream(f"{bucket}/{key}"), | ||
| self._fs.open_input_stream(location) if data is None else pa.BufferReader(data), | ||
| read_options=read_opts, | ||
| parse_options=parse_opts, | ||
| convert_options=csv.ConvertOptions( | ||
| strings_can_be_null=bool(binary_columns), | ||
| quoted_strings_can_be_null=False, | ||
| timestamp_parsers=self.timestamp_parsers, | ||
| column_types=column_types, | ||
| column_types={ | ||
| **column_types, | ||
| **{column_names[i]: pa.string() for i in timestamp_types}, | ||
| }, | ||
| ), | ||
| ) | ||
| for index, type_ in timestamp_types.items(): | ||
| table = table.set_column( | ||
| index, | ||
| pa.field(table.schema.field(index).name, type_), | ||
| _to_timestamp(table.column(index), type_), | ||
| ) | ||
| if has_duplicate_names: | ||
| table = table.rename_columns(names) | ||
| if binary_columns: | ||
|
|
@@ -377,7 +399,7 @@ def _read_csv(self) -> Table: | |
| table = table.set_column(index, field, table.column(index).fill_null("")) | ||
| return table | ||
| except Exception as e: | ||
| _logger.exception(f"Failed to read {bucket}/{key}.") | ||
| _logger.exception(f"Failed to read {location}.") | ||
| raise OperationalError(*e.args) from e | ||
|
|
||
| def _read_parquet(self) -> Table: | ||
|
|
@@ -406,27 +428,6 @@ def _as_arrow(self) -> Table: | |
| table = self._read_csv() | ||
| return table | ||
|
|
||
| def _as_arrow_from_api(self, converter: Converter | None = None) -> Table: | ||
| """Build an Arrow Table from GetQueryResults API. | ||
|
|
||
| Used as a fallback when ``output_location`` is not available | ||
| (e.g. managed query result storage). | ||
|
|
||
| Args: | ||
| converter: Type converter for result values. Defaults to | ||
| ``DefaultTypeConverter`` with json and time with time zone values kept as | ||
| text, as in the CSV result file. Arrow has no type for JSON values or for | ||
| times with a time zone. | ||
| """ | ||
| import pyarrow as pa | ||
|
|
||
| rows = self._fetch_all_rows(converter or _text_value_converter()) | ||
| if not rows: | ||
| return pa.Table.from_pydict({}) | ||
| description = self.description if self.description else [] | ||
| columns = [list(column) for column in zip(*rows, strict=True)] | ||
| return pa.table(columns, names=[d[0] for d in description]) | ||
|
|
||
| def as_arrow(self) -> Table: | ||
| """Return the query results as an Apache Arrow Table. | ||
|
|
||
|
|
||
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, callers, operations): FINDINGS, fixed in the PR body
Base a57325b, head 03c9f3f. Claims came from the PR body, the commit message, docstrings and comments,
docs/usage.md, anddocs/arrow.md.Findings and corrections (PR body only, no code change):
dateas a timestamp anddecimalas a string.DefaultPolarsTypeConvertermapsdatetopl.Dateanddecimaltopl.Decimal(p, s), so that was false for Polars. The note now describes Arrow'scolumn_typesand Polars'schema_overridesseparately._fetch_all_rows()/_fetch_all_rows_as_csv(). Rerun on 03c9f3f: the combined targeted run, 190 passed, and the Arrow suite, 116 passed. The body now cites only those.Verified claims:
GetQueryResultstext equals the CSV file for 28 types: measured earlier on Athena.pyathena/aio/{pandas,arrow,polars}/cursor.pyconstruct the sameAthena*ResultSetclasses.S3FSCursorandPandasCursor.result_set_type_hints: master's managed path passed hints through_get_rows(), and the S3 path never used them, so the note is accurate. The docs line in the Constraints section is updated.docs/aio.md(rows read insideexecute()) stays true.Callers and operations: the removed
_as_*_from_api(),_rows_to_columnar(), and_text_value_converter()are private, and only this package calls them.GetQueryResultstraffic is unchanged: the sameMaxResults=1000pagination from the start, after the one-row pre-fetch, with the same retry config. Peak memory is the CSV text (list, joined string, and bytes) in place of the Python objects of every converted value.