Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/arrow.md
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ class CustomArrowTypeConverter(Converter):
"char": pa.string(),
"varchar": pa.string(),
"string": pa.string(),
"timestamp": pa.timestamp("ms"),
"timestamp": pa.timestamp("us"),
"date": pa.timestamp("ms"),
"time": pa.string(),
"varbinary": pa.string(),
Expand Down
5 changes: 3 additions & 2 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ cursor = connect(work_group="YOUR_MANAGED_WORK_GROUP",
With managed query result storage, query results are retrieved via the `GetQueryResults` API
(1000 rows per request) instead of reading S3 files directly. This may be slower for large
result sets. For large datasets, consider using customer-managed storage or the `UNLOAD` statement.
The pandas, Arrow, and Polars cursors read these rows as they read a CSV result file, with the

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, 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, and docs/arrow.md.

Findings and corrections (PR body only, no code change):

  1. The release note said Arrow and Polars read date as a timestamp and decimal as a string. DefaultPolarsTypeConverter maps date to pl.Date and decimal to pl.Decimal(p, s), so that was false for Polars. The note now describes Arrow's column_types and Polars' schema_overrides separately.
  2. The Known limits said Polars fails on 7 to 12 digits. Only 12 digits was measured to fail; 9 digits reads. It is narrowed to the measurement, and the Arrow 9-digit failure is labeled as measured offline with pyarrow.
  3. The TEST evidence (114/116/12 passed) came from the tree before the final tidy of _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:

  • GetQueryResults text equals the CSV file for 28 types: measured earlier on Athena.
  • The async and aio cursors share the result sets: pyathena/aio/{pandas,arrow,polars}/cursor.py construct the same Athena*ResultSet classes.
  • Empty results have their columns on both paths: measured on Athena for all three cursors. The claim that master had none comes from the source and is labeled as such.
  • The labels-equal row on page 2 was dropped on master (1000 of 1001 rows) and is kept here: measured with S3FSCursor and PandasCursor.
  • 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: no other page states the old managed conversion. docs/aio.md (rows read inside execute()) 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. GetQueryResults traffic is unchanged: the same MaxResults=1000 pagination 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.

same types and converters, but not in chunks.
```

## Cursor iteration
Expand Down Expand Up @@ -709,8 +711,7 @@ column. You can mix both styles in the same dictionary.
JSON-formatted output, which is parsed reliably.
- **Arrow, Pandas, and Polars cursors** — These cursors accept `result_set_type_hints`
but their converters do not currently use the hints because they rely on their own
type systems. The parameter is passed through for forward compatibility and for
result sets that fall back to the default conversion path.
type systems. The parameter is passed through for forward compatibility.

### Breaking change in 3.30.0

Expand Down
30 changes: 28 additions & 2 deletions pyathena/arrow/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -20,6 +21,9 @@
)
from pyathena.util import override

if TYPE_CHECKING:
from pyarrow import ChunkedArray, TimestampType

_logger = logging.getLogger(__name__)


Expand All @@ -34,6 +38,28 @@
}


def _to_timestamp(column: ChunkedArray, type_: TimestampType) -> ChunkedArray:

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 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).

  • Round 1 (behavior):
    • _to_timestamp() moved byte for byte to pyathena/arrow/converter.py, and _to_datetimes() to pyathena/polars/converter.py; checked with a diff of the old and new function bodies. The only change is one comment that named execute(), now generic.
    • _TIMESTAMP_TEXT_LENGTHS is defined once in pyathena/converter.py (s/ms/us/ns). Polars looks up only ms/us/ns, the Datetime units.
    • The result sets import the converter modules; the converter modules do not import the result sets, so there is no import cycle (both import cleanly).
    • pyarrow and polars are still imported inside the functions, and only for type checking at module level, as for the other optional-dependency code.
    • They stay module functions, not methods of the default converters, because a custom converter need not subclass them.
  • Round 2 (claims): the commit message and the PR body (WHAT, tested commit) name the new locations. The unit tests moved to tests/pyathena/{arrow,polars}/test_converter.py, mirroring the source.
  • Tests on d376f51: offline 223 passed (the same tests, relocated). Live Arrow, Polars, and their aio variants filtered to managed/timestamp/duplicate/fetch: 57 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 review of d376f51 (relayed): CLEAN

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, codex exec -s read-only, session 01a10589-d007-7432-b6d5-dd9a514d0337. Static review of 1878a8c..d376f51 on a detached snapshot (left unchanged).

Covered all nine changed files. The helper bodies preserve behavior; Polars changes only a comment. The shared timestamp lengths keep all the values used before. The result set callers and the relocated tests import the helpers from the converter modules, and nothing references the old locations. No import cycles. Deferred annotations, TYPE_CHECKING imports, and lazy optional-dependency imports follow the existing conventions. The moved tests keep identical parameters, inputs, and assertions. The placement matches the converter modules' responsibilities and the source-mirroring test layout. CLEAN: no actionable findings.

"""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.

Expand Down Expand Up @@ -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(),
Expand Down
89 changes: 45 additions & 44 deletions pyathena/arrow/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()

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): FINDINGS (1)

Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort high, codex exec -s read-only, session 01a104da-4dc0-74d2-9bea-0f16596844ec. Static review only: no tests, builds, or network. Snapshot: detached worktree at head 03c9f3f, diff from base a57325b. The prompt had no PR number, description, or prior findings. The snapshot and the PR worktree stayed unchanged (both at 03c9f3f, clean).

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.

P2: Managed Arrow results regress for timestamps exceeding microsecond precision (pyathena/arrow/result_set.py:145). For a managed query returning TIMESTAMP(9) text such as 2020-01-01 00:00:00.123456000, the new CSV path uses timestamp[us]. Arrow's parser rejects fractions longer than six digits, even trailing zeros, so result-set construction raises OperationalError. Previously, the API fallback used _parse_datetime(), which truncates fractions to six digits. Classification: a newly introduced regression for managed results; the S3 CSV limitation already existed. The new comparison test covers only 3- and 6-digit timestamps.

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 DefaultTypeConverter (_parse_datetime() accepts 7 to 12 digits). pyarrow rejects more than 6 digits for timestamp[us], as measured offline. The PR body already lists this under Known limits; the resolution is pending the maintainer's decision.

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.

Repair: 46bf916 (the maintainer chose to fix this in this PR rather than defer it)

  • ArrowCursor (pyathena/arrow/result_set.py): when csv.read_csv() raises ArrowInvalid and the result has timestamp columns with a timestamp type, the data is read again with those columns as text. _to_timestamp() truncates the text to the unit's length (s 19, ms 23, us 26, ns 29), turns "" into NULL (the CSV NULL with strings_can_be_null=False), and casts it. The successful parse path is unchanged.
  • PolarsCursor (pyathena/polars/result_set.py): _get_timestamp_dtypes() picks the timestamp columns whose dtype is Datetime. They are read as pl.String, and _to_datetimes() slices and parses them (%Y-%m-%d %H:%M:%S%.f, the dtype's time unit) for read_csv() and lazily for chunked scan_csv(). Skipped: schema_overrides given to execute(), headerless .txt results, and columns absent from the frame.

Self-review of the repair

  • Round 1 (behavior):
    • A NULL timestamp is null with strings_can_be_null=True (binary columns present) and "" otherwise. Both give NULL (offline check with and without a varbinary column).
    • The binary fill_null("") loop runs after the conversion, so the converted columns are no longer strings and are skipped.
    • The retry opens a new stream (S3) or BufferReader (managed). A second failure is wrapped in OperationalError as before.
    • Polars failures inside the chunk iterator are raised within the existing try and wrapped in OperationalError.
    • Cost: the first design (always text) made Arrow's local whole-file parse +25–40% (1M rows × 8 columns), so Arrow reads as text only after a failure. Polars' explicit-format parse is faster than the inferred one (133–144 → 81–83 ms).
  • Round 2 (claims): the PR body now states the fix, the pre-PR thresholds (Arrow failed with more than 3 digits under timestamp[ms]; Polars with 12, measured, while 9 read), and the double read on failure, and drops the old known limit. The commit message's claim of no cost in the common case matches the unchanged first read.
  • Tests: test_managed_results_match_result_file (Arrow and Polars) now includes TIMESTAMP(12) on both paths and asserts the truncated values; removing the retry fails it. New offline tests: test_to_timestamp and test_to_datetimes. On 46bf916: 251 passed (Arrow and Polars, including aio), 105 passed (pandas, S3FS, aio pandas, filtered), 24 offline passed.

An independent follow-up review of the repair is pending.

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 review (relayed): FINDINGS (2)

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, codex exec -s read-only, session 01a104ee-cfcd-71e3-9817-e72a9564e3ac. Static review of 03c9f3f..46bf916 on a detached snapshot at 46bf916 (left unchanged). It found the earlier Arrow precision finding resolved for the default converter.

  1. P2: Renamed Polars timestamp columns remain strings (pyathena/polars/result_set.py:70). With new_columns=["renamed"], t is read as String, Polars renames it, and _to_datetimes() cannot find t, so as_polars()["renamed"] holds strings. Introduced by the repair.
  2. P2: The missing-string option now breaks NULL timestamps (pyathena/polars/result_set.py:77). With missing_utf8_is_empty_string=True, a NULL timestamp read as String becomes "", and strict str.to_datetime() raises, so the read fails with OperationalError (eager and chunked). Introduced by the repair.

Author verification: both reproduce offline on Polars 1.44.2. The same happens with empty_string_is_null=False, the name that missing_utf8_is_empty_string=True became in 1.43.

Repair: 301b22e. Polars now follows the Arrow design: the first read is unchanged, and only when read_csv() raises ComputeError (measured on Polars 1.39.0 and 1.44.2) are the timestamp columns read again as text. _to_datetimes() treats "" as NULL (finding 2). The text read is skipped when execute() was given new_columns or with_column_names (finding 1); such a read fails as before this PR instead of returning strings. Chunked scan_csv() is back to the master code, because managed results are never chunked. This leaves the pre-existing 12-digit failure of chunked S3 reads, listed under Known limits.

Self-review of the repair:

  • Round 1: common reads are byte-for-byte the old path. The retry rebuilds schema_overrides from the result set's dtypes, because the text read is skipped whenever the user gave schema_overrides. Errors in the retry are wrapped in OperationalError.
  • Round 2: the PR body (WHAT, release notes, Known limits, TEST) is updated. The commit message says empty_string_is_null; the exact case is empty_string_is_null=False, formerly missing_utf8_is_empty_string=True.
  • Tests: test_to_datetimes adds "". The new test_read_csv_truncates_timestamps drives the retry through _read_csv() offline, with columns and with new_columns. On 301b22e: live Polars, Arrow, and aio 251 passed; offline 24 passed.

A second independent follow-up on 46bf916..301b22e is pending.

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 review #2 (relayed): CLEAN

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, codex exec -s read-only, session 01a104f8-7629-7b71-b1b7-014f91a9df99. Static review of 46bf916..301b22e on a detached snapshot at 301b22e (left unchanged).

Covered: retry eligibility and OperationalError wrapping; S3 path vs in-memory bytes on the re-read; read kwargs precedence; column selection and renaming; Datetime classes and ms/us/ns/time-zone handling; the chunked-path reversion; Arrow's analogous retry; the changed tests.

Both previous findings are resolved: the initial read uses the original dtypes and new_columns disables the fallback (pyathena/polars/result_set.py:446); the fallback maps "" to null before strict parsing (:73). The retry is bounded to one additional read; S3 paths are reopened, and managed bytes are reused without refetching rows. The excluded override and headerless cases and the chunked precision failures keep their baseline limitations; they are not new regressions.

Test gaps noted: no direct coverage of the S3 retry, the empty-string read option end-to-end, nanoseconds, or zoned dtypes.

Author note on the gaps: test_managed_results_match_result_file (Arrow and Polars) reads TIMESTAMP(12) on the S3 path too, so it goes through the S3 retry live. The empty-string mapping is covered by test_to_datetimes; I did not add an end-to-end case, because the option's name differs across the supported Polars versions (missing_utf8_is_empty_string before 1.43, empty_string_is_null after). Nanosecond and zoned dtypes come only from custom converters and are not covered.

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 after rebase (relayed): FINDINGS (1)

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, codex exec -s read-only, session 01a1051e-1a39-7253-8146-7d763045a88a. Static full review of 8c9c201..ec2bad9 on a detached snapshot (unchanged). The prompt carried no PR number, description, or earlier findings.

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.

P2: Managed Arrow results now fail for duplicate labels with different types (pyathena/arrow/result_set.py:171). Managed results go through CSV parsing, whose dtype mapping is keyed by column name. With SELECT 'abc' AS x, 1 AS x, the second column overwrites the first column's dtype with int32, both CSV columns parse as integers, and execute() raises OperationalError. Previously, _as_arrow_from_api() built columns positionally and returned ('abc', 1). A newly introduced managed-results regression exposing a pre-existing S3-path defect.

Author verification: confirmed from the source. It is the S3-path defect filed as #1051 (pandas dtype and pyarrow column_types keyed by name). Its table shows the managed path working on master for pandas and Arrow. This PR moves managed reads onto the same reader, so pandas regresses the same way, not only Arrow. A fix of #1051 in the shared readers fixes both paths. Handling is pending the maintainer's decision; another session has a fix/1051-duplicate-column-types worktree.

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.

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.

An independent follow-up on ec2bad9 → 1878a8c is pending.

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 review (relayed): CLEAN

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, codex exec -s read-only, session 01a10569-0548-7af2-b7cf-1d20d4ced045. Static review of git range-diff 8c9c201a..ec2bad9d a89a4f7d..1878a8c8 and the full current diff a89a4f7d..1878a8c8b4349acadb61f0f390325537aeeb1466, on a detached snapshot at 1878a8c (left unchanged).

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; .txt parsing; related Polars changes and tests.

The previous finding is resolved for both managed pandas and Arrow results. Arrow keeps master's #1051 positional names, types, and binary columns, skips the CSV header as a parsed row, converts timestamps by position, then restores repeated names (pyathena/arrow/result_set.py:341). pandas applies the resolved labels to per-column options and keeps master's header-as-labels handling. The binary configuration comes before the header skipping, and the managed binary stream keeps the original header and newlines (pyathena/pandas/result_set.py:844). No actionable regression found in repeated timestamp names, timestamp and binary combinations, .txt handling, or the managed-source interactions.

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):
Expand Down Expand Up @@ -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))
Expand All @@ -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",
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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
Expand All @@ -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:
Expand All @@ -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:
Expand Down Expand Up @@ -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.

Expand Down
35 changes: 5 additions & 30 deletions pyathena/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,11 @@ def _to_datetime(varchar_value: str | None) -> datetime | None:
return _parse_datetime(varchar_value)


# The length of Athena timestamp text, such as "2020-01-02 03:04:05.123456", that holds
# the fraction a time unit can represent. Athena writes up to 12 fractional digits.
_TIMESTAMP_TEXT_LENGTHS: dict[str, int] = {"s": 19, "ms": 23, "us": 26, "ns": 29}


_UTC_OFFSET_PATTERN: re.Pattern[str] = re.compile(r"([+-])(\d{2}):(\d{2})")


Expand Down Expand Up @@ -803,33 +808,3 @@ 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]


# The types whose values the Arrow and Polars GetQueryResults fallbacks keep as text,
# as in a CSV result file, and convert when the rows are fetched.
_TEXT_VALUE_TYPES: tuple[str, ...] = ("json", "time with time zone", "timestamp with time zone")


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

Values nested in typed complex values keep only the time zone types as text,
because Arrow and Polars time and timestamp types hold one time zone per column;
nested JSON values are decoded as before.

Returns:
The converter.
"""
converter = DefaultTypeConverter()
for type_ in _TEXT_VALUE_TYPES:
converter.set(type_, _to_default)
converter._typed_converter = TypedValueConverter(
converters={
**_DEFAULT_CONVERTERS,
"time with time zone": _to_default,
"timestamp with time zone": _to_default,
},
default_converter=_to_default,
struct_parser=_to_struct,
)
return converter
Loading
Loading