-
Notifications
You must be signed in to change notification settings - Fork 116
Keep the whole pandas and Polars results reusable #1005
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
f41ca3c
a26c763
c857dd3
7ce2bf4
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 |
|---|---|---|
|
|
@@ -343,17 +343,27 @@ def __init__( | |
| d[0] for d in description if d[1] in ("time", "time with time zone") | ||
| ] | ||
|
|
||
| import pandas as pd | ||
|
|
||
| # The whole result when it was not read in chunks. | ||
| self._df: DataFrame | None = None | ||
| if self.state == AthenaQueryExecution.STATE_SUCCEEDED and self.output_location: | ||
| df = self._as_pandas() | ||
| result = self._as_pandas() | ||
| trunc_date = _no_trunc_date if self.is_unload else self._trunc_date | ||
| self._df_iter = PandasDataFrameIterator(df, trunc_date, self._csv_stream) | ||
| if isinstance(result, pd.DataFrame): | ||
| self._df = trunc_date(result) | ||
| else: | ||
| self._df_iter = PandasDataFrameIterator(result, trunc_date, self._csv_stream) | ||
| elif self.state == AthenaQueryExecution.STATE_SUCCEEDED: | ||
| df = self._as_pandas_from_api() | ||
| self._df_iter = PandasDataFrameIterator(df, self._trunc_date) | ||
| # GetQueryResults values are already converted and need no time truncation. | ||
| self._df = self._as_pandas_from_api() | ||
| else: | ||
| import pandas as pd | ||
|
|
||
| self._df_iter = PandasDataFrameIterator(pd.DataFrame(), _no_trunc_date) | ||
| self._df = pd.DataFrame() | ||
| if self._df is not None: | ||
| # A shallow copy keeps assignments to the DataFrame from as_pandas() | ||
| # out of the rows that the fetch methods return. Mutable values in its | ||
| # cells, such as lists from JSON columns, are still shared. | ||
| self._df_iter = PandasDataFrameIterator(self._df.copy(deep=False), _no_trunc_date) | ||
|
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 Covered (reviewer): S3 CSV, UNLOAD Parquet, API fallback, empty and non-SUCCEEDED construction; whole-result access, explicit/automatic chunking, fetch positions, close/reset, CSV streams and readers; mutation isolation, memory retention, pandas time truncation, docs; sync, future-based async, and aio wrappers; both new tests. Findings and author verification:
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. Repair: c857dd3 (on top of a26c763; same base a18ebda)
Validation: Self-review of the repair (round one, behavior): no findings. The fallback DataFrame types are unchanged, and only the fetch conversion and the truncation differ. Round two (claims): the commit message, comments, and PR body were rechecked against the measured errors (
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 (relayed): no regressions; 2 pre-existing findings Reviewer: Codex CLI 0.160.0, The reviewer resolved all four earlier comments: assignment isolation is documented with the shared mutable cells; fallback TIME bypasses Pre-existing findings and author disposition:
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 for Reviewer: Codex CLI 0.160.0, |
||
| self._iterrows = self._df_iter.iterrows() | ||
|
|
||
| def _get_parquet_engine(self) -> str: | ||
|
|
@@ -842,22 +852,30 @@ def as_pandas(self) -> PandasDataFrameIterator | DataFrame: | |
| """Return the query results as a DataFrame or an iterator of DataFrame chunks. | ||
|
|
||
| Returns: | ||
| If ``chunksize`` is None, one DataFrame that joins the chunks the result | ||
| iterator has not yet yielded (read in chunks when ``auto_optimize_chunksize`` | ||
| chose a chunk size), which is the whole result unless rows were already | ||
| fetched; otherwise the ``PandasDataFrameIterator`` that yields DataFrame chunks. | ||
| If ``chunksize`` is None, the DataFrame of the whole result, the same one | ||
| on every call. When ``auto_optimize_chunksize`` chose a chunk size, one | ||
| DataFrame that joins the chunks the result iterator has not yet yielded, | ||
| which is the whole result only if neither the fetch methods nor | ||
| ``iter_chunks()`` read from it before, and a later call returns an empty | ||
| DataFrame. If ``chunksize`` is set, the iterator that ``iter_chunks()`` | ||
| returns. | ||
| """ | ||
| if self._chunksize is None: | ||
| if self._df is not None: | ||
| return self._df | ||
| return self._df_iter.as_pandas() | ||
| return self._df_iter | ||
| return self.iter_chunks() | ||
|
|
||
| def iter_chunks(self) -> PandasDataFrameIterator: | ||
| """Iterate over result chunks as pandas DataFrames. | ||
|
|
||
| This method provides an iterator interface for processing large result sets. | ||
| When chunksize is specified, or ``auto_optimize_chunksize`` chose a chunk size | ||
| for a large CSV result, it yields DataFrames in chunks for memory-efficient | ||
| processing. Otherwise, it yields the entire result as a single DataFrame. | ||
| When a CSV result is read in chunks, because chunksize is specified or | ||
|
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, callers, operations): FINDINGS → repaired in a26c763 Base Claims checked:
Findings (repaired):
Callers and operations: no change in S3 or Athena requests (same reads, same places). No user guide in |
||
| ``auto_optimize_chunksize`` chose a chunk size, it yields DataFrames in chunks | ||
| for memory-efficient processing. These chunks come from the same iterator as | ||
| the fetch methods, so a chunk that one of them reads is not available to the | ||
| other. Otherwise, each call returns a new iterator that yields the entire | ||
| result as a single DataFrame, and the fetch methods keep their position. | ||
|
|
||
| Returns: | ||
| PandasDataFrameIterator that yields pandas DataFrames for each chunk | ||
|
|
@@ -876,6 +894,8 @@ def iter_chunks(self) -> PandasDataFrameIterator: | |
| >>> for df in cursor.iter_chunks(): | ||
| ... process(df) # Single DataFrame with all data | ||
| """ | ||
| if self._df is not None: | ||
| return PandasDataFrameIterator(self._df, _no_trunc_date) | ||
| return self._df_iter | ||
|
|
||
| @override | ||
|
|
@@ -884,6 +904,7 @@ def close(self) -> None: | |
|
|
||
| super().close() | ||
| self._df_iter.close() | ||
| self._df_iter = PandasDataFrameIterator(pd.DataFrame(), _no_trunc_date) | ||
| self._df = pd.DataFrame() | ||
| self._df_iter = PandasDataFrameIterator(self._df, _no_trunc_date) | ||
| self._iterrows = enumerate([]) | ||
| self._data_manifest = [] | ||
Uh oh!
There was an error while loading. Please reload this page.
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 (implementation behavior): CLEAN
Base
a18ebdaa80c0c85466eacbdb3eb3df76d7ed135a, headf41ca3c05cf61338ab2ab408541efdb6cf13b9fc.Covered:
AthenaPandasResultSet/AthenaPolarsResultSetconstruction for the S3 output, GetQueryResults fallback, and failed-state paths;as_pandas(),as_polars(),as_arrow(),iter_chunks(),close(); the sync/async/aio cursor wrappers (no subclasses;PandasCursor.iter_chunks()closes the iterator it receives, which is now a fresh one in non-chunked mode); explicitchunksizeandauto_optimize_chunksizechunked paths (unchanged single-pass); the two new tests.Checks:
_trunc_datein__init__could now raise on results that were never read. Refuted live: zero-rowtime, all-NULLtime, and mixedtime/NULL results convert identically on master and this head (master only differs in returning[]fromfetchall()afteras_pandas(), which is the bug).df["a"] = 0, in-placerename, anddf[0, "a"] = 0leaked into fetched rows; with them, they do not.Limitation (not a defect): in non-chunked mode the result set now holds the DataFrame until
close()or the nextexecute(), as before v2.13.0 for pandas, before v3.24.0 for Polars, and likeArrowCursor. Previously a fetch-only caller released it once the single chunk was exhausted.