From e52cc8a8e0c75bd81f3fbe431d0161c6c76a62c8 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:29:23 +0900 Subject: [PATCH 1/7] Keep integer and JSON values exact in PandasCursor results with NULL With managed query result storage, integer columns with NULL became float64, losing integers above 2**53, because the DataFrame was built from GetQueryResults rows with inferred dtypes. They now get the dtype that the CSV result file reads them with (Int64 by default), and json columns stay object. In the CSV result file, pandas inferred a numeric dtype from the values that the json converter returned, so JSON numbers with NULL became float64 on that path too. The json converter's values are now wrapped during read_csv() and unwrapped into an object column afterwards, for whole results, chunks, and an index column. Closes #1027 Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 100 +++++++++++++++++++++++++-- tests/pyathena/pandas/test_cursor.py | 46 ++++++++++++ 2 files changed, 141 insertions(+), 5 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 54a864e0..874e9946 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -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,39 @@ def _no_trunc_date(df: DataFrame) -> DataFrame: return df +class _CSVObject: + """A converted CSV value that pandas keeps as is instead of inferring a dtype from it.""" + + __slots__ = ("value",) + + def __init__(self, value: Any) -> None: + """Wrap a converted value. + + Args: + value: The value that a converter returned. + """ + self.value = value + + +def _convert_csv_object(converter: Callable[[str | None], Any], value: str) -> _CSVObject: + return _CSVObject(converter(value)) + + +def _unwrap_csv_objects(values: Series | Index) -> list[Any] | None: + """Unwrap the values of a column whose CSV values a converter kept as objects. + + Args: + values: The values of a column or an index. + + Returns: + The converted values, or None if the values are not wrapped. + """ + array = values.array + if values.dtype != object or not len(array) or not isinstance(array[0], _CSVObject): + return None + return [v.value if isinstance(v, _CSVObject) else v for v in array] + + class PandasDataFrameIterator(abc.Iterator): # type: ignore[type-arg] """Iterator for chunked DataFrame results from Athena queries. @@ -254,6 +287,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", @@ -349,7 +383,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 +550,47 @@ 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. + + Unwraps the values that converters kept as objects and truncates the time + columns. + + Args: + df: The DataFrame or chunk that ``pandas.read_csv()`` returned. + + Returns: + The same DataFrame. + """ + import pandas as pd + + for i in range(df.shape[1]): + if (values := _unwrap_csv_objects(df.iloc[:, i])) is not None: + df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) + if ( + not isinstance(df.index, pd.MultiIndex) + and (values := _unwrap_csv_objects(df.index)) is not None + ): + df.index = pd.Index(values, dtype=object, name=df.index.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. + + The values of json columns are wrapped, so that pandas does not infer a + numeric dtype from them, and ``_finish_csv_frame()`` unwraps them. + + Args: + type_: The Athena type of the column. + + Returns: + The conversion function. + """ + converter = self._converter.get(type_) + if type_ == "json": + return partial(_convert_csv_object, 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. @@ -626,7 +701,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 +819,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 +927,18 @@ 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 values stay as decoded, so that NULL does not make them floats. + dtypes: dict[str, Any] = {} + for d in description: + if d[1] in self._INTEGER_TYPES: + dtypes[d[0]] = self._converter.get_dtype(d[1], d[4], d[5]) + elif d[1] == "json": + dtypes[d[0]] = object + return pd.DataFrame( + {name: pd.Series(values, dtype=dtypes.get(name)) 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. diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index d4d78c9c..15df4c20 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -1671,3 +1671,49 @@ 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", "kwargs"), + [ + pytest.param({}, {}, id="default"), + pytest.param({}, {"chunksize": 1}, id="chunked"), + pytest.param({}, {"index_col": "col_json_index"}, id="index"), + 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')), + (2, NULL, NULL, NULL, json_parse('9007199254740997')) + ) AS t(col_int, col_bigint, col_json, col_binary, col_json_index) + ORDER BY col_int + """, + **kwargs, + ) + df = pandas_cursor.as_pandas() + if "chunksize" in kwargs: + df = pd.concat(list(df)) + if "index_col" in kwargs: + assert df.index.dtype == np.object_ + assert df.index.tolist() == [9007199254740995, 9007199254740997] + else: + assert df["col_json_index"].dtype == np.object_ + assert df["col_json_index"].tolist() == [9007199254740995, 9007199254740997] + assert df["col_int"].dtype == pd.Int64Dtype() + assert df["col_bigint"].dtype == pd.Int64Dtype() + assert df["col_json"].dtype == np.object_ + assert df["col_bigint"].tolist() == [9007199254740993, pd.NA] + assert df["col_json"].tolist() == [9007199254740993, None] + assert df["col_binary"].tolist() == [b"\x01", None] From 5cc88f32c710ac40302a581d887712a103492800 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:34:48 +0900 Subject: [PATCH 2/7] Unwrap JSON values in MultiIndex levels and cover column name resolution An index_col list with a json column left the wrapped values in its MultiIndex level. The regression test also covers usecols, which rebuilds the converters in _configure_binary_csv_read(), and a MultiIndex. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 24 +++++++++++++++------ tests/pyathena/pandas/test_cursor.py | 32 +++++++++++++++++----------- 2 files changed, 38 insertions(+), 18 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 874e9946..da3bd952 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -111,7 +111,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 @@ -567,11 +567,23 @@ def _finish_csv_frame(self, df: DataFrame) -> DataFrame: for i in range(df.shape[1]): if (values := _unwrap_csv_objects(df.iloc[:, i])) is not None: df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) - if ( - not isinstance(df.index, pd.MultiIndex) - and (values := _unwrap_csv_objects(df.index)) is not None - ): - df.index = pd.Index(values, dtype=object, name=df.index.name) + index = df.index + levels = ( + [index.get_level_values(i) for i in range(index.nlevels)] + if isinstance(index, pd.MultiIndex) + else [index] + ) + unwrapped = [_unwrap_csv_objects(level) for level in levels] + if any(values is not None for values in unwrapped): + levels = [ + level if values is None else pd.Index(values, dtype=object, name=level.name) + for level, values in zip(levels, unwrapped, strict=True) + ] + df.index = ( + pd.MultiIndex.from_arrays(levels, names=index.names) + if isinstance(index, pd.MultiIndex) + else levels[0] + ) return self._trunc_date(df) def _get_csv_converter(self, type_: str) -> Callable[[str | None], Any]: diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 15df4c20..38e6a543 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -1678,6 +1678,12 @@ def test_read_options(self, execute_kwargs): pytest.param({}, {}, id="default"), pytest.param({}, {"chunksize": 1}, id="chunked"), pytest.param({}, {"index_col": "col_json_index"}, id="index"), + pytest.param({}, {"index_col": ["col_int", "col_json_index"]}, id="multi_index"), + pytest.param( + {}, + {"usecols": ["col_int", "col_bigint", "col_json", "col_binary", "col_json_index"]}, + id="usecols", + ), pytest.param( {"work_group": ENV.managed_work_group, "s3_staging_dir": ""}, {}, @@ -1705,15 +1711,17 @@ def test_integer_and_json_with_null(self, pandas_cursor, kwargs): df = pandas_cursor.as_pandas() if "chunksize" in kwargs: df = pd.concat(list(df)) - if "index_col" in kwargs: - assert df.index.dtype == np.object_ - assert df.index.tolist() == [9007199254740995, 9007199254740997] - else: - assert df["col_json_index"].dtype == np.object_ - assert df["col_json_index"].tolist() == [9007199254740995, 9007199254740997] - assert df["col_int"].dtype == pd.Int64Dtype() - assert df["col_bigint"].dtype == pd.Int64Dtype() - assert df["col_json"].dtype == np.object_ - assert df["col_bigint"].tolist() == [9007199254740993, pd.NA] - assert df["col_json"].tolist() == [9007199254740993, None] - assert df["col_binary"].tolist() == [b"\x01", None] + + 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_ + assert column("col_json_index").dtype == np.object_ + assert column("col_bigint").tolist() == [9007199254740993, pd.NA] + assert column("col_json").tolist() == [9007199254740993, None] + assert column("col_json_index").tolist() == [9007199254740995, 9007199254740997] + assert column("col_binary").tolist() == [b"\x01", None] From b7d07290da974e996ce5af93364db4bb799a1b4e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 02:00:15 +0900 Subject: [PATCH 3/7] Keep unequal managed columns an error and cover get_chunk() Building the managed columns as Series aligned them on their index, so a duplicate column name, whose values _rows_to_columnar() collects into one list, filled extra rows with NA instead of raising ValueError as before. The typed columns are now pandas arrays, which are not aligned. The chunked test reads the first chunk with get_chunk(), which unwraps JSON values and truncates time columns as iteration does. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 5 ++++- tests/pyathena/pandas/test_cursor.py | 23 ++++++++++++++++++----- 2 files changed, 22 insertions(+), 6 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index da3bd952..6256f02a 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -949,7 +949,10 @@ def _as_pandas_from_api(self, converter: Converter | None = None) -> DataFrame: elif d[1] == "json": dtypes[d[0]] = object return pd.DataFrame( - {name: pd.Series(values, dtype=dtypes.get(name)) for name, values in columnar.items()} + { + name: values if name not in dtypes else pd.array(values, dtype=dtypes[name]) + for name, values in columnar.items() + } ) def as_pandas(self) -> PandasDataFrameIterator | DataFrame: diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 38e6a543..44e81c23 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -1681,7 +1681,16 @@ def test_read_options(self, execute_kwargs): pytest.param({}, {"index_col": ["col_int", "col_json_index"]}, id="multi_index"), pytest.param( {}, - {"usecols": ["col_int", "col_bigint", "col_json", "col_binary", "col_json_index"]}, + { + "usecols": [ + "col_int", + "col_bigint", + "col_json", + "col_binary", + "col_json_index", + "col_time", + ] + }, id="usecols", ), pytest.param( @@ -1701,16 +1710,16 @@ def test_integer_and_json_with_null(self, pandas_cursor, kwargs): """ SELECT * FROM (VALUES (1, BIGINT '9007199254740993', json_parse('9007199254740993'), X'01', - json_parse('9007199254740995')), - (2, NULL, NULL, NULL, json_parse('9007199254740997')) - ) AS t(col_int, col_bigint, col_json, col_binary, col_json_index) + 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_index, col_time) ORDER BY col_int """, **kwargs, ) df = pandas_cursor.as_pandas() if "chunksize" in kwargs: - df = pd.concat(list(df)) + df = pd.concat([df.get_chunk(1), *df]) def column(name): if name in df.index.names: @@ -1725,3 +1734,7 @@ def column(name): assert column("col_json").tolist() == [9007199254740993, None] assert column("col_json_index").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(), + ] From 602a64c2c5023830c8e1296c803a558c185d84f3 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 02:12:02 +0900 Subject: [PATCH 4/7] Infer managed integer columns when the converter gives no dtype pd.array() infers Int64 for integers with NULL where the CSV result file and the earlier pd.Series() infer float64, so integer columns without a converter dtype keep their list and pandas' inference. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 3 ++- tests/pyathena/pandas/test_cursor.py | 34 ++++++++++++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 6256f02a..eea94d24 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -945,7 +945,8 @@ def _as_pandas_from_api(self, converter: Converter | None = None) -> DataFrame: dtypes: dict[str, Any] = {} for d in description: if d[1] in self._INTEGER_TYPES: - dtypes[d[0]] = self._converter.get_dtype(d[1], d[4], d[5]) + if (dtype := self._converter.get_dtype(d[1], d[4], d[5])) is not None: + dtypes[d[0]] = dtype elif d[1] == "json": dtypes[d[0]] = object return pd.DataFrame( diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 44e81c23..3cf86bea 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -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)] @@ -1672,6 +1678,34 @@ def test_read_options(self, execute_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"), [ From 0b77627329f772a689955bf5658a1af3d25b2cc0 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 09:23:07 +0900 Subject: [PATCH 5/7] Gather the CSV value wrapping in _CSVObject The converter wrapper and the unwrapping of columns and index levels become _CSVObject.wrap() and _CSVObject.unwrap(), so that the class holds everything about values that pandas.read_csv() must not infer a dtype from. Behavior is unchanged. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 105 +++++++++++++++++++++------------- 1 file changed, 65 insertions(+), 40 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index eea94d24..b4d74cc6 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -44,7 +44,12 @@ def _no_trunc_date(df: DataFrame) -> DataFrame: class _CSVObject: - """A converted CSV value that pandas keeps as is instead of inferring a dtype from it.""" + """A converted CSV value that ``pandas.read_csv()`` keeps as is. + + pandas infers a dtype from the values that a converter returns, so JSON numbers + with NULL become float64. A converter from ``wrap()`` returns its values in this + class, and ``unwrap()`` restores them as object columns after reading. + """ __slots__ = ("value",) @@ -56,24 +61,66 @@ def __init__(self, value: Any) -> None: """ self.value = value + @classmethod + def wrap(cls, converter: Callable[[str | None], Any]) -> Callable[[str | None], _CSVObject]: + """Wrap a converter so that it returns its values in this class. -def _convert_csv_object(converter: Callable[[str | None], Any], value: str) -> _CSVObject: - return _CSVObject(converter(value)) + Args: + converter: The conversion function. + Returns: + The conversion function that wraps the converted values. + """ + return lambda value: cls(converter(value)) -def _unwrap_csv_objects(values: Series | Index) -> list[Any] | None: - """Unwrap the values of a column whose CSV values a converter kept as objects. + @classmethod + def unwrap(cls, df: DataFrame) -> DataFrame: + """Restore the wrapped values in the columns and the index of a DataFrame. - Args: - values: The values of a column or an index. + Args: + df: The DataFrame or chunk that ``pandas.read_csv()`` returned. - Returns: - The converted values, or None if the values are not wrapped. - """ - array = values.array - if values.dtype != object or not len(array) or not isinstance(array[0], _CSVObject): - return None - return [v.value if isinstance(v, _CSVObject) else v for v in array] + Returns: + The same DataFrame, with the wrapped values in object columns and levels. + """ + import pandas as pd + + for i in range(df.shape[1]): + if (values := cls._unwrap_values(df.iloc[:, i])) is not None: + df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) + index = df.index + levels = ( + [index.get_level_values(i) for i in range(index.nlevels)] + if isinstance(index, pd.MultiIndex) + else [index] + ) + unwrapped = [cls._unwrap_values(level) for level in levels] + if any(values is not None for values in unwrapped): + levels = [ + level if values is None else pd.Index(values, dtype=object, name=level.name) + for level, values in zip(levels, unwrapped, strict=True) + ] + df.index = ( + pd.MultiIndex.from_arrays(levels, names=index.names) + if isinstance(index, pd.MultiIndex) + else levels[0] + ) + return df + + @classmethod + def _unwrap_values(cls, values: Series | Index) -> list[Any] | None: + """Unwrap the values of a column or an index level. + + Args: + values: The values of a column or an index level. + + Returns: + The converted values, or None if the values are not wrapped. + """ + array = values.array + if values.dtype != object or not len(array) or not isinstance(array[0], cls): + return None + return [v.value if isinstance(v, cls) else v for v in array] class PandasDataFrameIterator(abc.Iterator): # type: ignore[type-arg] @@ -562,35 +609,13 @@ def _finish_csv_frame(self, df: DataFrame) -> DataFrame: Returns: The same DataFrame. """ - import pandas as pd - - for i in range(df.shape[1]): - if (values := _unwrap_csv_objects(df.iloc[:, i])) is not None: - df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) - index = df.index - levels = ( - [index.get_level_values(i) for i in range(index.nlevels)] - if isinstance(index, pd.MultiIndex) - else [index] - ) - unwrapped = [_unwrap_csv_objects(level) for level in levels] - if any(values is not None for values in unwrapped): - levels = [ - level if values is None else pd.Index(values, dtype=object, name=level.name) - for level, values in zip(levels, unwrapped, strict=True) - ] - df.index = ( - pd.MultiIndex.from_arrays(levels, names=index.names) - if isinstance(index, pd.MultiIndex) - else levels[0] - ) - return self._trunc_date(df) + return self._trunc_date(_CSVObject.unwrap(df)) def _get_csv_converter(self, type_: str) -> Callable[[str | None], Any]: """Get the converter that ``pandas.read_csv()`` applies to a column type. - The values of json columns are wrapped, so that pandas does not infer a - numeric dtype from them, and ``_finish_csv_frame()`` unwraps them. + The values of json columns are wrapped in ``_CSVObject``, so that pandas does + not infer a numeric dtype from them, and ``_finish_csv_frame()`` unwraps them. Args: type_: The Athena type of the column. @@ -600,7 +625,7 @@ def _get_csv_converter(self, type_: str) -> Callable[[str | None], Any]: """ converter = self._converter.get(type_) if type_ == "json": - return partial(_convert_csv_object, converter) + return _CSVObject.wrap(converter) return converter def _trunc_date(self, df: DataFrame) -> DataFrame: From c5dd5b76cb402102a42be4191cdbf5195e1d6502 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 10:14:57 +0900 Subject: [PATCH 6/7] Check only object columns and unwrap through NumPy arrays _CSVObject.unwrap() selected the columns of every frame through iloc, which made results without JSON 24% slower to read in 1,000-row chunks of 200 columns; it now reads df.dtypes first and touches only object columns. Iterating values.array was also several times slower than the object ndarray from to_numpy() (172 ms vs 39 ms per million values). Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index b4d74cc6..6094b549 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -84,9 +84,10 @@ def unwrap(cls, df: DataFrame) -> DataFrame: The same DataFrame, with the wrapped values in object columns and levels. """ import pandas as pd + from pandas.api.types import is_object_dtype - for i in range(df.shape[1]): - if (values := cls._unwrap_values(df.iloc[:, i])) is not None: + for i, dtype in enumerate(df.dtypes): + if is_object_dtype(dtype) and (values := cls._unwrap_values(df.iloc[:, i])) is not None: df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) index = df.index levels = ( @@ -117,8 +118,11 @@ def _unwrap_values(cls, values: Series | Index) -> list[Any] | None: Returns: The converted values, or None if the values are not wrapped. """ - array = values.array - if values.dtype != object or not len(array) or not isinstance(array[0], cls): + if values.dtype != object: + return None + # Iterating an object ndarray is much faster than iterating values.array. + array = values.to_numpy() + if not len(array) or not isinstance(array[0], cls): return None return [v.value if isinstance(v, cls) else v for v in array] From 88cd9cde004bf73a7d3ec05aa7e461ac4a702b1f Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 10:56:05 +0900 Subject: [PATCH 7/7] Return a shared placeholder for NULL json values instead of wrapping each value Wrapping every converted json value in an object made json columns about 30% slower to read, half of it garbage collection. pandas only turns JSON numbers into floats when the values mix numbers and None, so the json converter now returns one shared placeholder in place of None, which keeps the column object, and _finish_csv_frame() puts None back in the columns or index levels whose converter returned it. Reading json columns now costs within noise of master. json columns without NULL keep the dtype that pandas infers, as on master, and the managed path makes json columns object only when they have NULL, to match. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 125 +++++++++++++-------------- tests/pyathena/pandas/test_cursor.py | 22 +++-- 2 files changed, 73 insertions(+), 74 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 6094b549..d700fa72 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -43,88 +43,73 @@ def _no_trunc_date(df: DataFrame) -> DataFrame: return df -class _CSVObject: - """A converted CSV value that ``pandas.read_csv()`` keeps as is. +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. A converter from ``wrap()`` returns its values in this - class, and ``unwrap()`` restores them as object columns after reading. + with NULL become float64. This converter returns ``NULL`` in place of None, which + keeps the column object, and ``restore()`` puts None back after reading. """ - __slots__ = ("value",) + NULL: ClassVar[object] = object() - def __init__(self, value: Any) -> None: - """Wrap a converted value. + __slots__ = ("_converter", "_has_null") + + def __init__(self, converter: Callable[[str | None], Any]) -> None: + """Wrap a json conversion function. Args: - value: The value that a converter returned. + converter: The conversion function, which returns None for NULL. """ - self.value = value + self._converter = converter + self._has_null = False - @classmethod - def wrap(cls, converter: Callable[[str | None], Any]) -> Callable[[str | None], _CSVObject]: - """Wrap a converter so that it returns its values in this class. + def __call__(self, value: str | None) -> Any: + """Convert a CSV value. Args: - converter: The conversion function. + value: The value as text. Returns: - The conversion function that wraps the converted values. + The converted value, or ``NULL`` in place of None. """ - return lambda value: cls(converter(value)) + converted = self._converter(value) + if converted is None: + self._has_null = True + return self.NULL + return converted - @classmethod - def unwrap(cls, df: DataFrame) -> DataFrame: - """Restore the wrapped values in the columns and the index of a DataFrame. + def restore(self, df: DataFrame, name: Any) -> None: + """Put None back in place of ``NULL`` in a DataFrame that ``read_csv()`` returned. - Args: - df: The DataFrame or chunk that ``pandas.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. - Returns: - The same DataFrame, with the wrapped values in object columns and levels. + 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 - from pandas.api.types import is_object_dtype - for i, dtype in enumerate(df.dtypes): - if is_object_dtype(dtype) and (values := cls._unwrap_values(df.iloc[:, i])) is not None: - df.isetitem(i, pd.Series(values, index=df.index, dtype=object)) index = df.index - levels = ( - [index.get_level_values(i) for i in range(index.nlevels)] - if isinstance(index, pd.MultiIndex) - else [index] - ) - unwrapped = [cls._unwrap_values(level) for level in levels] - if any(values is not None for values in unwrapped): - levels = [ - level if values is None else pd.Index(values, dtype=object, name=level.name) - for level, values in zip(levels, unwrapped, strict=True) - ] - df.index = ( - pd.MultiIndex.from_arrays(levels, names=index.names) - if isinstance(index, pd.MultiIndex) - else levels[0] - ) - return df + 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 _unwrap_values(cls, values: Series | Index) -> list[Any] | None: - """Unwrap the values of a column or an index level. - - Args: - values: The values of a column or an index level. - - Returns: - The converted values, or None if the values are not wrapped. - """ - if values.dtype != object: - return None - # Iterating an object ndarray is much faster than iterating values.array. - array = values.to_numpy() - if not len(array) or not isinstance(array[0], cls): - return None - return [v.value if isinstance(v, cls) else v for v in array] + 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] @@ -423,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 [] @@ -604,8 +591,7 @@ def parse_dates(self) -> list[Any | None]: def _finish_csv_frame(self, df: DataFrame) -> DataFrame: """Finish a DataFrame read from the CSV result file. - Unwraps the values that converters kept as objects and truncates the time - columns. + Puts None back in the json columns and truncates the time columns. Args: df: The DataFrame or chunk that ``pandas.read_csv()`` returned. @@ -613,13 +599,16 @@ def _finish_csv_frame(self, df: DataFrame) -> DataFrame: Returns: The same DataFrame. """ - return self._trunc_date(_CSVObject.unwrap(df)) + 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. - The values of json columns are wrapped in ``_CSVObject``, so that pandas does - not infer a numeric dtype from them, and ``_finish_csv_frame()`` unwraps them. + 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. @@ -629,7 +618,7 @@ def _get_csv_converter(self, type_: str) -> Callable[[str | None], Any]: """ converter = self._converter.get(type_) if type_ == "json": - return _CSVObject.wrap(converter) + return _JSONConverter(converter) return converter def _trunc_date(self, df: DataFrame) -> DataFrame: @@ -692,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. @@ -970,13 +960,14 @@ def _as_pandas_from_api(self, converter: Converter | None = None) -> DataFrame: columns = [d[0] for d in description] columnar = self._rows_to_columnar(rows, columns) # Integer columns get the dtype that the CSV result file reads them with, - # and json values stay as decoded, so that NULL does not make them floats. + # 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: if (dtype := self._converter.get_dtype(d[1], d[4], d[5])) is not None: dtypes[d[0]] = dtype - elif d[1] == "json": + elif d[1] == "json" and None in columnar[d[0]]: dtypes[d[0]] = object return pd.DataFrame( { diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 3cf86bea..9d0f1279 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -1711,8 +1711,8 @@ def test_integer_without_dtype(self, pandas_cursor): [ pytest.param({}, {}, id="default"), pytest.param({}, {"chunksize": 1}, id="chunked"), - pytest.param({}, {"index_col": "col_json_index"}, id="index"), - pytest.param({}, {"index_col": ["col_int", "col_json_index"]}, id="multi_index"), + pytest.param({}, {"index_col": "col_json"}, id="index"), + pytest.param({}, {"index_col": ["col_int", "col_json"]}, id="multi_index"), pytest.param( {}, { @@ -1721,7 +1721,7 @@ def test_integer_without_dtype(self, pandas_cursor): "col_bigint", "col_json", "col_binary", - "col_json_index", + "col_json_not_null", "col_time", ] }, @@ -1746,7 +1746,7 @@ def test_integer_and_json_with_null(self, pandas_cursor, kwargs): (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_index, col_time) + ) AS t(col_int, col_bigint, col_json, col_binary, col_json_not_null, col_time) ORDER BY col_int """, **kwargs, @@ -1763,10 +1763,18 @@ def column(name): assert column("col_int").dtype == pd.Int64Dtype() assert column("col_bigint").dtype == pd.Int64Dtype() assert column("col_json").dtype == np.object_ - assert column("col_json_index").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] - assert column("col_json").tolist() == [9007199254740993, None] - assert column("col_json_index").tolist() == [9007199254740995, 9007199254740997] + 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(),