-
Notifications
You must be signed in to change notification settings - Fork 116
Read CSV results for the pandas PyArrow engine with pyarrow.csv #1057
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
5339aff
e272719
84e2fca
fbceb11
15399b5
2302049
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 |
|---|---|---|
|
|
@@ -43,6 +43,104 @@ def _no_trunc_date(df: DataFrame) -> DataFrame: | |
| return df | ||
|
|
||
|
|
||
| def _read_csv_with_pyarrow(source: str | IOBase, read_csv_kwargs: dict[str, Any]) -> DataFrame: | ||
|
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 (implementation behavior): base Covered: Edge cases compared with Result: FINDINGS
Simplifications already applied before publishing: the reader is a module function (it uses no result-set state), and the non-mapping
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): Codex CLI 0.160.0, model Reviewer coverage: the five-file diff, CSV/TSV conversion, engine selection, dtype/date/NULL handling, sync and async callers, stream ownership, fsspec options, tests, and the pandas/Polars docs, compared with pandas 3.0.6, pyarrow 25.0.1, and Polars 1.44.2. Verdict: FINDINGS (5 × P2 regressions)
Author verification: all five confirmed. Findings 1–4 were run against pandas 3.0.6: findings 1 and 3 reproduce as described, master returns Repair in 84e2fca, recorded in the reply below.
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 84e2fca (old head
Self-review of the repair, round one (behavior): parity harness all equal (types, TSV, Round two (claims): the PR text no longer lists C-engine fallbacks as release-note items. Its only behavior change is the multi-line fix on the default path. The independent follow-up on this repair is pending.
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 1 (relayed): Codex CLI 0.160.0, model The reviewer found the five earlier findings resolved: dates followed by a final cast; duplicate names via short-circuit and positional conversion; non-default NA settings, other options, and Verdict: FINDINGS
Author verification: both reproduced against pandas 3.0.6 ( 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 2 (relayed): Codex CLI 0.160.0, model Verdict: CLEAN. Both points are resolved, and the reviewer found no new defect.
Validation at fbceb11:
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. Merge of master (15399b5): master gained #1044 ( Upstream contracts checked: |
||
| """Read a CSV result with pyarrow as ``pandas.read_csv(engine="pyarrow")`` does. | ||
|
|
||
| pandas' PyArrow engine does not set ``newlines_in_values``, so pyarrow | ||
| raises or returns wrong values when a quoted value containing a newline | ||
| crosses one of its read blocks. This reads the file with that option and | ||
| converts the table as pandas does for the options that | ||
| ``AthenaPandasResultSet._reads_csv_with_pyarrow()`` accepts: NULL-typed | ||
| columns become float64, integer columns without a ``dtype`` entry get | ||
| NumPy integer types, the ``dtype`` mapping is applied before and after the | ||
| ``parse_dates`` columns are parsed with ``pandas.to_datetime()``, and | ||
| without ``future.infer_string``, strings become objects. | ||
|
|
||
| Args: | ||
| source: The result file. | ||
| read_csv_kwargs: The pandas.read_csv() options built by | ||
| ``AthenaPandasResultSet._get_csv_read_options()``. | ||
|
|
||
| Returns: | ||
| The result as a DataFrame. | ||
| """ | ||
| import pandas as pd | ||
| import pyarrow as pa | ||
| from pyarrow import csv as pyarrow_csv | ||
|
|
||
| header = read_csv_kwargs["header"] | ||
| names = read_csv_kwargs["names"] | ||
| null_values = list(read_csv_kwargs["na_values"]) | ||
| table = pyarrow_csv.read_csv( | ||
| source, | ||
| read_options=pyarrow_csv.ReadOptions(autogenerate_column_names=header is None), | ||
| parse_options=pyarrow_csv.ParseOptions( | ||
| delimiter=read_csv_kwargs["sep"], | ||
| ignore_empty_lines=read_csv_kwargs["skip_blank_lines"], | ||
| newlines_in_values=True, | ||
| ), | ||
| convert_options=pyarrow_csv.ConvertOptions( | ||
| null_values=null_values, strings_can_be_null="" in null_values | ||
| ), | ||
| ) | ||
| schema = table.schema | ||
| for index, type_ in enumerate(schema.types): | ||
| if pa.types.is_null(type_): | ||
| schema = schema.set(index, schema.field(index).with_type(pa.float64())) | ||
| integer_dtypes = { | ||
| pa.int8(): pd.Int8Dtype(), | ||
| pa.int16(): pd.Int16Dtype(), | ||
| pa.int32(): pd.Int32Dtype(), | ||
| pa.int64(): pd.Int64Dtype(), | ||
| } | ||
| # Integers become nullable dtypes first so that a dtype entry converts them | ||
| # without going through float64. | ||
| df = table.cast(schema).to_pandas(types_mapper=integer_dtypes.get) | ||
| if header is None: | ||
| # pandas names the columns beyond the given names by their positions. | ||
| df.columns = [str(index) for index in range(len(df.columns) - len(names))] + names | ||
|
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. Additional review requested by the maintainer (relayed): Claude Code 2.1.287, model Reviewer coverage: the helper, the predicate and Verdict: FINDINGS (all low, no blocker)
Author verification and repair 2302049:
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 on the repair (relayed): Codex CLI 0.160.0, model Verdict: CLEAN. All three points are resolved, and the reviewer found no new defect.
AWS CI at 2302049: |
||
| dtype = dict(read_csv_kwargs["dtype"]) | ||
| for column in df.columns: | ||
| # Integer columns without a dtype entry get NumPy integer types. | ||
| if column not in dtype and df[column].dtype in integer_dtypes.values(): | ||
| dtype[column] = df[column].dtype.numpy_dtype | ||
| dtype = { | ||
| column: pd.api.types.pandas_dtype(value) | ||
| for column, value in dtype.items() | ||
| if column in df.columns | ||
| } | ||
| df = df.astype(dtype) | ||
| if not pd.get_option("future.infer_string"): | ||
| # Without the string dtype, pandas returns strings, and the string | ||
| # categories of categorical columns, as objects. | ||
| for index in range(len(df.columns)): | ||
| values = df.iloc[:, index] | ||
| if values.dtype == "str": | ||
| df.isetitem(index, values.astype(object).fillna(None)) | ||
| elif isinstance(values.dtype, pd.CategoricalDtype) and ( | ||
| values.dtype.categories.dtype == "str" | ||
| ): | ||
| categories = values.dtype.categories.astype(object) | ||
| df.isetitem( | ||
| index, | ||
| values.astype(pd.CategoricalDtype(categories, ordered=values.dtype.ordered)), | ||
| ) | ||
| for column in read_csv_kwargs["parse_dates"]: | ||
| if isinstance(column, int) and column not in df.columns: | ||
| column = df.columns[column] | ||
| if df[column].dtype.kind in "Mm": | ||
| continue | ||
| values = df[column].astype("string") | ||
| try: | ||
| df[column] = pd.to_datetime(values, utc=False) | ||
| except (ValueError, TypeError): | ||
| # pandas keeps the column as strings if it cannot parse it. | ||
| df[column] = values.to_numpy(dtype=object, na_value=float("nan")) | ||
| # pandas applies the dtype mapping again after parsing dates. | ||
| df = df.astype(dtype) | ||
| return df | ||
|
|
||
|
|
||
| class _JSONConverter: | ||
| """A json converter for ``pandas.read_csv()`` that keeps NULL from making values floats. | ||
|
|
||
|
|
@@ -329,6 +427,8 @@ class AthenaPandasResultSet(AthenaResultSet): | |
| "time", | ||
| "timestamp", | ||
| ] | ||
| # The pandas.read_csv() options given to execute() that _read_csv_with_pyarrow() reads. | ||
| _PYARROW_READ_CSV_OPTIONS: ClassVar[frozenset[str]] = frozenset({"dtype", "parse_dates"}) | ||
|
|
||
| def __init__( | ||
| self, | ||
|
|
@@ -696,7 +796,10 @@ def _read_csv(self) -> TextFileReader | DataFrame: | |
| source = self._csv_stream = stack.enter_context( | ||
| self._fs.open(self.output_location, mode="rb") | ||
| ) | ||
| result = pd.read_csv(source, **read_csv_kwargs) | ||
| if csv_engine == "pyarrow" and self._reads_csv_with_pyarrow(): | ||
| result = _read_csv_with_pyarrow(source, read_csv_kwargs) | ||
| else: | ||
| result = pd.read_csv(source, **read_csv_kwargs) | ||
| if not isinstance(result, pd.DataFrame): | ||
| # The chunk iterator takes ownership of the stream. | ||
| stack.pop_all() | ||
|
|
@@ -716,6 +819,25 @@ def _read_csv(self) -> TextFileReader | DataFrame: | |
| _logger.exception(f"Failed to read {self.output_location}.") | ||
| raise OperationalError(*e.args) from e | ||
|
|
||
| def _reads_csv_with_pyarrow(self) -> bool: | ||
| """Whether ``_read_csv_with_pyarrow()`` reads the CSV result for the PyArrow engine. | ||
|
|
||
| It reproduces ``pandas.read_csv(engine="pyarrow")`` for PyAthena's default NA | ||
| values and for ``dtype`` as a mapping and ``parse_dates`` as a list given to | ||
| ``execute()``. With other options, pandas reads the file. | ||
|
|
||
| Returns: | ||
| True if ``_read_csv_with_pyarrow()`` reads the result. | ||
| """ | ||
| return ( | ||
| not self._keep_default_na | ||
| and isinstance(self._na_values, (list, tuple)) | ||
| and list(self._na_values) == [""] | ||
| and self._kwargs.keys() <= self._PYARROW_READ_CSV_OPTIONS | ||
| and isinstance(self._kwargs.get("dtype", {}), dict) | ||
| and isinstance(self._kwargs.get("parse_dates", []), list) | ||
| ) | ||
|
|
||
| def _get_csv_read_options(self, csv_engine: str, chunksize: int | None) -> dict[str, Any]: | ||
| """Build pandas options for Athena CSV or tab-separated results.""" | ||
| if self.output_location and self.output_location.endswith(".txt"): | ||
|
|
||
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 two (claims, compatibility, operations): base
e49d576ce1e8cb60d74e1fde56de9c72ff3cd8e1, heade272719123fce03c291f22f015ba00ff04b8eb0d.Claims checked:
ParseOptionsfromdelimiter,quote_char,escape_char, andignore_empty_linesonly: confirmed in pandas 3.0.6arrow_parser_wrapper.py, which also mapson_bad_linestoinvalid_row_handler. The PR text now names it.STR_NA_VALUESlives inpandas._libs.parsersand is not exported bypandas.io.parsers, hence thekeep_default_na=Falsecondition.pyarrow.csvdirectly (newlines_in_values=False), and all were correct withTrue.docs/polars.mdsaid "withuse_pyarrow=Trueandschema_overrides=None" as if those two were sufficient. Polars 1.44.2 also requiresn_rows,n_threads, andnull_valuesto be unset andlow_memory=False, so the text now namesschema_overrides=Noneand "Polars' other conditions" (e272719). Live,PolarsCursorwith those options raisesOperationalError: CSV parser got out of sync with chunkeron the 2,000-row query, while the default reader returns all rows correctly.docs/pandas.md: the engine conditions match_get_csv_engine(), and the conversion sentence matches the function.Existing callers: results without multi-line values give frames equal to before (
assert_frame_equal(check_exact=True)over every eligible type,infer_stringon and off, and 200 random files). Callers that passed other read_csv options orkeep_default_na=Truewithengine="pyarrow"now get the C engine, which is listed as a release-note item. No AWS requests are added (one GET of the result file, as before).Result: CLEAN after the wording corrections. The duplicate-name interaction from round one stays deferred to #1050.