From 70a1b877bb8abb3d25f07fd1459ab9dc6ebbe548 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 22:40:37 +0900 Subject: [PATCH 1/2] Do not convert ArrowCursor fallback values twice With no S3 result file, AthenaArrowResultSet builds its table from GetQueryResults rows that are already converted, and the fetch methods converted them again. TIME, VARBINARY, and JSON object columns raised TypeError, and JSON string scalars were decoded twice. The result set now chooses the conversion functions for the fetch methods once, and uses none for the GetQueryResults fallback. Co-Authored-By: Claude Opus 5.5 --- pyathena/arrow/result_set.py | 7 +++++-- tests/pyathena/arrow/test_cursor.py | 15 +++++++++++++-- 2 files changed, 18 insertions(+), 4 deletions(-) diff --git a/pyathena/arrow/result_set.py b/pyathena/arrow/result_set.py index ea11a419f..38b7ddc06 100644 --- a/pyathena/arrow/result_set.py +++ b/pyathena/arrow/result_set.py @@ -146,6 +146,9 @@ def __init__( import pyarrow as pa self._table = pa.Table.from_pydict({}) + # The conversion functions that the fetch methods apply to the table values. + # GetQueryResults values are already converted. + self._row_converters = self.converters if self.output_location else {} self._batches = iter(self._table.to_batches(arraysize)) def _create_s3_file_system(self): @@ -251,9 +254,9 @@ def _fetch(self) -> None: return else: dict_rows = rows.to_pydict() - column_names = dict_rows.keys() + converters = [self._row_converters.get(k) for k in dict_rows] processed_rows = [ - tuple(self.converters[k](v) for k, v in zip(column_names, row, strict=False)) + tuple(c(v) if c else v for c, v in zip(converters, row, strict=False)) for row in zip(*dict_rows.values(), strict=False) ] self._rows.extend(processed_rows) diff --git a/tests/pyathena/arrow/test_cursor.py b/tests/pyathena/arrow/test_cursor.py index af1abca69..f70e99f65 100644 --- a/tests/pyathena/arrow/test_cursor.py +++ b/tests/pyathena/arrow/test_cursor.py @@ -969,5 +969,16 @@ def test_null_vs_empty_string(self, arrow_cursor): indirect=["arrow_cursor"], ) def test_fetch_all_rows(self, arrow_cursor): - arrow_cursor.execute("SELECT 1 AS col") - assert arrow_cursor.fetchall() == [(1,)] + arrow_cursor.execute( + """ + SELECT + 1 AS col + ,CAST('12:34:56' AS TIME) AS col_time + ,X'0102' AS col_varbinary + ,json_parse('{"a": 1}') AS col_json + ,CAST('{"a": 1}' AS JSON) AS col_json_string + """ + ) + assert arrow_cursor.fetchall() == [ + (1, datetime(2017, 1, 1, 12, 34, 56).time(), b"\x01\x02", {"a": 1}, '{"a": 1}') + ] From 808cc34ffaec7b2d433d7232966a6716bd0bdc5a Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 22:49:28 +0900 Subject: [PATCH 2/2] Resolve the Arrow converters when the rows are fetched Choosing the converters in __init__ ignored converter.set() calls made after execute(). The result set now only records whether the rows come from a result file, and builds the converters once per record batch. Co-Authored-By: Claude Opus 5.5 --- pyathena/arrow/result_set.py | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/pyathena/arrow/result_set.py b/pyathena/arrow/result_set.py index 38b7ddc06..77978264a 100644 --- a/pyathena/arrow/result_set.py +++ b/pyathena/arrow/result_set.py @@ -146,9 +146,9 @@ def __init__( import pyarrow as pa self._table = pa.Table.from_pydict({}) - # The conversion functions that the fetch methods apply to the table values. + # The fetch methods convert only the values read from a result file. # GetQueryResults values are already converted. - self._row_converters = self.converters if self.output_location else {} + self._convert_rows = bool(self.output_location) self._batches = iter(self._table.to_batches(arraysize)) def _create_s3_file_system(self): @@ -254,11 +254,15 @@ def _fetch(self) -> None: return else: dict_rows = rows.to_pydict() - converters = [self._row_converters.get(k) for k in dict_rows] - processed_rows = [ - tuple(c(v) if c else v for c, v in zip(converters, row, strict=False)) - for row in zip(*dict_rows.values(), strict=False) - ] + if self._convert_rows: + converters = self.converters + column_names = dict_rows.keys() + processed_rows = [ + tuple(converters[k](v) for k, v in zip(column_names, row, strict=False)) + for row in zip(*dict_rows.values(), strict=False) + ] + else: + processed_rows = list(zip(*dict_rows.values(), strict=False)) self._rows.extend(processed_rows) @override