Keep result-set filesystems out of the fsspec instance cache and read pandas results through one filesystem - #1001
Conversation
| return [q.strip() for q in template.render(**kwargs).split(";") if q and q.strip()] | ||
|
|
||
|
|
||
| def cached_file_systems(connection): |
There was a problem hiding this comment.
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(), thepd.read_csv()storage_options(also used by the binary-columnfilesystem_open()stream) and thepd.read_parquet()storage_options. fsspec's_Cached.__call__popsskip_instance_cachebefore__init__, and pandas passesstorage_optionstofsspec.open()/url_to_fs(), so the key never reachesS3FileSystem.__init__orS3File.pyathena/s3fs/result_set.py:filesystem_classis typedtype[AbstractFileSystem], so the fsspec metaclass handles the flag for any accepted class.pyathena/polars/result_set.py: onlypl.read_csv()goes through fsspec (polars 1.44.2io/csv/functions.py); Parquet andscan_csv()use object_store and are unaffected.pyathena/filesystem/s3_async.py: the internalS3FileSystemfollows its owner's lifetime;**kwargscannot carryskip_instance_cachebecause 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.
| read_csv_kwargs["converters"] = converters | ||
| return binary_columns | ||
|
|
||
| def _open_output_location(self, storage_options: dict[str, Any] | None, **kwargs: Any) -> Any: |
There was a problem hiding this comment.
Self-review round 1 follow-up (behavior and implementation) for the round-two repair d5326d0 — CLEAN
Scope: git diff dcc9600ee756d9608953233999a4aa4cd09cec1c..d5326d096ee6115917b62a995b06095134b2ff7c (pyathena/pandas/result_set.py, tests/pyathena/pandas/test_cursor.py), traced into _read_csv(), _open_binary_csv_stream(), _read_parquet(), PandasDataFrameIterator.close() and AthenaPandasResultSet.close().
Checked:
- Plain CSV path:
self._fs.open(path, mode="rb")gives a binaryS3File; pandas wraps a binary handle in text for the C/Python engines and passes it to pyarrow for the pyarrow engine, as it does for the file it opens from a path. Non-chunked reads close the file onExitStackexit; chunked reads hand it to the iterator (stack.pop_all()), which closes it on exhaustion, error orclose().S3File.close()is idempotent, so the later close of an already closed file is harmless. - Binary path:
AbstractFileSystem.open()wraps the binary file inTextIOWrapper(encoding="utf-8", newline="")formode="rt", matching the previousfsspec.open()call.S3FileSystemdoes not overrideopen(). - User
storage_options:None(absent) uses the result set's filesystem; any given dict, including{}, goes throughfsspec.open()aspd.read_csv()did before. No other caller of_get_csv_read_options()readsstorage_options. - Parquet: with an fsspec
filesystem, pandas 3.0.6_get_path_or_handle()keeps the path, and pyarrow rejected thes3://form, so the path isbucket/keyas in_read_parquet_schema(). A userfilesystemorstorage_optionskwarg raised before (duplicate keyword or pandasValueError) and still raises. - Errors from
self._fs.open()are inside thetryand still becomeOperationalError.
Validation: pytest -n 4 tests/pyathena/pandas/ tests/pyathena/aio/pandas/ on this source: 288 passed. New test_result_set_file_system: 3 passed; with dcc9600's pandas source, 3 failed (assert 2 == 1).
| default_block_size=self._block_size, | ||
| default_cache_type=self._cache_type, | ||
| max_workers=self._max_workers, | ||
| skip_instance_cache=True, |
There was a problem hiding this comment.
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__addsthreading.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:105already 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
AioS3FileSystemno longer shares itsdircachewith anS3FileSystemcreated with the same arguments" was misleading. The internal instance's token contains every keywordAioS3FileSystempasses, so only anS3FileSystemwith exactly those keywords shared it. Repair: the claim now names onlyskip_instance_cache=Trueinstances with the same arguments. - Comment "fsspec caches this instance itself" in
s3_async.pywas ambiguous inside theS3FileSystem(...)call; it now namesAioS3FileSystem(dcc9600). - Polars
read_csv()goes through fsspec (polars 1.44.2io/csv/functions.py) andscan_csv()/Parquet through object_store (pyathena/polars/result_set.py:619-626). Confirmed;docs/polars.md:15stays 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.
d5326d0 to
4484e16
Compare
| self._open_binary_csv_stream(binary_columns, storage_options) | ||
| ) | ||
| source = self._csv_stream | ||
| elif "storage_options" not in read_csv_kwargs: |
There was a problem hiding this comment.
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
-
P2 — Explicit CSV storage options still retain connections. pyathena/pandas/result_set.py:751
Withstorage_options={"connection": conn, "default_cache_type": "none"}, this branch creates a filesystem withoutskip_instance_cache=True. Its cachedstorage_optionsretainsconn, and itsdircachesurvives disposal of the result set. The existingtest_binary_as_pandasuses exactly these options but does not check cache membership. This is the original defect left unfixed for explicit options, rather than a new regression. -
P2 — Chunked CSV reads discard the owner responsible for closing underlying streams. pyathena/pandas/result_set.py:571
For gzip-compressed content at a.csvoutput location, usingchunksize=1andstorage_options={"connection": conn, "compression": "gzip", "skip_instance_cache": True}returns aGzipFilefrom fsspec’sOpenFile. The subsequentstack.pop_all()discards theOpenFilecleanup context. Iterator exhaustion, close, or a read exception closes only the gzip wrapper;GzipFile.close()does not close its supplied underlyingS3File. That file and its executor remain open while the retained wrapper references them. Previously pandas retained and closed theOpenFile. New regression for ordinary CSV reads; analogous binary-path behavior predates this diff. Transfer the cleanup context alongside the stream. -
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, andstorage_options={}, metadata lookup succeeds through the connection. Previously pandas retried the separate filesystem open withanon=TrueafterNoCredentialsError. Opening outside pandas now raisesOperationalErrorinstead. 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:
- Explicit
storage_optionswithconnectionstay cached. Confirmed, but not acted on. The options are the caller's: pandas also cached them before this PR, and addingskip_instance_cache: Trueto them avoids it, as the newtest_csv_storage_optionsdoes. Pre-existing and outside Cached S3FileSystem instances keep every connection alive and grow their dircache without bound #978, which covers the filesystems PyAthena creates. - Chunked read with
compressionin the userstorage_optionsloses theOpenFileowner. Confirmed for d5326d0:stack.enter_context(OpenFile)returns theGzipFile,pop_all()drops theOpenFile, andGzipFile.close()leaves the suppliedfileobjopen. This was a regression for non-binary reads. The same pattern in the binary-column path predates this PR and is unchanged. - pandas' anonymous retry is lost for user
storage_options. Confirmed: pandas 3.0.6_get_filepath_or_buffer()retriesfsspec.open()withanon=TrueafterNoCredentialsError/PermissionErrorfors3://paths, and opening the file ourselves skipped that. This was a regression for non-binary reads with user options. The default path never benefited, becauseS3FileSystem(connection=...)ignoresanon.
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.
There was a problem hiding this comment.
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._fsand the stream ownership is unchanged from d5326d0. - Non-binary reads with user
storage_optionsare again exactly master's call:pd.read_csv(path, storage_options=...). pandas owns its handles there, and_csv_streamstaysNoneas on master. - The binary path with given options calls
fsspec.open()with them, as on master. The only difference is an explicitstorage_options=None, which now uses the result set's filesystem instead of fsspec's defaultS3FileSystemwithout 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.
There was a problem hiding this comment.
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_optionsabsent,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
-
P2 — Binary CSV still conflates explicit
storage_options=Nonewith omission.
pyathena/pandas/result_set.py:567 loses that distinction, so line 745 selects the connection-backed filesystem. Before the change, explicitNonebecame{}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-Nonecoverage; 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:
- Round 1: only the binary branch's option resolution changed; no other caller.
- Round 2: the PR description's statement that given options go through fsspec as before now holds for an explicit
Nonetoo. - Upstream since 55af09a (Make S3Object follow the mapping contract #991 S3Object mapping, Reflect the field and value types of top-level STRUCT and MAP columns #995 SQLAlchemy reflection, Treat dircache entries removed by another thread as cache misses #1002 dircache eviction race) does not touch filesystem construction or instance caching.
There was a problem hiding this comment.
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
e1f4d470and commits 1–5 are unchanged. _read_csv,_open_binary_csv_stream, andPandasDataFrameIterator: plain/binary CSV, full/chunked reads, and C/Python/PyArrow engine selection.storage_optionsomitted,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.pyand related existing lifecycle tests. - Default Parquet filesystem reuse, result-set cache isolation, and
AioS3FileSysteminternal isolation. - Upstream filesystem changes from
55af09a2to9b2f0337, including dircache operations andS3Objectmapping 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.
dbab34f to
e1f4d47
Compare
… 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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
e1f4d47 to
75ef9c2
Compare
WHAT
Result sets and
AioS3FileSystemno longer put the filesystems they create into fsspec's instance cache, and the pandas result set reads through a single filesystem:S3FileSystemwithskip_instance_cache=Trueand reads through it:pd.read_csv()as a file, for both the plain and the binary-column paths. Whenstorage_optionsare given in the read options (execute(..., storage_options=...)), evenNone, the path and these options go topd.read_csv()(or, for binary columns,fsspec.open()) as before.pd.read_parquet()receives that filesystem asfilesystem, with the unload location asbucket/key(pyarrow rejects thes3://form with an fsspec filesystem: "GetFileInfo() yielded path ... which is outside base dir").pd.read_csv()andpd.read_parquet()each built anotherS3FileSystemfromstorage_options.filesystem_class) withskip_instance_cache=True.skip_instance_cache=Truein thestorage_optionsofpl.read_csv(). The Parquet andscan_csv()paths use Polars' native object store, not fsspec, and are unchanged.AioS3FileSystemalways creates its internalS3FileSystemwithskip_instance_cache=True. fsspec already caches theAioS3FileSystemitself unlessskip_instance_cache=Trueis given to it, so the internal instance follows the lifetime of its owner.A result set's filesystem, its S3 client and its
dircacheare now released with the result set, and a closed connection can be garbage-collected. Each result set builds oneS3FileSystem(measured: one per query for pandas CSV, pandasunload=True, Polars and S3FS; pandas previously built two).Behavior changes (release-note):
pyathena/result_set.py:105), so a query now opens one more HTTPS connection to S3 than before. Measured from a laptop to us-west-2: the first request of a new client took a median of about 450 ms, a request on an open connection about 150 ms; building the client took about 2.3 ms. The cost inside the region was not measured.AioS3FileSystem(skip_instance_cache=True)instances created with the same arguments no longer share their internalS3FileSystemanddircache.storage_options, the pandas result set owns the CSV file it reads for every CSV result, not only for binary columns, and closes it with the chunk iterator (PandasDataFrameIterator.close()).storage_optionsgiven by the user are used as given; aconnectionin them is cached by fsspec unless they also containskip_instance_cache: True(unchanged).WHY
Closes #978.
S3FileSystemis cachable and fsspec keys its instance cache on the constructor arguments, which include theconnectionobject. A program that opens a connection per query kept one cached filesystem, with its connection, boto3 session and clients, per connection; one long-lived connection reused a single cached instance whosedircachegained an entry for every result file. Closed #417 (memory growth with about 1,000PandasCursorqueries per hour) may have been this problem.TEST
Tested commit:
75ef9c266b513ca16fe613929e74c87a85db4f88(based on4e7b55ff).just lint: passed.TestAioS3FileSystem.test_internal_file_system_not_cached(offline): twoskip_instance_cache=Trueinstances do not share the internal filesystem ordircache, and no internal filesystem entersS3FileSystem._cache.TestPandasCursor.test_result_set_file_system(AWS; CSV, chunked CSV,unload=True): oneS3FileSystemis built per result set, no filesystem holding the connection is in theS3FileSystemorAioS3FileSysteminstance cache, and the CSV file is closed afterfetchall().TestPandasCursor.test_csv_storage_options(AWS; plain/binary ×None/options): withstorage_optionsin the read options, the result set's filesystem never opens the output, given options reach the filesystem that does, and pandas opens the plain file itself.test_result_set_file_system_not_cachedinTestPolarsCursorandTestS3FSCursor(AWS): no filesystem holding the connection is in either instance cache.skip_instance_cache-only pandas change, the pandas filesystem test fails withassert 2 == 1in all three cases.test_csv_storage_optionsfails on d232af0's pandas source (none-binary,none-plain,options-plain) and on c42670b's (none-binary).TestAioS3FileSystemrun first in the same process (it registersAioS3FileSystemfors3and does not restore it), the Polars test fails when itsskip_instance_cacheoption is removed.uv run --env-file .env 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.pyone1f4d470(same patches, previous base9b2f0337): 471 passed, 1 skipped.4e7b55ff(the only range-diff change is the context of the AIO test),pytest -n 4 tests/pyathena/filesystem/test_s3_async.py tests/pyathena/{pandas,polars,s3fs}/test_cursor.py: 339 passed, 1 skipped.🤖 Generated with Claude Code