-
Notifications
You must be signed in to change notification settings - Fork 116
Let execute() override cursor read options and accept them in the AsyncPandasCursor constructor #1022
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
Let execute() override cursor read options and accept them in the AsyncPandasCursor constructor #1022
Changes from all commits
30548fd
180e716
000e5c2
289b0f7
84d3ab3
b88146b
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 |
|---|---|---|
|
|
@@ -79,6 +79,9 @@ def __init__( | |
| chunksize: int | None = None, | ||
| result_reuse_enable: bool = False, | ||
| result_reuse_minutes: int = CursorIterator.DEFAULT_RESULT_REUSE_MINUTES, | ||
| block_size: int | None = None, | ||
| cache_type: str | None = None, | ||
| auto_optimize_chunksize: bool = False, | ||
| **kwargs, | ||
| ) -> None: | ||
| """Initialize an AsyncPandasCursor. | ||
|
|
@@ -99,8 +102,13 @@ def __init__( | |
| unload: Whether to wrap queries in ``UNLOAD`` and read the Parquet output. | ||
| engine: Parsing engine (``auto``, ``c``, ``python``, or ``pyarrow``). | ||
| chunksize: Number of rows per DataFrame chunk when reading CSV results. | ||
| If set, it takes precedence over ``auto_optimize_chunksize``. | ||
| result_reuse_enable: Whether to enable Athena query result reuse. | ||
| result_reuse_minutes: Maximum age of a reused query result in minutes. | ||
| block_size: Default block size of the S3 filesystem that reads the results. | ||
| cache_type: Default cache type of the S3 filesystem that reads the results. | ||
| auto_optimize_chunksize: Whether to choose a chunk size from the size of the | ||
| CSV result file when ``chunksize`` is None. | ||
| **kwargs: Other cursor arguments, such as ``connection`` and ``converter``, | ||
| passed to ``AsyncCursor.__init__``. | ||
| """ | ||
|
|
@@ -122,6 +130,9 @@ def __init__( | |
| self._unload = unload | ||
| self._engine = engine | ||
| self._chunksize = chunksize | ||
| self._block_size = block_size | ||
| self._cache_type = cache_type | ||
| self._auto_optimize_chunksize = auto_optimize_chunksize | ||
|
|
||
| @staticmethod | ||
| @override | ||
|
|
@@ -170,6 +181,11 @@ def _collect_result_set( | |
| unload_location=unload_location, | ||
| engine=kwargs.pop("engine", self._engine), | ||
| chunksize=kwargs.pop("chunksize", self._chunksize), | ||
| block_size=kwargs.pop("block_size", self._block_size), | ||
| cache_type=kwargs.pop("cache_type", self._cache_type), | ||
| auto_optimize_chunksize=kwargs.pop( | ||
| "auto_optimize_chunksize", self._auto_optimize_chunksize | ||
| ), | ||
| result_set_type_hints=result_set_type_hints, | ||
| **kwargs, | ||
| ) | ||
|
|
@@ -215,6 +231,12 @@ def execute( | |
| :class:`~pyathena.options.ExecuteOptions` instance. Individual | ||
| keyword arguments take precedence over ``options`` fields. | ||
| **kwargs: Additional pandas read_csv/read_parquet parameters. | ||
| ``engine``, ``chunksize``, ``block_size``, ``cache_type``, and | ||
| ``auto_optimize_chunksize`` override the cursor's values for this query. | ||
| ``storage_options`` and, for UNLOAD results, ``filesystem`` replace | ||
| PyAthena's S3 filesystem (see | ||
| :class:`~pyathena.pandas.result_set.AthenaPandasResultSet`). | ||
| ``max_workers`` sets the number of S3 read workers for this query. | ||
|
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 two (claims, existing callers, documentation, evidence) — FINDINGS (PR description only; no code change) Claims checked:
Findings (repaired in the PR description):
Limits: the rest of the pandas/Polars suites and |
||
|
|
||
| Returns: | ||
| Tuple of (query_id, future) where future resolves to AthenaPandasResultSet. | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -310,6 +310,10 @@ def __init__( | |
| result_set_type_hints: Athena type signatures for complex-type columns, | ||
| keyed by column name (case-insensitive) or zero-based column index. | ||
| **kwargs: Additional arguments passed to pandas.read_csv/read_parquet. | ||
| A given ``storage_options``, even None, replaces PyAthena's S3 filesystem | ||
|
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 two (pandas UNLOAD) — FINDINGS (PR description only)
Finding: the per-call |
||
| for reading the result files, and so does ``filesystem`` for UNLOAD results. | ||
| The UNLOAD manifest is still read with the connection's S3 client, and the | ||
|
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), pandas UNLOAD change: FINDINGS (P2 pre-existing, P3 docstring) → P3 repaired in c8d4f2d Reviewer: Codex CLI 0.160.0, model Covered (reviewer): keyword forwarding and exception propagation in all three pandas cursors; default and overridden Parquet reads; CSV parity; manifest/schema reads; the locked pandas 3.0.6 and pyarrow 25.0.1 sources; both new test groups. No functional defect in the argument selection; the four non-default mocked cases and both live cases fail on the previous head, and the tests assert the pandas call boundary and the returned records. Findings:
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 (
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 (relayed): CLEAN
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. Rebase record:
|
||
| schema with PyAthena's filesystem. | ||
| """ | ||
| super().__init__( | ||
| connection=connection, | ||
|
|
@@ -778,23 +782,22 @@ def _read_parquet(self, engine) -> DataFrame: | |
| self._unload_location = "/".join(self._data_manifest[0].split("/")[:-1]) + "/" | ||
|
|
||
| if engine == "pyarrow": | ||
| # 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, | ||
| } | ||
| kwargs: dict[str, Any] = {"use_threads": True, **self._kwargs} | ||
| # Given storage_options, even None, pandas opens the files itself, | ||
| # as for CSV results. | ||
| if "filesystem" not in kwargs and "storage_options" not in 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. Self-review round one (pandas UNLOAD Covered:
|
||
| kwargs["filesystem"] = self._fs | ||
| if kwargs.get("filesystem") is None: | ||
| unload_location = self._unload_location | ||
| else: | ||
| # pyarrow takes the path without the scheme with a filesystem. | ||
| bucket, key = parse_output_location(self._unload_location) | ||
| unload_location = f"{bucket}/{key}" | ||
| else: | ||
| raise ProgrammingError("Engine must be `pyarrow`.") | ||
| kwargs.update(self._kwargs) | ||
|
|
||
| try: | ||
| return pd.read_parquet( | ||
| unload_location, | ||
| engine=self._engine, | ||
| filesystem=self._fs, | ||
| **kwargs, | ||
| ) | ||
| return pd.read_parquet(unload_location, engine=self._engine, **kwargs) | ||
| except Exception as e: | ||
| _logger.exception(f"Failed to read {self.output_location}.") | ||
| raise OperationalError(*e.args) from e | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -189,6 +189,11 @@ def execute( | |
| :class:`~pyathena.options.ExecuteOptions` instance. Individual | ||
| keyword arguments take precedence over ``options`` fields. | ||
| **kwargs: Additional execution parameters passed to Polars read functions. | ||
| ``block_size``, ``cache_type``, ``max_workers``, and ``chunksize`` | ||
| override the cursor's values for this query. | ||
| Read function arguments replace the ones the result set chooses, such as | ||
| ``separator``, ``has_header``, ``schema_overrides``, and ``storage_options`` | ||
| (see :class:`~pyathena.polars.result_set.AthenaPolarsResultSet`). | ||
|
|
||
| Returns: | ||
| Self reference for method chaining. | ||
|
|
@@ -229,10 +234,10 @@ def execute( | |
| retry_config=self._retry_config, | ||
| unload=self._unload, | ||
| unload_location=unload_location, | ||
| block_size=self._block_size, | ||
| cache_type=self._cache_type, | ||
| max_workers=self._max_workers, | ||
| chunksize=self._chunksize, | ||
| block_size=kwargs.pop("block_size", self._block_size), | ||
|
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 one — out-of-scope note (not changed in this PR): the same "cursor attribute plus
|
||
| cache_type=kwargs.pop("cache_type", self._cache_type), | ||
| max_workers=kwargs.pop("max_workers", self._max_workers), | ||
| chunksize=kwargs.pop("chunksize", self._chunksize), | ||
| result_set_type_hints=options.result_set_type_hints, | ||
| **kwargs, | ||
| ) | ||
|
|
||
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 one (behavior, failure paths, data/resource boundaries, simplicity and test quality) — FINDINGS
Scope:
git diff 785be364e380a7088b82ecf3b5bcb3a5bf01b9ca..bb006761a94daa69f858347479d6d5b03934e295— all 12 files: the six pandas/Polars cursors and their test files.Finding (repaired in d542edb): the new
block_sizeandcache_typeparameters were inserted beforeresult_reuse_enable, so a caller passingresult_reuse_enableorresult_reuse_minutespositionally toAsyncPandasCursor(...)would have setblock_size/cache_typeinstead. All three new parameters are now appended afterresult_reuse_minutes; keyword callers (includingconnect(cursor_kwargs=...)) are unaffected either way.Checked without findings:
kwargs.pop()reads a per-call dict (execute()'s own**kwargs, or the dict submitted to the executor forAsyncPandasCursor/AsyncPolarsCursor), so popping does not leak into later calls.**kwargsbefore, so per-call values that already worked (e.g.AsyncPandasCursor.execute(block_size=...)) behave the same.AsyncPandasCursor/AsyncPolarsCursor:max_workersgiven toexecute()sets only the S3 read workers; the executor keeps the constructor size.execute(block_size=None)now overrides the cursor value with None, matching the existingPandasCursorpop semantics.test_read_optionscases fail on the base (5TypeErrors, plus theAsyncPandasCursorconstructor case) and assert the actual result-set arguments; the livetest_auto_optimize_chunksizecovers the real chunked read throughcursor_kwargs.