Skip to content
138 changes: 132 additions & 6 deletions pyathena/pandas/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
from pyathena.util import RetryConfig, override, parse_output_location

if TYPE_CHECKING:
from pandas import DataFrame
from pandas import DataFrame, Index, Series
from pandas.io.parsers import TextFileReader

from pyathena.connection import Connection
Expand All @@ -43,6 +43,75 @@ def _no_trunc_date(df: DataFrame) -> DataFrame:
return df


class _JSONConverter:

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.

Design change (maintainer decision): _CSVObject → _JSONConverter (NULL placeholder), 88cd9cd.

Measured: wrapping every json value cost about 30% on json columns, half of it GC (+35% with GC on, +18% off). Fable 5.1 and Codex were consulted read-only on the alternatives:

Design Measured cost on json columns
Decode raw strings after reading +8–32%
Wrap only scalars +11% on objects, +18% on unrelated object columns (full scans)
Per-column value buffers (Codex) +15–25%
NULL placeholder (Fable) about noise

Both reviewers agree pandas has no option that stops inference on converter output. The maintainer chose the placeholder.

Self-review, round one (git diff e49d576ce1e8cb60d74e1fde56de9c72ff3cd8e1..88cd9cde004bf73a7d3ec05aa7e461ac4a702b1f): CLEAN.

  • Converter instances live per result set, created in _get_csv_converter(). _csv_converters takes the final read_csv_kwargs["converters"]: after _configure_binary_csv_read() (which rebuilds converters on name resolution), and after user converters kwargs, which replace ours, so nothing is restored then.
  • The flag is reset in restore(), which _finish_csv_frame() runs right after each frame or chunk (__next__ and get_chunk()).
  • JSON null literals give None from _to_json and are handled like NULL.
  • usecols that drops a json column never calls its converter.
  • The python engine never converts index_col columns (pre-existing), so their flag stays unset.
  • Mutations: returning None instead of the placeholder fails 4 params; skipping the restore fails 5; skipping the index restore fails index and multi_index; making every managed json column object fails managed.

Round two (claims): the commit message's "within noise of master" holds for the combined benchmarks (−1.1%, +3.3%, +2.3%), but single-column cases measured up to +13.8% (numbers without NULL) and +16.4% (objects with NULL in 50k chunks). The PR body states per-value costs instead (0.06 µs per call, 0.04 µs per restored value).

New user-visible effect, now in the release notes: large low_memory reads with sparse NULL emit DtypeWarning with exact values. Reproduced offline: 600,000 rows, one NULL at the end. Before, the column was silently float64.

Evidence: 317 passed on tests/pyathena/pandas/ + tests/pyathena/aio/pandas/ locally; CI test/run (3.14) passed on 88cd9cd.

"""A json converter for ``pandas.read_csv()`` that keeps NULL from making values floats.

pandas infers a dtype from the values that a converter returns, so JSON numbers
with NULL become float64. This converter returns ``NULL`` in place of None, which
keeps the column object, and ``restore()`` puts None back after reading.
"""

NULL: ClassVar[object] = object()

__slots__ = ("_converter", "_has_null")

def __init__(self, converter: Callable[[str | None], Any]) -> None:
"""Wrap a json conversion function.

Args:
converter: The conversion function, which returns None for NULL.
"""
self._converter = converter
self._has_null = False

def __call__(self, value: str | None) -> Any:
"""Convert a CSV value.

Args:
value: The value as text.

Returns:
The converted value, or ``NULL`` in place of None.
"""
converted = self._converter(value)
if converted is None:
self._has_null = True
return self.NULL
return converted

def restore(self, df: DataFrame, name: Any) -> None:
"""Put None back in place of ``NULL`` in a DataFrame that ``read_csv()`` returned.

Only does anything if this converter returned ``NULL`` since the last call,
so call it for each DataFrame or chunk right after reading it.

Args:
df: The DataFrame or chunk.
name: The name of the column that this converter converted, which can
also be an index level.
"""
if not self._has_null:
return
self._has_null = False

import pandas as pd

index = df.index
if name in df.columns:
df[name] = pd.Series(self._restore_values(df[name]), index=index, dtype=object)
elif isinstance(index, pd.MultiIndex) and name in index.names:
levels = [index.get_level_values(i) for i in range(index.nlevels)]
i = index.names.index(name)
levels[i] = pd.Index(self._restore_values(levels[i]), dtype=object, name=name)
df.index = pd.MultiIndex.from_arrays(levels, names=index.names)
elif index.name == name:
df.index = pd.Index(self._restore_values(index), dtype=object, name=name)

@classmethod
def _restore_values(cls, values: Series | Index) -> list[Any]:
return [None if v is cls.NULL else v for v in values.to_numpy()]


class PandasDataFrameIterator(abc.Iterator): # type: ignore[type-arg]
"""Iterator for chunked DataFrame results from Athena queries.

Expand Down Expand Up @@ -78,7 +147,7 @@ def __init__(

Args:
reader: Either a TextFileReader (for chunked) or a single DataFrame.
trunc_date: Function to apply date truncation to each chunk.
trunc_date: Function to apply to each chunk, such as date truncation.
csv_stream: Optional CSV stream owned and closed by this iterator.
"""
from pandas import DataFrame
Expand Down Expand Up @@ -254,6 +323,7 @@ class AthenaPandasResultSet(AthenaResultSet):
AUTO_CHUNK_SIZE_LARGE: int = 100_000
AUTO_CHUNK_SIZE_MEDIUM: int = 50_000

_INTEGER_TYPES: ClassVar[tuple[str, ...]] = ("tinyint", "smallint", "integer", "bigint")
_PARSE_DATES: ClassVar[list[str]] = [
"date",
"time",
Expand Down Expand Up @@ -338,6 +408,8 @@ def __init__(
self._kwargs = kwargs
self._fs = self._create_s3_file_system()
self._csv_stream: IOBase | None = None
# The converters that pandas.read_csv() applies, keyed by column name.
self._csv_converters: dict[Any, Callable[[str | None], Any]] = {}

# Cache time column names for efficient _trunc_date processing
description = self.description if self.description else []
Expand All @@ -349,7 +421,7 @@ def __init__(
self._df: DataFrame | None = None
if self.state == AthenaQueryExecution.STATE_SUCCEEDED and self.output_location:
result = self._as_pandas()
trunc_date = _no_trunc_date if self.is_unload else self._trunc_date
trunc_date = _no_trunc_date if self.is_unload else self._finish_csv_frame
if isinstance(result, pd.DataFrame):
self._df = trunc_date(result)
else:
Expand Down Expand Up @@ -516,6 +588,39 @@ def parse_dates(self) -> list[Any | None]:
description = self.description if self.description else []
return [d[0] for d in description if d[1] in self._PARSE_DATES]

def _finish_csv_frame(self, df: DataFrame) -> DataFrame:
"""Finish a DataFrame read from the CSV result file.

Puts None back in the json columns and truncates the time columns.

Args:
df: The DataFrame or chunk that ``pandas.read_csv()`` returned.

Returns:
The same DataFrame.
"""
for name, converter in self._csv_converters.items():
if isinstance(converter, _JSONConverter):
converter.restore(df, name)
return self._trunc_date(df)

def _get_csv_converter(self, type_: str) -> Callable[[str | None], Any]:
"""Get the converter that ``pandas.read_csv()`` applies to a column type.

json columns get a ``_JSONConverter``, so that NULL does not make pandas infer
a numeric dtype, and ``_finish_csv_frame()`` puts None back.

Args:
type_: The Athena type of the column.

Returns:
The conversion function.
"""
converter = self._converter.get(type_)
if type_ == "json":
return _JSONConverter(converter)
return converter

def _trunc_date(self, df: DataFrame) -> DataFrame:
if self._time_columns:
# A NULL is None, as with the GetQueryResults fallback and the other types.
Expand Down Expand Up @@ -576,6 +681,7 @@ def _read_csv(self) -> TextFileReader | DataFrame:
with ExitStack() as stack:
source: str | IOBase = self.output_location
binary_columns = self._configure_binary_csv_read(read_csv_kwargs, pd.read_csv)
self._csv_converters = read_csv_kwargs.get("converters") or {}
if binary_columns:
# Given storage_options, even None, open the file through fsspec
# as pandas does.
Expand Down Expand Up @@ -626,7 +732,11 @@ def _get_csv_read_options(self, csv_engine: str, chunksize: int | None) -> dict[
"header": header,
"names": names,
"dtype": self.dtypes,
"converters": self.converters,
"converters": {
d[0]: self._get_csv_converter(d[1])
for d in self.description or []
if d[1] in self._converter.mappings
},
"parse_dates": self.parse_dates,
"skip_blank_lines": False,
"keep_default_na": self._keep_default_na,
Expand Down Expand Up @@ -740,7 +850,7 @@ def _configure_binary_csv_read(
if len(column_names) != len(description):
return set()
converters = {
name: self._converter.get(d[1])
name: self._get_csv_converter(d[1])
for name, d in zip(column_names, description, strict=True)
if d[1] in self._converter.mappings and name in selected_names
}
Expand Down Expand Up @@ -848,7 +958,23 @@ def _as_pandas_from_api(self, converter: Converter | None = None) -> DataFrame:
return pd.DataFrame()
description = self.description if self.description else []
columns = [d[0] for d in description]
return pd.DataFrame(self._rows_to_columnar(rows, columns))
columnar = self._rows_to_columnar(rows, columns)
# Integer columns get the dtype that the CSV result file reads them with,
# and json columns with NULL stay objects as there, so that NULL does not
# make their values floats.
dtypes: dict[str, Any] = {}
for d in description:
if d[1] in self._INTEGER_TYPES:

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 two (claims, callers, operations): FINDINGS. I corrected the PR body; the code is unchanged.

Scope: git diff a73c6db457110ea13eae8f7d1e7670b9095605dd..eff8530ca31909b48703275631794a7147348913.

Claims checked:

  • Managed integers and JSON numbers with NULL were float64, and S3-path JSON numbers too: measured on Athena before and after the change, for 22 types on both storage modes.
  • An all-integer JSON column was int64: measured.
  • With pandas 3, all-string JSON was str: offline read_csv() and pd.DataFrame() checks.
  • read_csv() re-infers converter values even with dtype=object: offline check, pandas 3.0.6.
  • The async and aio cursors share the result set: pyathena/pandas/async_cursor.py and pyathena/aio/pandas/cursor.py import AthenaPandasResultSet.

Finding (claim): the release note said a non-nullable custom integer dtype now fails on managed storage "as it already did" on the S3 path. Verified offline: managed raises a raw TypeError from pd.Series(...) here, while the S3 path raises OperationalError wrapping ValueError: Integer column has NA values. I corrected the PR body to state both exceptions. The managed path has never wrapped _fetch_all_rows() errors, so I left the exception type as is.

Callers: user converters/dtype kwargs still override PyAthena's (read_csv_kwargs.update(self._kwargs)), so they never see wrapped values. The pyarrow CSV engine is only chosen without converters. Repeated as_pandas() returns the same _df. Chunks are unwrapped one at a time. There are no AWS API changes: the same GetQueryResults calls and CSV reads.

Docs: docs/pandas.md and docs/null_handling.md state no json or integer dtype that this changes, so nothing became obsolete.

Evidence: published head eff8530, tests/pyathena/pandas/ + tests/pyathena/aio/pandas/: 300 passed locally (managed params included). The SQLAlchemy suites have not run yet; CI runs them on Ready.

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.

Rebase check: rebased a73c6db..05b08d9 onto master ec5323e, giving head cfc3976. git range-diff shows the four commits unchanged apart from test-file context: #1022's test_read_options and #1033's CONVERTED_VALUES_* imports landed at the same spot as test_integer_and_json_with_null, and I kept both sides.

Upstream contracts checked against this patch:

After the rebase: just lint passed, and 12 passed for -k "test_integer_and_json_with_null or test_integer_without_dtype or test_fetch_all_rows or test_read_options". The full pandas + aio pandas suite is running on cfc3976.

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.

Rebase check 2: rebased cfc3976..e25b9f0 onto master e49d576, giving head 0b77627.

#1054 (#1041) independently made PandasDataFrameIterator.get_chunk() apply the per-chunk function, so the conflict in that hunk resolves to upstream's code, which behaves the same. b7d0729 (was fbeb26a) now holds only the pd.array() change and the test, and I reworded its message. The get_chunk() release-note item moved out of this PR's body to #1054.

#1054's _trunc_date() now returns None for NULL time; the test's col_time has no NULL. Its _to_array parser changes don't touch json columns.

After the rebase: just lint passed, and 62 passed for -k "test_integer_and_json_with_null or test_integer_without_dtype or test_fetch_all_rows or test_complex_as_pandas or binary or get_chunk".

if (dtype := self._converter.get_dtype(d[1], d[4], d[5])) is not None:
dtypes[d[0]] = dtype
elif d[1] == "json" and None in columnar[d[0]]:

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, gpt-6-astra, effort high, read-only, session 01a104a0-9478-7a22-87d7-b77f10273eda. Full static review of e49d576..88cd9cd, run because the design changed, on a detached snapshot at 88cd9cd. The prompt gave the intended behavior and accepted limits, without PR text or prior findings.

Coverage:

  • managed dtypes, SQL NULL and JSON null, and empty results;
  • column and index restore;
  • iteration, get_chunk(), and iterator as_pandas();
  • converters/dtype/names/usecols overrides, engines, low_memory, and varbinary;
  • async/aio, UNLOAD, converter lifetime, performance structure, and tests.

Result: one P2, pre-existing. On managed storage, SELECT 1 AS x, 2 AS x puts both values in one column, through _rows_to_columnar() (pyathena/result_set.py:764). That is #1032, unchanged by this PR; my prompt overstated a duplicate-name guarantee. No new regressions found.

dtypes[d[0]] = object

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 one — checked and no change: pd.Series(values, dtype=get_dtype(...)) uses the same types mapping that read_csv(dtype=...) uses. A custom non-nullable integer dtype now raises on NULL on managed storage as on the S3 path (release note in the PR body). Duplicate column names still collapse in _rows_to_columnar(); that is the existing #1032, not changed here. Wrapper detection checks the first value only. The C and Python parsers call converters on every cell, including NA (''), so a json column is either fully wrapped or not converted.

return pd.DataFrame(
{
name: values if name not in dtypes else pd.array(values, dtype=dtypes[name])

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): Finding 2 (P2): with managed storage, SELECT 1 AS x, 2 AS x, 3 AS y gives x=[1, 2], y=[3] from _rows_to_columnar(), because the duplicate names collide (existing #1032). Building the columns as Series aligned them on the index, so the result had two rows with NA instead of raising ValueError as at the base.

Verified offline (pandas 3.0.6): dict of lists gives ValueError: All arrays must be of the same length; dict of Series gives a 2-row frame with NaN. Repaired in fbeb26a: typed columns use pd.array(values, dtype=...), which is not aligned and raises the same ValueError. Dtype results are unchanged (Int64, object, float, str; int64 with NULL raises the same TypeError). I didn't add a test, because one would pin the #1032 ValueError, which that issue will replace.

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, effort high, read-only, session 01a102b6-373c-7f43-a351-549a16767249. Static review of the delta eff8530..fbeb26a on a detached snapshot at fbeb26a. Coverage: get_chunk()/__next__()/collection and sync/async callers (CSV chunks finished once; non-chunked, UNLOAD, and managed use the no-op function); unequal column lengths raise ValueError again; pd.array() vs pd.Series() dtypes. Result: FINDINGS (1 × P2).

Finding: when a custom converter's types omits an integer type, get_dtype() returns None. pd.array([1, None], dtype=None) infers Int64, where pd.Series() and read_csv() (S3 path) infer float64/NaN. Verified. Repaired in 05b08d9: integer columns without a converter dtype keep their list and pandas' inference. New test_integer_without_dtype (default + managed, converter without a bigint dtype) expects float64 and NaN on both. Recording None again fails [managed].

Declined: the reviewer also suggested a managed duplicate-name test. The ValueError it would assert is the defect tracked in #1032, which will replace that error, so I didn't pin it.

Self-review of the repairs (fbeb26a, 05b08d9), round one: get_chunk's callback is _no_trunc_date except for CSV chunks; explicit dtypes (Int64, object, float, str) are unchanged under pd.array(). Round two: the commit messages' claims (ValueError restored, float64 inference) are checked offline and by the test; the PR body is updated with the get_chunk() time truncation release note.

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, effort high, read-only, session 01a102c0-e95f-70a0-a6c4-c4fd0ea6bfdb. Static review of the delta fbeb26a..05b08d9 on a detached snapshot at 05b08d9. Coverage: API/CSV dtype handling, converter propagation, regression-test assertions, helper and fixture state isolation. Result: CLEAN. Omitting None dtypes restores float64/NaN inference, and explicit dtypes and JSON handling are unchanged. The managed case would fail against the previous Int64 result. Each helper call builds its own converter, so removing bigint doesn't leak to other tests.

for name, values in columnar.items()
}
)

def as_pandas(self) -> PandasDataFrameIterator | DataFrame:
"""Return the query results as a DataFrame or an iterator of DataFrame chunks.
Expand Down
109 changes: 109 additions & 0 deletions tests/pyathena/pandas/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,12 @@
from tests.pyathena.util import CONVERTED_VALUES_QUERY, CONVERTED_VALUES_ROW, cached_file_systems


def _pandas_converter_without_bigint_dtype():
converter = DefaultPandasTypeConverter()
del converter.types["bigint"]
return converter


class TestPandasCursor:
@pytest.mark.parametrize(
("engine", "chunksize"), [("auto", None), ("c", 2), ("python", 2), ("pyarrow", None)]
Expand Down Expand Up @@ -1671,3 +1677,106 @@ def test_read_options(self, execute_kwargs):
kwargs = result_set_class.call_args.kwargs
expected = {**cursor_kwargs, **execute_kwargs}
assert {key: kwargs[key] for key in expected} == expected

@pytest.mark.parametrize(
"pandas_cursor",
[
pytest.param({"converter": _pandas_converter_without_bigint_dtype()}, id="default"),
pytest.param(
{
"work_group": ENV.managed_work_group,
"s3_staging_dir": "",
"converter": _pandas_converter_without_bigint_dtype(),
},
id="managed",
marks=pytest.mark.skipif(
not ENV.managed_work_group,
reason="AWS_ATHENA_MANAGED_WORKGROUP not set",
),
),
],
indirect=["pandas_cursor"],
)
def test_integer_without_dtype(self, pandas_cursor):
pandas_cursor.execute(
"SELECT * FROM (VALUES BIGINT '1', NULL) AS t(col_bigint) ORDER BY col_bigint"
)
df = pandas_cursor.as_pandas()
assert df["col_bigint"].dtype == np.float64
assert df["col_bigint"].iloc[0] == 1.0
assert math.isnan(df["col_bigint"].iloc[1])

@pytest.mark.parametrize(
("pandas_cursor", "kwargs"),
[
pytest.param({}, {}, id="default"),
pytest.param({}, {"chunksize": 1}, id="chunked"),
pytest.param({}, {"index_col": "col_json"}, id="index"),
pytest.param({}, {"index_col": ["col_int", "col_json"]}, id="multi_index"),
pytest.param(
{},
{
"usecols": [
"col_int",
"col_bigint",
"col_json",
"col_binary",
"col_json_not_null",
"col_time",
]
},
id="usecols",
),
pytest.param(
{"work_group": ENV.managed_work_group, "s3_staging_dir": ""},
{},
id="managed",
marks=pytest.mark.skipif(
not ENV.managed_work_group,
reason="AWS_ATHENA_MANAGED_WORKGROUP not set",
),
),
],
indirect=["pandas_cursor"],
)
def test_integer_and_json_with_null(self, pandas_cursor, kwargs):
pandas_cursor.execute(
"""
SELECT * FROM (VALUES
(1, BIGINT '9007199254740993', json_parse('9007199254740993'), X'01',
json_parse('9007199254740995'), CAST('01:02:03' AS TIME)),
(2, NULL, NULL, NULL, json_parse('9007199254740997'), CAST('04:05:06' AS TIME))
) AS t(col_int, col_bigint, col_json, col_binary, col_json_not_null, col_time)
ORDER BY col_int
""",
**kwargs,
)
df = pandas_cursor.as_pandas()
if "chunksize" in kwargs:
df = pd.concat([df.get_chunk(1), *df])

def column(name):
if name in df.index.names:
return df.index.get_level_values(name)
return df[name]

assert column("col_int").dtype == pd.Int64Dtype()
assert column("col_bigint").dtype == pd.Int64Dtype()
assert column("col_json").dtype == np.object_
# A json column without NULL keeps the dtype that pandas infers.
assert column("col_json_not_null").dtype == np.int64
assert column("col_bigint").tolist() == [9007199254740993, pd.NA]
if isinstance(df.index, pd.MultiIndex):
# A MultiIndex level holds None as NaN.
json_values = column("col_json").tolist()
assert json_values[0] == 9007199254740993
assert isinstance(json_values[0], int)
assert pd.isna(json_values[1])
else:
assert column("col_json").tolist() == [9007199254740993, None]
assert column("col_json_not_null").tolist() == [9007199254740995, 9007199254740997]
assert column("col_binary").tolist() == [b"\x01", None]
assert column("col_time").tolist() == [
datetime(2017, 1, 1, 1, 2, 3).time(),
datetime(2017, 1, 1, 4, 5, 6).time(),
]
Loading