-
Notifications
You must be signed in to change notification settings - Fork 116
Keep result-set filesystems out of the fsspec instance cache and read pandas results through one filesystem #1001
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
9375a34
9f5ab6f
7b82394
0f2a226
c42670b
75ef9c2
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 [] | ||
|
|
@@ -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 | ||
|
|
@@ -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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed) — FINDINGS Reviewer: Codex CLI 0.160.0, model Reviewer output (verbatim)Covered surfaces:
FINDINGS
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:
Repair for 2 and 3 (4484e16, on this line): when
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review of the repair 4484e16 (rounds 1 and 2) — CLEAN Scope: Round 1 (behavior):
Round 2 (claims): the PR description now says that user options keep master's path, and that caching of a Validation on 4484e16:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 Reviewer output (verbatim)Covered at
FINDINGS
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 Repair (dbab34f, rebased as e1f4d47 onto 9b2f033): the binary path uses the result set's filesystem only when the read options have no
Self-review of the repair:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 Reviewer output (verbatim)Covered:
CLEAN No actionable findings in this follow-up scope. The explicit- 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. | ||
|
|
@@ -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, | ||
| } | ||
|
|
@@ -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="") | ||
|
|
@@ -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, | ||
| } | ||
|
|
@@ -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: | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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): | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round 1 (behavior and implementation) — FINDINGS, repaired Scope: Covered:
Finding: Out of scope: the leaked registration itself is pre-existing test isolation behavior. Validation after repair: |
||
| """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") | ||
|
|
||
|
|
||
|
|
||
There was a problem hiding this comment.
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..28738a8dc7c963d8c28393ad08772d34d219512efor the initial pass; the PR description, commit messages and code comments.Claims checked:
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._Cached.__call__addsthreading.get_ident()to the token for sync filesystems. Confirmed.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).
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.s3_async.pywas ambiguous inside theS3FileSystem(...)call; it now namesAioS3FileSystem(dcc9600).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.Compatibility: no public signature changed. Users passing
storage_optionskeep the fsspec path, andfilesystem_classis typedtype[AbstractFileSystem]. The docs (docs/pandas.md,docs/s3fs.md,docs/aio.md) make no claim about instance caching.