-
Notifications
You must be signed in to change notification settings - Fork 116
Keep integer and JSON values exact in PandasCursor results with NULL #1044
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
e52cc8a
5cc88f3
b7d0729
602a64c
0b77627
c5dd5b7
88cd9cd
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 |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -43,6 +43,75 @@ def _no_trunc_date(df: DataFrame) -> DataFrame: | |
| return df | ||
|
|
||
|
|
||
| class _JSONConverter: | ||
| """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. | ||
|
|
||
|
|
@@ -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 | ||
|
|
@@ -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", | ||
|
|
@@ -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 [] | ||
|
|
@@ -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: | ||
|
|
@@ -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. | ||
|
|
@@ -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. | ||
|
|
@@ -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, | ||
|
|
@@ -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 | ||
| } | ||
|
|
@@ -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: | ||
|
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 two (claims, callers, operations): FINDINGS. I corrected the PR body; the code is unchanged. Scope: Claims checked:
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 Callers: user Docs: Evidence: published head eff8530,
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. Rebase check: rebased a73c6db..05b08d9 onto master ec5323e, giving head cfc3976. Upstream contracts checked against this patch:
After the rebase:
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. Rebase check 2: rebased cfc3976..e25b9f0 onto master e49d576, giving head 0b77627. #1054 (#1041) independently made #1054's After the rebase: |
||
| 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]]: | ||
|
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, Coverage:
Result: one P2, pre-existing. On managed storage, |
||
| dtypes[d[0]] = object | ||
|
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 one — checked and no change: |
||
| return pd.DataFrame( | ||
| { | ||
| name: values if name not in dtypes else pd.array(values, dtype=dtypes[name]) | ||
|
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): Finding 2 (P2): with managed storage, Verified offline (pandas 3.0.6): dict of lists gives
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, Finding: when a custom converter's Declined: the reviewer also suggested a managed duplicate-name test. The Self-review of the repairs (fbeb26a, 05b08d9), round one: get_chunk's callback is
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, |
||
| 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. | ||
|
|
||
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.
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:
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._get_csv_converter()._csv_converterstakes the finalread_csv_kwargs["converters"]: after_configure_binary_csv_read()(which rebuilds converters on name resolution), and after userconverterskwargs, which replace ours, so nothing is restored then.restore(), which_finish_csv_frame()runs right after each frame or chunk (__next__andget_chunk()).nullliterals give None from_to_jsonand are handled like NULL.usecolsthat drops a json column never calls its converter.index_colcolumns (pre-existing), so their flag stays unset.indexandmulti_index; making every managed json columnobjectfailsmanaged.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_memoryreads with sparse NULL emitDtypeWarningwith exact values. Reproduced offline: 600,000 rows, one NULL at the end. Before, the column was silentlyfloat64.Evidence: 317 passed on
tests/pyathena/pandas/+tests/pyathena/aio/pandas/locally; CI test/run (3.14) passed on 88cd9cd.