From 9375a34d86abc7c6d7d1f0233f8073ba16bf7d2e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 18:36:55 +0900 Subject: [PATCH 1/6] Keep result-set and internal S3FileSystem instances out of the fsspec cache The result sets created their filesystems with the connection as a constructor argument, so fsspec's instance cache kept every connection and the shared dircache grew with every result file read. The internal S3FileSystem of AioS3FileSystem was also cached, so instances created with skip_instance_cache=True still shared it. Create these filesystems with skip_instance_cache=True so that they are released with the result set or the AioS3FileSystem that owns them. Closes #978 Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3_async.py | 2 ++ pyathena/pandas/result_set.py | 5 +++++ pyathena/polars/result_set.py | 3 +++ pyathena/s3fs/result_set.py | 3 +++ tests/pyathena/filesystem/test_s3_async.py | 14 ++++++++++++++ tests/pyathena/pandas/test_cursor.py | 17 +++++++++++++++++ tests/pyathena/polars/test_cursor.py | 12 ++++++++++++ tests/pyathena/s3fs/test_cursor.py | 12 ++++++++++++ 8 files changed, 68 insertions(+) diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index eb1d78b92..92564fb64 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -123,6 +123,8 @@ def __init__( allow_bucket_creation=allow_bucket_creation, allow_bucket_deletion=allow_bucket_deletion, version_aware=version_aware, + # fsspec caches this instance itself when caching is wanted. + skip_instance_cache=True, **kwargs, ) # Share dircache for cache coherence between async and sync instances diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 55c077261..de75fe4b8 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -465,11 +465,14 @@ def _create_s3_file_system(self): """ from pyathena.filesystem.s3 import S3FileSystem + # Not cached by fsspec so that the connection and the dircache are + # released with the result set. return S3FileSystem( connection=self.connection, default_block_size=self._block_size, default_cache_type=self._cache_type, max_workers=self._max_workers, + skip_instance_cache=True, ) @property @@ -613,6 +616,7 @@ def _get_csv_read_options(self, csv_engine: str, chunksize: int | None) -> dict[ "default_block_size": self._block_size, "default_cache_type": self._cache_type, "max_workers": self._max_workers, + "skip_instance_cache": True, }, "chunksize": chunksize, "engine": csv_engine, @@ -782,6 +786,7 @@ def _read_parquet(self, engine) -> DataFrame: "default_block_size": self._block_size, "default_cache_type": self._cache_type, "max_workers": self._max_workers, + "skip_instance_cache": True, }, **kwargs, ) diff --git a/pyathena/polars/result_set.py b/pyathena/polars/result_set.py index b6a2979bb..1fced8450 100644 --- a/pyathena/polars/result_set.py +++ b/pyathena/polars/result_set.py @@ -289,6 +289,9 @@ def _csv_storage_options(self) -> dict[str, Any]: "default_block_size": self._block_size, "default_cache_type": self._cache_type, "max_workers": self._max_workers, + # Not cached by fsspec so that the connection and the dircache are + # released with the result set. + "skip_instance_cache": True, } @property diff --git a/pyathena/s3fs/result_set.py b/pyathena/s3fs/result_set.py index 151bf82fc..74dff6c7b 100644 --- a/pyathena/s3fs/result_set.py +++ b/pyathena/s3fs/result_set.py @@ -131,9 +131,12 @@ def __init__( def _create_s3_file_system(self) -> AbstractFileSystem: """Create S3FileSystem using connection settings.""" + # Not cached by fsspec so that the connection and the dircache are + # released with the result set. return self._filesystem_class( connection=self.connection, default_block_size=self._block_size, + skip_instance_cache=True, ) def _init_csv_reader(self) -> None: diff --git a/tests/pyathena/filesystem/test_s3_async.py b/tests/pyathena/filesystem/test_s3_async.py index 59604bb97..cf9441b5c 100644 --- a/tests/pyathena/filesystem/test_s3_async.py +++ b/tests/pyathena/filesystem/test_s3_async.py @@ -503,6 +503,20 @@ def upload_part_copy(**kw): sync_fs._complete_multipart_upload.call_args.kwargs["RequestPayer"] == "requester" ) + def test_internal_file_system_not_cached(self): + # GH-978: the internal S3FileSystem was kept in the fsspec instance + # cache, so skip_instance_cache=True instances shared it. + connection = mock.MagicMock() + fs1 = AioS3FileSystem(connection=connection, skip_instance_cache=True) + fs2 = AioS3FileSystem(connection=connection, skip_instance_cache=True) + assert fs1._sync_fs is not fs2._sync_fs + assert fs1.dircache is not fs2.dircache + + fs3 = AioS3FileSystem(connection=connection) + assert AioS3FileSystem(connection=connection) is fs3 + for fs in (fs1, fs2, fs3): + assert fs._sync_fs not in S3FileSystem._cache.values() + @pytest.fixture(scope="class") def fs(self, request): if not hasattr(request, "param"): diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 18faf5d5a..173f5d56d 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -15,6 +15,7 @@ from pandas.io.parsers import TextFileReader from pyathena.error import DatabaseError, ProgrammingError +from pyathena.filesystem.s3 import S3FileSystem from pyathena.pandas.converter import DefaultPandasTypeConverter from pyathena.pandas.cursor import PandasCursor from pyathena.pandas.result_set import AthenaPandasResultSet, PandasDataFrameIterator @@ -230,6 +231,22 @@ def test_fetchone(self, pandas_cursor, parquet_engine, chunksize): assert pandas_cursor.rownumber == 1 assert pandas_cursor.fetchone() is None + @pytest.mark.parametrize( + "pandas_cursor", + [{"cursor_kwargs": {"unload": False}}, {"cursor_kwargs": {"unload": True}}], + indirect=True, + ) + def test_result_set_file_system_not_cached(self, pandas_cursor): + # GH-978: the filesystems that read the results were kept in the fsspec + # instance cache with the connection, so the connection was never freed. + pandas_cursor.execute("SELECT * FROM one_row") + assert pandas_cursor.fetchall() == [(1,)] + assert not [ + fs + for fs in S3FileSystem._cache.values() + if fs.storage_options.get("connection") is pandas_cursor.connection + ] + @pytest.mark.parametrize( ("pandas_cursor", "parquet_engine", "chunksize"), [ diff --git a/tests/pyathena/polars/test_cursor.py b/tests/pyathena/polars/test_cursor.py index c9da61696..facb89df9 100644 --- a/tests/pyathena/polars/test_cursor.py +++ b/tests/pyathena/polars/test_cursor.py @@ -17,6 +17,7 @@ import pytest from pyathena.error import DatabaseError, ProgrammingError +from pyathena.filesystem.s3 import S3FileSystem from pyathena.polars.cursor import PolarsCursor from pyathena.polars.result_set import AthenaPolarsResultSet from tests import ENV @@ -36,6 +37,17 @@ def test_fetchone(self, polars_cursor): assert polars_cursor.rownumber == 1 assert polars_cursor.fetchone() is None + def test_result_set_file_system_not_cached(self, polars_cursor): + # GH-978: the filesystem that read the CSV results was kept in the fsspec + # instance cache with the connection, so the connection was never freed. + polars_cursor.execute("SELECT * FROM one_row") + assert polars_cursor.fetchall() == [(1,)] + assert not [ + fs + for fs in S3FileSystem._cache.values() + if fs.storage_options.get("connection") is polars_cursor.connection + ] + @pytest.mark.parametrize( "polars_cursor", [{"cursor_kwargs": {"unload": False}}, {"cursor_kwargs": {"unload": True}}], diff --git a/tests/pyathena/s3fs/test_cursor.py b/tests/pyathena/s3fs/test_cursor.py index 40edaa072..df99feb17 100644 --- a/tests/pyathena/s3fs/test_cursor.py +++ b/tests/pyathena/s3fs/test_cursor.py @@ -9,6 +9,7 @@ import pytest from pyathena.error import DatabaseError, ProgrammingError +from pyathena.filesystem.s3 import S3FileSystem from pyathena.s3fs.cursor import S3FSCursor from pyathena.s3fs.reader import AthenaCSVReader, DefaultCSVReader from pyathena.s3fs.result_set import AthenaS3FSResultSet @@ -24,6 +25,17 @@ def test_fetchone(self, s3fs_cursor): assert s3fs_cursor.rownumber == 1 assert s3fs_cursor.fetchone() is None + def test_result_set_file_system_not_cached(self, s3fs_cursor): + # GH-978: the filesystem that read the results was kept in the fsspec + # instance cache with the connection, so the connection was never freed. + s3fs_cursor.execute("SELECT * FROM one_row") + assert s3fs_cursor.fetchall() == [(1,)] + assert not [ + fs + for fs in S3FileSystem._cache.values() + if fs.storage_options.get("connection") is s3fs_cursor.connection + ] + def test_fetchmany(self, s3fs_cursor): s3fs_cursor.execute("SELECT * FROM many_rows LIMIT 15") assert len(s3fs_cursor.fetchmany(10)) == 10 From 9f5ab6ff754614037895a2ac44166e4c98e065c9 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 18:43:33 +0900 Subject: [PATCH 2/6] Check the AioS3FileSystem instance cache in the result-set tests TestAioS3FileSystem registers AioS3FileSystem for "s3" and does not restore the registration, so pandas and Polars read through AioS3FileSystem in a worker that ran it first. Look for the connection in both instance caches through a shared helper. Co-Authored-By: Claude Opus 5.5 --- tests/pyathena/pandas/test_cursor.py | 8 ++------ tests/pyathena/polars/test_cursor.py | 8 ++------ tests/pyathena/s3fs/test_cursor.py | 8 ++------ tests/pyathena/util.py | 12 ++++++++++++ 4 files changed, 18 insertions(+), 18 deletions(-) diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 173f5d56d..dbe1349be 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -15,12 +15,12 @@ from pandas.io.parsers import TextFileReader from pyathena.error import DatabaseError, ProgrammingError -from pyathena.filesystem.s3 import S3FileSystem from pyathena.pandas.converter import DefaultPandasTypeConverter from pyathena.pandas.cursor import PandasCursor from pyathena.pandas.result_set import AthenaPandasResultSet, PandasDataFrameIterator from tests import ENV from tests.pyathena.conftest import connect +from tests.pyathena.util import cached_file_systems class TestPandasCursor: @@ -241,11 +241,7 @@ def test_result_set_file_system_not_cached(self, pandas_cursor): # instance cache with the connection, so the connection was never freed. pandas_cursor.execute("SELECT * FROM one_row") assert pandas_cursor.fetchall() == [(1,)] - assert not [ - fs - for fs in S3FileSystem._cache.values() - if fs.storage_options.get("connection") is pandas_cursor.connection - ] + assert not cached_file_systems(pandas_cursor.connection) @pytest.mark.parametrize( ("pandas_cursor", "parquet_engine", "chunksize"), diff --git a/tests/pyathena/polars/test_cursor.py b/tests/pyathena/polars/test_cursor.py index facb89df9..343cefb92 100644 --- a/tests/pyathena/polars/test_cursor.py +++ b/tests/pyathena/polars/test_cursor.py @@ -17,11 +17,11 @@ import pytest from pyathena.error import DatabaseError, ProgrammingError -from pyathena.filesystem.s3 import S3FileSystem from pyathena.polars.cursor import PolarsCursor from pyathena.polars.result_set import AthenaPolarsResultSet from tests import ENV from tests.pyathena.conftest import connect +from tests.pyathena.util import cached_file_systems class TestPolarsCursor: @@ -42,11 +42,7 @@ def test_result_set_file_system_not_cached(self, polars_cursor): # instance cache with the connection, so the connection was never freed. polars_cursor.execute("SELECT * FROM one_row") assert polars_cursor.fetchall() == [(1,)] - assert not [ - fs - for fs in S3FileSystem._cache.values() - if fs.storage_options.get("connection") is polars_cursor.connection - ] + assert not cached_file_systems(polars_cursor.connection) @pytest.mark.parametrize( "polars_cursor", diff --git a/tests/pyathena/s3fs/test_cursor.py b/tests/pyathena/s3fs/test_cursor.py index df99feb17..777303f4e 100644 --- a/tests/pyathena/s3fs/test_cursor.py +++ b/tests/pyathena/s3fs/test_cursor.py @@ -9,12 +9,12 @@ import pytest from pyathena.error import DatabaseError, ProgrammingError -from pyathena.filesystem.s3 import S3FileSystem from pyathena.s3fs.cursor import S3FSCursor from pyathena.s3fs.reader import AthenaCSVReader, DefaultCSVReader from pyathena.s3fs.result_set import AthenaS3FSResultSet from tests import ENV from tests.pyathena.conftest import connect +from tests.pyathena.util import cached_file_systems class TestS3FSCursor: @@ -30,11 +30,7 @@ def test_result_set_file_system_not_cached(self, s3fs_cursor): # instance cache with the connection, so the connection was never freed. s3fs_cursor.execute("SELECT * FROM one_row") assert s3fs_cursor.fetchall() == [(1,)] - assert not [ - fs - for fs in S3FileSystem._cache.values() - if fs.storage_options.get("connection") is s3fs_cursor.connection - ] + assert not cached_file_systems(s3fs_cursor.connection) def test_fetchmany(self, s3fs_cursor): s3fs_cursor.execute("SELECT * FROM many_rows LIMIT 15") diff --git a/tests/pyathena/util.py b/tests/pyathena/util.py index 4c6a9bf4d..b2e0cf400 100644 --- a/tests/pyathena/util.py +++ b/tests/pyathena/util.py @@ -15,6 +15,8 @@ from jinja2 import Environment, FileSystemLoader from sqlalchemy import types +from pyathena.filesystem.s3 import S3FileSystem +from pyathena.filesystem.s3_async import AioS3FileSystem from pyathena.glue import GlueMetadataClient from pyathena.model import AthenaCalculationExecutionStatus, AthenaQueryExecution @@ -28,6 +30,16 @@ def read_query(name, **kwargs): return [q.strip() for q in template.render(**kwargs).split(";") if q and q.strip()] +def cached_file_systems(connection): + """Return the filesystems in the fsspec instance cache that hold the connection.""" + return [ + fs + for cls in (S3FileSystem, AioS3FileSystem) + for fs in cls._cache.values() + if fs.storage_options.get("connection") is connection + ] + + METADATA_OPERATIONS = ("get_table_metadata", "list_table_metadata", "list_databases") From 7b82394afbefafd174be60a3eafa662e05c2555d Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 18:45:36 +0900 Subject: [PATCH 3/6] Name the instance that fsspec caches in AioS3FileSystem Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3_async.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index 92564fb64..4157c26a6 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -123,7 +123,7 @@ def __init__( allow_bucket_creation=allow_bucket_creation, allow_bucket_deletion=allow_bucket_deletion, version_aware=version_aware, - # fsspec caches this instance itself when caching is wanted. + # fsspec caches the AioS3FileSystem itself when caching is wanted. skip_instance_cache=True, **kwargs, ) From 0f2a2269ac0ed919eca97a6c182edc323dcb506a Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 19:07:00 +0900 Subject: [PATCH 4/6] Read pandas results through the result set's filesystem pd.read_csv() and pd.read_parquet() built a second S3FileSystem, with its own S3 client and connections, from storage_options for every result set. Open the CSV output through the result set's filesystem and pass it to pd.read_parquet() as filesystem, so each result set builds one S3 client for its reads. storage_options given in the read options still open the CSV output through fsspec with these options. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 61 ++++++++++++++-------------- tests/pyathena/pandas/test_cursor.py | 25 +++++++++--- 2 files changed, 50 insertions(+), 36 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index de75fe4b8..f9eb2054a 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -8,7 +8,7 @@ from collections.abc import Callable, Iterable, Iterator from contextlib import ExitStack from functools import partial -from io import BufferedReader, StringIO, TextIOWrapper +from io import BufferedReader, IOBase, StringIO, TextIOWrapper from multiprocessing import cpu_count from typing import ( TYPE_CHECKING, @@ -72,7 +72,7 @@ def __init__( self, reader: TextFileReader | DataFrame, trunc_date: Callable[[DataFrame], DataFrame], - csv_stream: TextIOWrapper | None = None, + csv_stream: IOBase | None = None, ) -> None: """Initialize the iterator. @@ -335,7 +335,7 @@ def __init__( self._data_manifest: list[str] = [] self._kwargs = kwargs self._fs = self._create_s3_file_system() - self._csv_stream: TextIOWrapper | None = None + self._csv_stream: IOBase | None = None # Cache time column names for efficient _trunc_date processing description = self.description if self.description else [] @@ -561,15 +561,17 @@ def _read_csv(self) -> TextFileReader | DataFrame: try: with ExitStack() as stack: - source: str | TextIOWrapper = self.output_location binary_columns = self._configure_binary_csv_read(read_csv_kwargs, pd.read_csv) + storage_options = read_csv_kwargs.pop("storage_options", None) if binary_columns: - storage_options = read_csv_kwargs.pop("storage_options", None) or {} self._csv_stream = stack.enter_context( self._open_binary_csv_stream(binary_columns, storage_options) ) - source = self._csv_stream - result = pd.read_csv(source, **read_csv_kwargs) + else: + self._csv_stream = stack.enter_context( + self._open_output_location(storage_options, mode="rb") + ) + result = pd.read_csv(self._csv_stream, **read_csv_kwargs) if not isinstance(result, pd.DataFrame): # The chunk iterator takes ownership of the stream. stack.pop_all() @@ -611,13 +613,6 @@ def _get_csv_read_options(self, csv_engine: str, chunksize: int | None) -> dict[ "keep_default_na": self._keep_default_na, "na_values": self._na_values, "quoting": self._quoting, - "storage_options": { - "connection": self.connection, - "default_block_size": self._block_size, - "default_cache_type": self._cache_type, - "max_workers": self._max_workers, - "skip_instance_cache": True, - }, "chunksize": chunksize, "engine": csv_engine, } @@ -739,19 +734,29 @@ def _configure_binary_csv_read( read_csv_kwargs["converters"] = converters return binary_columns + def _open_output_location(self, storage_options: dict[str, Any] | None, **kwargs: Any) -> Any: + """Open the CSV output location for reading. + + Args: + storage_options: The ``storage_options`` given in the read options. Without + them, the file is opened through the filesystem of this result set; + with them, through fsspec with these options. + **kwargs: The mode and text options to open the file with. + + Returns: + A context manager that returns the open file. + """ + if storage_options is None: + return self._fs.open(self.output_location, **kwargs) + return filesystem_open(self.output_location, **kwargs, **storage_options) + def _open_binary_csv_stream( - self, binary_columns: set[int], storage_options: dict[str, Any] + self, binary_columns: set[int], storage_options: dict[str, Any] | None ) -> TextIOWrapper: """Open a stream that preserves binary NULL fields and original CSV newlines.""" with ExitStack() as stack: source = stack.enter_context( - filesystem_open( - self.output_location, - mode="rt", - encoding="utf-8", - newline="", - **storage_options, - ) + self._open_output_location(storage_options, mode="rt", encoding="utf-8", newline="") ) reader = stack.enter_context(BinaryCSVReader(source, binary_columns)) buffer = stack.enter_context(BufferedReader(reader)) @@ -769,7 +774,9 @@ def _read_parquet(self, engine) -> DataFrame: self._unload_location = "/".join(self._data_manifest[0].split("/")[:-1]) + "/" if engine == "pyarrow": - unload_location = self._unload_location + # pyarrow takes the path without the scheme with an fsspec filesystem. + bucket, key = parse_output_location(self._unload_location) + unload_location = f"{bucket}/{key}" kwargs = { "use_threads": True, } @@ -781,13 +788,7 @@ def _read_parquet(self, engine) -> DataFrame: return pd.read_parquet( unload_location, engine=self._engine, - storage_options={ - "connection": self.connection, - "default_block_size": self._block_size, - "default_cache_type": self._cache_type, - "max_workers": self._max_workers, - "skip_instance_cache": True, - }, + filesystem=self._fs, **kwargs, ) except Exception as e: diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index dbe1349be..2007b4697 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -15,6 +15,7 @@ from pandas.io.parsers import TextFileReader from pyathena.error import DatabaseError, ProgrammingError +from pyathena.filesystem.s3 import S3FileSystem from pyathena.pandas.converter import DefaultPandasTypeConverter from pyathena.pandas.cursor import PandasCursor from pyathena.pandas.result_set import AthenaPandasResultSet, PandasDataFrameIterator @@ -232,16 +233,28 @@ def test_fetchone(self, pandas_cursor, parquet_engine, chunksize): assert pandas_cursor.fetchone() is None @pytest.mark.parametrize( - "pandas_cursor", - [{"cursor_kwargs": {"unload": False}}, {"cursor_kwargs": {"unload": True}}], - indirect=True, + ("pandas_cursor", "chunksize"), + [ + ({"cursor_kwargs": {"unload": False}}, None), + ({"cursor_kwargs": {"unload": False}}, 1_000), + ({"cursor_kwargs": {"unload": True}}, None), + ], + indirect=["pandas_cursor"], ) - def test_result_set_file_system_not_cached(self, pandas_cursor): + def test_result_set_file_system(self, pandas_cursor, chunksize): # GH-978: the filesystems that read the results were kept in the fsspec # instance cache with the connection, so the connection was never freed. - pandas_cursor.execute("SELECT * FROM one_row") - assert pandas_cursor.fetchall() == [(1,)] + # The result set reads through its own filesystem instead of creating + # another one from storage_options. + with patch.object( + S3FileSystem, "__init__", autospec=True, side_effect=S3FileSystem.__init__ + ) as init: + pandas_cursor.execute("SELECT * FROM one_row", chunksize=chunksize) + assert pandas_cursor.fetchall() == [(1,)] + assert init.call_count == 1 assert not cached_file_systems(pandas_cursor.connection) + if not pandas_cursor.result_set.is_unload: + assert pandas_cursor.result_set._csv_stream.closed @pytest.mark.parametrize( ("pandas_cursor", "parquet_engine", "chunksize"), From c42670bf7e081df16e1de0eca71755da117497c7 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 19:24:22 +0900 Subject: [PATCH 5/6] Let pandas open the CSV output when storage_options are given Opening the output with fsspec.open() for user storage_options dropped pandas' anonymous-access retry and, with a compression option, the owner that closes the underlying file of a chunked read. Pass the path and the storage_options to pd.read_csv() as before, and open the output through the result set's filesystem only without them. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 40 +++++++++++----------------- tests/pyathena/pandas/test_cursor.py | 18 +++++++++++++ 2 files changed, 33 insertions(+), 25 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index f9eb2054a..82e4cbc3f 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -561,17 +561,19 @@ def _read_csv(self) -> TextFileReader | DataFrame: try: with ExitStack() as stack: + source: str | IOBase = self.output_location binary_columns = self._configure_binary_csv_read(read_csv_kwargs, pd.read_csv) - storage_options = read_csv_kwargs.pop("storage_options", None) if binary_columns: - self._csv_stream = stack.enter_context( + storage_options = read_csv_kwargs.pop("storage_options", None) + source = self._csv_stream = stack.enter_context( self._open_binary_csv_stream(binary_columns, storage_options) ) - else: - self._csv_stream = stack.enter_context( - self._open_output_location(storage_options, mode="rb") + elif "storage_options" not in read_csv_kwargs: + # With storage_options, pandas opens the file through fsspec. + source = self._csv_stream = stack.enter_context( + self._fs.open(self.output_location, mode="rb") ) - result = pd.read_csv(self._csv_stream, **read_csv_kwargs) + result = pd.read_csv(source, **read_csv_kwargs) if not isinstance(result, pd.DataFrame): # The chunk iterator takes ownership of the stream. stack.pop_all() @@ -734,30 +736,18 @@ def _configure_binary_csv_read( read_csv_kwargs["converters"] = converters return binary_columns - def _open_output_location(self, storage_options: dict[str, Any] | None, **kwargs: Any) -> Any: - """Open the CSV output location for reading. - - Args: - storage_options: The ``storage_options`` given in the read options. Without - them, the file is opened through the filesystem of this result set; - with them, through fsspec with these options. - **kwargs: The mode and text options to open the file with. - - Returns: - A context manager that returns the open file. - """ - if storage_options is None: - return self._fs.open(self.output_location, **kwargs) - return filesystem_open(self.output_location, **kwargs, **storage_options) - def _open_binary_csv_stream( self, binary_columns: set[int], storage_options: dict[str, Any] | None ) -> TextIOWrapper: """Open a stream that preserves binary NULL fields and original CSV newlines.""" + text_options: dict[str, Any] = {"mode": "rt", "encoding": "utf-8", "newline": ""} with ExitStack() as stack: - source = stack.enter_context( - self._open_output_location(storage_options, mode="rt", encoding="utf-8", newline="") - ) + if storage_options is None: + source = stack.enter_context(self._fs.open(self.output_location, **text_options)) + else: + source = stack.enter_context( + filesystem_open(self.output_location, **text_options, **storage_options) + ) reader = stack.enter_context(BinaryCSVReader(source, binary_columns)) buffer = stack.enter_context(BufferedReader(reader)) stream = TextIOWrapper(buffer, encoding="utf-8", newline="") diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 2007b4697..54b9bce7b 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -256,6 +256,24 @@ def test_result_set_file_system(self, pandas_cursor, chunksize): if not pandas_cursor.result_set.is_unload: assert pandas_cursor.result_set._csv_stream.closed + def test_csv_storage_options(self, pandas_cursor): + # Given storage_options, pandas opens the CSV output through fsspec with them. + with patch.object( + S3FileSystem, "__init__", autospec=True, side_effect=S3FileSystem.__init__ + ) as init: + pandas_cursor.execute( + "SELECT * FROM one_row", + storage_options={ + "connection": pandas_cursor.connection, + "default_cache_type": "none", + "skip_instance_cache": True, + }, + ) + assert pandas_cursor.fetchall() == [(1,)] + assert init.call_count == 2 + assert init.call_args.kwargs["default_cache_type"] == "none" + assert pandas_cursor.result_set._csv_stream is None + @pytest.mark.parametrize( ("pandas_cursor", "parquet_engine", "chunksize"), [ From 75ef9c266b513ca16fe613929e74c87a85db4f88 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 19:46:40 +0900 Subject: [PATCH 6/6] Open binary CSV output through fsspec for explicit None storage_options An explicit storage_options=None opened binary-column results through the result set's filesystem, while plain results and the code before this change open them through fsspec's default filesystem. Use the result set's filesystem only when the read options have no storage_options key. Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/result_set.py | 6 +++- tests/pyathena/pandas/test_cursor.py | 49 +++++++++++++++++++--------- 2 files changed, 38 insertions(+), 17 deletions(-) diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 82e4cbc3f..10576a678 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -564,7 +564,11 @@ def _read_csv(self) -> TextFileReader | DataFrame: source: str | IOBase = self.output_location binary_columns = self._configure_binary_csv_read(read_csv_kwargs, pd.read_csv) if binary_columns: - storage_options = read_csv_kwargs.pop("storage_options", None) + # Given storage_options, even None, open the file through fsspec + # as pandas does. + storage_options = None + if "storage_options" in read_csv_kwargs: + storage_options = read_csv_kwargs.pop("storage_options") or {} source = self._csv_stream = stack.enter_context( self._open_binary_csv_stream(binary_columns, storage_options) ) diff --git a/tests/pyathena/pandas/test_cursor.py b/tests/pyathena/pandas/test_cursor.py index 54b9bce7b..1ebbe2d25 100644 --- a/tests/pyathena/pandas/test_cursor.py +++ b/tests/pyathena/pandas/test_cursor.py @@ -256,23 +256,40 @@ def test_result_set_file_system(self, pandas_cursor, chunksize): if not pandas_cursor.result_set.is_unload: assert pandas_cursor.result_set._csv_stream.closed - def test_csv_storage_options(self, pandas_cursor): - # Given storage_options, pandas opens the CSV output through fsspec with them. + @pytest.mark.parametrize( + ("query", "expected", "binary"), + [ + ("SELECT * FROM one_row", [(1,)], False), + ("SELECT X'01' AS value", [(b"\x01",)], True), + ], + ids=["plain", "binary"], + ) + @pytest.mark.parametrize("with_options", [False, True], ids=["none", "options"]) + def test_csv_storage_options(self, pandas_cursor, query, expected, binary, with_options): + # Given storage_options, even None, the CSV output is opened through fsspec + # with them, as pandas does, instead of the result set's filesystem. + storage_options = ( + { + "connection": pandas_cursor.connection, + "default_cache_type": "none", + "skip_instance_cache": True, + } + if with_options + else None + ) with patch.object( - S3FileSystem, "__init__", autospec=True, side_effect=S3FileSystem.__init__ - ) as init: - pandas_cursor.execute( - "SELECT * FROM one_row", - storage_options={ - "connection": pandas_cursor.connection, - "default_cache_type": "none", - "skip_instance_cache": True, - }, - ) - assert pandas_cursor.fetchall() == [(1,)] - assert init.call_count == 2 - assert init.call_args.kwargs["default_cache_type"] == "none" - assert pandas_cursor.result_set._csv_stream is None + S3FileSystem, "open", autospec=True, side_effect=S3FileSystem.open + ) as open_: + pandas_cursor.execute(query, storage_options=storage_options) + assert pandas_cursor.fetchall() == expected + file_systems = [c.args[0] for c in open_.call_args_list] + assert not [fs for fs in file_systems if fs is pandas_cursor.result_set._fs] + if with_options: + assert file_systems + assert all(fs.default_cache_type == "none" for fs in file_systems) + if not binary: + # pandas opens and closes the file itself. + assert pandas_cursor.result_set._csv_stream is None @pytest.mark.parametrize( ("pandas_cursor", "parquet_engine", "chunksize"),