Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,8 @@ def __init__(
allow_bucket_creation=allow_bucket_creation,
allow_bucket_deletion=allow_bucket_deletion,
version_aware=version_aware,
# fsspec caches the AioS3FileSystem itself when caching is wanted.
skip_instance_cache=True,
**kwargs,
)
# Share dircache for cache coherence between async and sync instances
Expand Down
58 changes: 29 additions & 29 deletions pyathena/pandas/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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 []
Expand Down Expand Up @@ -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,

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 2 (claims, compatibility, operations) — FINDINGS, repaired

Scope: git diff 0b656b8d6de41b1f87bdb4d1e670e00a480df895..28738a8dc7c963d8c28393ad08772d34d219512e for the initial pass; the PR description, commit messages and code comments.

Claims checked:

  • "Result sets build one filesystem per query / pandas builds two": measured by wrapping S3FileSystem.__init__ around live queries. Before the repair it was pandas 2 per query (CSV and unload), Polars 1 and S3FS 1; after d5326d0 it is 1 for all four.
  • "fsspec caches per connection and thread": fsspec 2026.9.0 _Cached.__call__ adds threading.get_ident() to the token for sync filesystems. Confirmed.
  • Client build cost about 2.3 ms: measured locally.
  • Finding: the description presented per-query client creation as a pure CPU cost. It also gives up HTTPS connection reuse. Measured from a laptop to us-west-2: median 447 ms for the first request on a new client, 154 ms on an open connection. It also missed that pyathena/result_set.py:105 already builds a new S3 client per result set, so the real increase is one new connection per query. With the pandas second filesystem it was two.
    Repair: the maintainer chose to reduce pandas to one filesystem (d5326d0). The description now states the measured connection cost and its limits (not measured inside the region).
  • Finding: "An AioS3FileSystem no longer shares its dircache with an S3FileSystem created with the same arguments" was misleading. The internal instance's token contains every keyword AioS3FileSystem passes, so only an S3FileSystem with exactly those keywords shared it. Repair: the claim now names only skip_instance_cache=True instances with the same arguments.
  • Comment "fsspec caches this instance itself" in s3_async.py was ambiguous inside the S3FileSystem(...) call; it now names AioS3FileSystem (dcc9600).
  • Polars read_csv() goes through fsspec (polars 1.44.2 io/csv/functions.py) and scan_csv()/Parquet through object_store (pyathena/polars/result_set.py:619-626). Confirmed; docs/polars.md:15 stays accurate.
  • PyAthena - Memory Issue #417 is cited only as "may have been this problem".

Compatibility: no public signature changed. Users passing storage_options keep the fsspec path, and filesystem_class is typed type[AbstractFileSystem]. The docs (docs/pandas.md, docs/s3fs.md, docs/aio.md) make no claim about instance caching.

)

@property
Expand Down Expand Up @@ -558,14 +561,22 @@ def _read_csv(self) -> TextFileReader | DataFrame:

try:
with ExitStack() as stack:
source: str | TextIOWrapper = self.output_location
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) or {}
self._csv_stream = stack.enter_context(
# 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)
)
source = self._csv_stream
elif "storage_options" not in read_csv_kwargs:

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) — FINDINGS

Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, sandbox read-only, session 01a1013c-da33-7313-b62b-ed4961777a1f.
Scope: git diff 0b656b8d6de41b1f87bdb4d1e670e00a480df895..d5326d096ee6115917b62a995b06095134b2ff7c, reviewed in a detached snapshot at d5326d0. The prompt carried the issue's intended behavior but no PR number, description, commit messages or earlier findings. The reviewer was told not to edit, build, run tests, use GitHub or the network. The review snapshot and the PR worktree were both unchanged afterwards. This is a static review.

Reviewer output (verbatim)

Covered surfaces:

  • pandas CSV: plain/chunked/binary reads, C/Python/pyarrow engines, .txt, read kwargs, and stream cleanup.
  • pandas Parquet: filesystem=, scheme-less paths, and schema reads.
  • S3FS/Polars result sets and sync, Future-based async, and aio callers.
  • fsspec cache semantics, installed dependency sources, supported ranges, and new tests’ failure and ordering behavior.

FINDINGS

  1. P2 — Explicit CSV storage options still retain connections. pyathena/pandas/result_set.py:751
    With storage_options={"connection": conn, "default_cache_type": "none"}, this branch creates a filesystem without skip_instance_cache=True. Its cached storage_options retains conn, and its dircache survives disposal of the result set. The existing test_binary_as_pandas uses exactly these options but does not check cache membership. This is the original defect left unfixed for explicit options, rather than a new regression.

  2. P2 — Chunked CSV reads discard the owner responsible for closing underlying streams. pyathena/pandas/result_set.py:571
    For gzip-compressed content at a .csv output location, using chunksize=1 and storage_options={"connection": conn, "compression": "gzip", "skip_instance_cache": True} returns a GzipFile from fsspec’s OpenFile. The subsequent stack.pop_all() discards the OpenFile cleanup context. Iterator exhaustion, close, or a read exception closes only the gzip wrapper; GzipFile.close() does not close its supplied underlying S3File. That file and its executor remain open while the retained wrapper references them. Previously pandas retained and closed the OpenFile. New regression for ordinary CSV reads; analogous binary-path behavior predates this diff. Transfer the cleanup context alongside the stream.

  3. P2 — Ordinary CSV reads lose pandas’ anonymous-access fallback. pyathena/pandas/result_set.py:572
    Given a connection with explicit credentials, publicly readable results, no ambient AWS credentials, and storage_options={}, metadata lookup succeeds through the connection. Previously pandas retried the separate filesystem open with anon=True after NoCredentialsError. Opening outside pandas now raises OperationalError instead. New regression for ordinary CSV reads. Preserve that fallback when handling explicit storage options.

The new cache assertions would catch the original default-path defects. Fresh connection identities and checking both filesystem caches avoid the identified registry-order concern under pytest-xdist. They do not cover the explicit-options cases above.

Static review only; no edits, builds, tests, GitHub calls, or network access.

Author verification:

  1. Explicit storage_options with connection stay cached. Confirmed, but not acted on. The options are the caller's: pandas also cached them before this PR, and adding skip_instance_cache: True to them avoids it, as the new test_csv_storage_options does. Pre-existing and outside Cached S3FileSystem instances keep every connection alive and grow their dircache without bound #978, which covers the filesystems PyAthena creates.
  2. Chunked read with compression in the user storage_options loses the OpenFile owner. Confirmed for d5326d0: stack.enter_context(OpenFile) returns the GzipFile, pop_all() drops the OpenFile, and GzipFile.close() leaves the supplied fileobj open. This was a regression for non-binary reads. The same pattern in the binary-column path predates this PR and is unchanged.
  3. pandas' anonymous retry is lost for user storage_options. Confirmed: pandas 3.0.6 _get_filepath_or_buffer() retries fsspec.open() with anon=True after NoCredentialsError/PermissionError for s3:// paths, and opening the file ourselves skipped that. This was a regression for non-binary reads with user options. The default path never benefited, because S3FileSystem(connection=...) ignores anon.

Repair for 2 and 3 (4484e16, on this line): when storage_options are in the read options, the path and the options go to pd.read_csv() unchanged, as on master. Only without them is the output opened through the result set's filesystem. The binary path keeps master's behavior for given options; without them it uses the result set's filesystem. New test_csv_storage_options fails on d5326d0 (_csv_stream is the self-opened file) and passes on 4484e16.

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 of the repair 4484e16 (rounds 1 and 2) — CLEAN

Scope: git range-diff 0b656b8d6de41b1f87bdb4d1e670e00a480df895..d5326d096ee6115917b62a995b06095134b2ff7c 55af09a2..4484e16b90110154637b22657082c33dbcb82cc1. Commits 1-4 are unchanged (=), commit 5 is new. Upstream between the bases (#998, #999) touches s3.py/s3_async.py range reads and transaction part limits, not filesystem construction or caching.

Round 1 (behavior):

  • Without user options, the CSV output is still opened through self._fs and the stream ownership is unchanged from d5326d0.
  • Non-binary reads with user storage_options are again exactly master's call: pd.read_csv(path, storage_options=...). pandas owns its handles there, and _csv_stream stays None as on master.
  • The binary path with given options calls fsspec.open() with them, as on master. The only difference is an explicit storage_options=None, which now uses the result set's filesystem instead of fsspec's default S3FileSystem without the connection.
  • _open_output_location() was removed, so no caller is left.

Round 2 (claims): the PR description now says that user options keep master's path, and that caching of a connection inside user options is unchanged (independent-review item 1). The test list and the tested commit are current.

Validation on 4484e16: just lint passed. pytest -n 4 tests/pyathena/pandas/ tests/pyathena/aio/pandas/ tests/pyathena/filesystem/test_s3_async.py tests/pyathena/polars/test_cursor.py tests/pyathena/s3fs/test_cursor.py: 468 passed, 1 skipped.

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 review 1 (relayed) — FINDINGS

Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, sandbox read-only, session 01a10153-3899-75a2-a769-7b2709c32935.
Scope: the patch series 0b656b8d..d5326d09 compared with 55af09a2..4484e16b (range-diff, the new commit and the full diff), reviewed in a detached snapshot at 4484e16. The prompt carried no PR number, description or earlier findings. The snapshot and the PR worktree were both unchanged afterwards. This is a static review.

Reviewer output (verbatim)

Covered at 4484e16b:

  • Supplied range-diff and follow-up diff; commits 1–4 are unchanged.
  • Plain, chunked, and binary CSV paths; C, Python, and PyArrow engine selection.
  • storage_options absent, None, {}, and non-empty; pandas fallbacks and handle ownership.
  • Stream cleanup on completion, exhaustion, explicit close, cursor reuse, and errors.
  • New tests and existing binary lifecycle tests.
  • Parquet filesystem reuse, instance-cache isolation, and upstream filesystem changes. The rebase adds no conflicting cache or ownership behavior.

FINDINGS

  1. P2 — Binary CSV still conflates explicit storage_options=None with omission.
    pyathena/pandas/result_set.py:567 loses that distinction, so line 745 selects the connection-backed filesystem. Before the change, explicit None became {} and opened through fsspec defaults.

    Concrete failure: consume a large, chunked VARBINARY result with storage_options=None; the connection uses fixed temporary credentials that expire during consumption, while ambient credentials remain valid. Subsequent S3 reads now fail instead of continuing through the ambient credentials. Plain CSV preserves the previous behavior after this follow-up.

    Origin: introduced by commit 4 (d232af0a), pre-existing relative to commit 5 and still unresolved. Preserve key presence in the binary branch and add explicit-None coverage; the new test exercises only a non-empty options dictionary.

No newly introduced defect was found in commit 5. This was static review only; no files were changed, tests/builds run, or network accessed.

Author verification: confirmed. On master, an explicit storage_options=None became {} in the binary path (fsspec defaults), and plain reads after 4484e16 pass None to pandas (also fsspec defaults). The binary path at 4484e16 used the result set's filesystem instead. The regression came from d232af0.

Repair (dbab34f, rebased as e1f4d47 onto 9b2f033): the binary path uses the result set's filesystem only when the read options have no storage_options key; an explicit None becomes {} as on master. test_csv_storage_options now covers plain/binary × None/options:

  • it asserts that the result set's filesystem never opens the file, that given options reach the filesystem that does (default_cache_type == "none"), and that pandas opens the plain file itself (_csv_stream is None);
  • it fails on 4484e16 (none-binary) and on d232af0 (none-binary, none-plain, options-plain), and passes on the repair.

Self-review of the repair:

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 review 2 (relayed) — CLEAN

Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, sandbox read-only, session 01a10160-8741-7b70-8509-b66970b82e69.
Scope: the patch series 55af09a2..4484e16b compared with 9b2f0337..e1f4d470 (range-diff: commits 1-5 unchanged, commit 6 new), the new commit, and the full diff 9b2f0337..e1f4d470, reviewed in a detached snapshot at e1f4d47. The prompt carried no PR number, description or earlier findings. The snapshot was unchanged afterwards. This is a static review.

Reviewer output (verbatim)

Covered:

  • Supplied patch comparison and new commit; confirmed HEAD is e1f4d470 and commits 1–5 are unchanged.
  • _read_csv, _open_binary_csv_stream, and PandasDataFrameIterator: plain/binary CSV, full/chunked reads, and C/Python/PyArrow engine selection.
  • storage_options omitted, None, {}, and non-empty; compared opening, fallback, and ownership behavior with the pre-series source.
  • Stream cleanup on construction/read failures, exhaustion, early close, context exit, and cursor reuse.
  • New/changed tests in tests/pyathena/pandas/test_cursor.py and related existing lifecycle tests.
  • Default Parquet filesystem reuse, result-set cache isolation, and AioS3FileSystem internal isolation.
  • Upstream filesystem changes from 55af09a2 to 9b2f0337, including dircache operations and S3Object mapping behavior, against installed dependency sources.

CLEAN

No actionable findings in this follow-up scope. The explicit-None repair restores the pre-series binary CSV opening behavior without disrupting omitted-option filesystem reuse or stream ownership.

Static review only; no files changed, tests/builds run, or network access performed.

# 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(source, **read_csv_kwargs)
if not isinstance(result, pd.DataFrame):
# The chunk iterator takes ownership of the stream.
Expand Down Expand Up @@ -608,12 +619,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,
},
"chunksize": chunksize,
"engine": csv_engine,
}
Expand Down Expand Up @@ -736,19 +741,17 @@ def _configure_binary_csv_read(
return binary_columns

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."""
text_options: dict[str, Any] = {"mode": "rt", "encoding": "utf-8", "newline": ""}
with ExitStack() as stack:
source = stack.enter_context(
filesystem_open(
self.output_location,
mode="rt",
encoding="utf-8",
newline="",
**storage_options,
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="")
Expand All @@ -765,7 +768,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,
}
Expand All @@ -777,12 +782,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,
},
filesystem=self._fs,
**kwargs,
)
except Exception as e:
Expand Down
3 changes: 3 additions & 0 deletions pyathena/polars/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions pyathena/s3fs/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
14 changes: 14 additions & 0 deletions tests/pyathena/filesystem/test_s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"):
Expand Down
61 changes: 61 additions & 0 deletions tests/pyathena/pandas/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,13 @@
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:
Expand Down Expand Up @@ -230,6 +232,65 @@ 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", "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(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.
# 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(
("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, "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"),
[
Expand Down
8 changes: 8 additions & 0 deletions tests/pyathena/polars/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
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:
Expand All @@ -36,6 +37,13 @@ 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 cached_file_systems(polars_cursor.connection)

@pytest.mark.parametrize(
"polars_cursor",
[{"cursor_kwargs": {"unload": False}}, {"cursor_kwargs": {"unload": True}}],
Expand Down
8 changes: 8 additions & 0 deletions tests/pyathena/s3fs/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
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:
Expand All @@ -24,6 +25,13 @@ 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 cached_file_systems(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
Expand Down
12 changes: 12 additions & 0 deletions tests/pyathena/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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):

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 1 (behavior and implementation) — FINDINGS, repaired

Scope: git diff 0b656b8d6de41b1f87bdb4d1e670e00a480df895..eebef4099242b968b0db6f565b5c679b219cc3d3 (full initial diff, 8 files).

Covered:

  • pyathena/pandas/result_set.py: _create_s3_file_system(), the pd.read_csv() storage_options (also used by the binary-column filesystem_open() stream) and the pd.read_parquet() storage_options. fsspec's _Cached.__call__ pops skip_instance_cache before __init__, and pandas passes storage_options to fsspec.open() / url_to_fs(), so the key never reaches S3FileSystem.__init__ or S3File.
  • pyathena/s3fs/result_set.py: filesystem_class is typed type[AbstractFileSystem], so the fsspec metaclass handles the flag for any accepted class.
  • pyathena/polars/result_set.py: only pl.read_csv() goes through fsspec (polars 1.44.2 io/csv/functions.py); Parquet and scan_csv() use object_store and are unaffected.
  • pyathena/filesystem/s3_async.py: the internal S3FileSystem follows its owner's lifetime; **kwargs cannot carry skip_instance_cache because fsspec already consumed it.
  • Sync, async and aio cursors build the same result-set classes, so they share the change. Arrow uses pyarrow's own filesystem and is not affected.

Finding: TestAioS3FileSystem registers AioS3FileSystem for s3/s3a with a class-scoped fixture (tests/pyathena/filesystem/test_s3_async.py:31-38) and never restores it. In an xdist worker that ran it first, pandas and Polars read through a cached AioS3FileSystem, which the new tests did not inspect (S3FileSystem._cache only), so a regression in the Polars storage_options could pass there.
Repair (28738a8): this cached_file_systems() helper checks both instance caches. Verified with pytest -n 0 running TestAioS3FileSystem::test_parse_path first: the Polars test fails with its skip_instance_cache line removed and passes with it.

Out of scope: the leaked registration itself is pre-existing test isolation behavior.

Validation after repair: just lint passed; pytest -n 2 tests/pyathena/{pandas,polars,s3fs}/test_cursor.py -k "file_system_not_cached or binary_null_vs_empty": 8 passed.

"""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")


Expand Down
Loading