Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 26 additions & 5 deletions pyathena/arrow/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -302,10 +302,24 @@ def _read_csv(self) -> Table:
):
return pa.Table.from_pydict({})
length = self._get_content_length()
binary_columns = {d[0] for d in self.description or [] if d[1] == "varbinary"}
description = self.description if self.description else []
names = [d[0] for d in description]
# pyarrow types every column with a name by its column_types entry, so columns
# with the same name are read under their positions and get their names back
# after reading.
has_duplicate_names = len(set(names)) != len(names)
if has_duplicate_names:
column_names = [str(i) for i in range(len(names))]
column_types = {
str(i): dtype
for i, d in enumerate(description)
if (dtype := self._converter.get_dtype(d[1], d[4], d[5])) is not None
}
else:
column_names = names
column_types = self.column_types

Copy link
Copy Markdown
Member Author

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 and implementation) — base 8c9c201ac53a06c2f54d7f282ec4de23b0fe449f, head 6af16a0dfc8c6810e634d8c4b1390435c101c354. Result: FINDINGS (1, repaired in 936a692).

Finding: for results with unique column names, _read_csv() built its own positional column_types dict instead of reading the public column_types property, so a subclass overriding that property would no longer affect reading, a change outside #1051.
Repair: unique names use self.column_types and the master read options exactly; only repeated names take the positional path. Rerun: pytest -n 4 tests/pyathena/arrow tests/pyathena/aio/arrow → 119 passed.

Covered without findings:

  • Positional column_names/column_types and rename_columns(names); zip(strict=True) lengths always match by construction.
  • skip_rows_after_names=1 skips a parsed row (verified locally on pyarrow 25.0.1 with a quoted newline in the header and small block sizes). A header-only file gives an empty table with the renamed columns.
  • .txt results with repeated names, varbinary NULL handling by position, and async (AsyncArrowCursor shares the result set).

binary_columns = {i for i, d in enumerate(description) if d[1] == "varbinary"}
if length and self.output_location.endswith(".txt"):
description = self.description if self.description else []
column_names = [d[0] for d in description]
read_opts = csv.ReadOptions(
skip_rows=0,
column_names=column_names,
Expand All @@ -320,6 +334,11 @@ def _read_csv(self) -> Table:
)
elif length and self.output_location.endswith(".csv"):
read_opts = csv.ReadOptions(skip_rows=0, block_size=self._block_size, use_threads=True)
if has_duplicate_names:
read_opts.column_names = column_names
# Skips the header as a parsed row; skip_rows would split a quoted name
# that contains a newline.
read_opts.skip_rows_after_names = 1
parse_opts = csv.ParseOptions(
delimiter=",",
quote_char='"',
Expand All @@ -344,12 +363,14 @@ def _read_csv(self) -> Table:
strings_can_be_null=bool(binary_columns),
quoted_strings_can_be_null=False,
timestamp_parsers=self.timestamp_parsers,
column_types=self.column_types,
column_types=column_types,
),
)
if has_duplicate_names:
table = table.rename_columns(names)
if binary_columns:
for index, field in enumerate(table.schema):
if field.name not in binary_columns and (
if index not in binary_columns and (
pa.types.is_string(field.type) or pa.types.is_binary(field.type)
):
# Preserve the existing CSV behavior for non-binary Athena columns.
Expand Down
64 changes: 64 additions & 0 deletions pyathena/pandas/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -429,6 +429,32 @@ class AthenaPandasResultSet(AthenaResultSet):
]
# 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"})
# The pandas.read_csv() options given to execute() that do not change how pandas reads
# the header row, with which _read_csv_header_as_labels() replaces the header.
_LABELED_HEADER_READ_CSV_OPTIONS: ClassVar[frozenset[str]] = frozenset(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, -s read-only.

Round 1, session 01a10539-de2c-7ee3-8038-4e4c63e5e6b8, reviewed 936a6924..e1567d15 on a clean snapshot at e1567d1. Result: FINDINGS (3):

  1. With sep=None and a populated result, the header was kept, so the dtype copy still failed. This is the same behavior as master.
  2. comment changed the header labels pandas reads (e.g. index_col="z" for "z#suffix"), but the replaced names did not.
  3. encoding changed the decoded header labels, but the replaced names did not.

Findings 2 and 3 were regressions from master for those options; I verified both offline. Repaired in af023a2: the header is replaced only when every option given to execute() is in the allow-list _LABELED_HEADER_READ_CSV_OPTIONS, which lists options that do not change how pandas reads the header row. Any other option (dtype, names, sep, comment, encoding, escapechar, ...) keeps master's header reading, including master's limitation for same-named columns with different types. This narrows the change rather than extending it.

Repair self-review (both perspectives):

  • Offline: comment + index_col, usecols + index_col, user converters, and no options all read as expected.
  • On Athena: pytest -n 4 tests/pyathena/pandas/test_cursor.py tests/pyathena/aio/pandas -k "duplicate_column or binary or header_options or read_options" gave 59 passed.
  • The PR description now states the allow-list condition and the master fallback.

Round 2, session 01a10547-5690-77e0-8675-ca22c9c3b059, reviewed e1567d15..af023a29 on a clean snapshot at af023a2. Covered: all 20 allow-listed options against names + header=None + skiprows for the C and python engines, empty and multiline-header results, guard order, and the fallback for excluded options. Result: CLEAN. This was a static review against pandas 3.0.6 Python sources; the compiled C tokenizer source was unavailable. No tests were run.

{
"cache_dates",
"converters",
"date_format",
"dayfirst",
"decimal",
"dtype_backend",
"false_values",
"float_precision",
"index_col",
"keep_default_na",
"low_memory",
"na_filter",
"na_values",
"nrows",
"on_bad_lines",
"parse_dates",
"skipfooter",
"thousands",
"true_values",
"usecols",
}
)

def __init__(
self,
Expand Down Expand Up @@ -808,6 +834,9 @@ def _read_csv(self) -> TextFileReader | DataFrame:
with ExitStack() as stack:
source: str | IOBase = self.output_location
binary_columns = self._configure_binary_csv_read(read_csv_kwargs, labels)
if labels is not None:
# After _configure_binary_csv_read(), which checks for the header row.
self._read_csv_header_as_labels(read_csv_kwargs, csv_engine)
self._csv_converters = read_csv_kwargs.get("converters") or {}
if binary_columns:
# Given storage_options, even None, open the file through fsspec
Expand Down Expand Up @@ -1058,6 +1087,41 @@ def _key_csv_columns_by_labels(
]
self._time_columns = [label for label, d in columns if d[1] == "time"]

def _read_csv_header_as_labels(self, read_csv_kwargs: dict[str, Any], csv_engine: str) -> None:

Copy link
Copy Markdown
Member Author

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 and implementation) — base 8c9c201ac53a06c2f54d7f282ec4de23b0fe449f, head 6af16a0d…. pandas side: CLEAN.

Covered:

  • Inheritance is real in both engines (pandas 3.0.6): the C engine copies the dtype even with header=0 and names; the python engine copies it to non-date columns, which is why the newline test adds an interval column. No dtype value stands in for "no dtype" (object changes inference and date dtypes; None differs between engines).
  • Skipping the header: the C engine skips parsed rows, and the python engine skips physical lines. pandas opens text with newline="", so universal-newline counting matches for \n, \r, \r\n, and consecutive newlines (verified locally, also with chunksize).
  • Guards: runs only when labels is not None (standard parsing, not the pyarrow engine, labels matching the columns) and neither dtype nor names was given to execute(). User usecols, index_col, and chunksize work with the replaced header (checked on Athena). engine is popped from kwargs by the cursor, so csv_engine is the engine used.
  • Ordering: called after _configure_binary_csv_read(), whose _is_standard_csv_parsing() check requires header=0. BinaryCSVReader passes the header record through unchanged, and test_binary_dataframe_column_names[duplicate_names-c|python] passes.
  • Tests: the new and extended tests fail on master (checked by reverting the source diff).

Copy link
Copy Markdown
Member Author

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, callers, operations, evidence): base 8c9c201ac53a06c2f54d7f282ec4de23b0fe449f, head 936a6924722f375c75891302b6d87640ae6bf5be. Result: FINDINGS (PR description only, corrected). No code changes.

Claims checked:

  • "pandas copies the dtype ... in both engines, even when names is given with header=0": measured locally on pandas 3.0.6 (C: time after Int64; python: interval after Int64 failed with header=0 + names). The comment and docstring here state the same thing.
  • "skip_rows counts physical lines" (Arrow) and "the python engine skips lines": measured locally with a quoted \n header (pyarrow 25.0.1, pandas 3.0.6).
  • Athena allows newlines in column names and writes them in the header: measured on Athena ("a\nb","a\nb").
  • "Results with unique names are read as before": true after the round-one repair (Arrow uses self.column_types and the master options; pandas returns early).
  • "PolarsCursor ... already correct": was based only on the issue's table. Now measured on Athena with the issue's query and the newline-named int/time/interval columns, with and without chunksize=1.
  • Varbinary positional claim: test_binary_null_vs_empty fails on master and passes here.

Callers and operations: there are no API signature or default changes. Public column_types/dtypes/converters keep their keys. With PyAthena-built dtypes, renamed columns lose only an inherited dtype that was never theirs, and pandas no longer emits a "Both a converter and dtype" ParserWarning for those columns. There are no new S3 or Athena requests; pandas adds one in-memory header parse only when names repeat.

Corrections: the TEST section claimed 6af16a0 for a run that preceded a comment-only edit, and did not cover the repair head. It now lists the run per revision (473 on the pre-comment tree; 119 Arrow sync and async tests on 936a692). "Measured with this branch" became a dated Athena measurement, and the Polars evidence was added.
Deferred: .txt results and custom quoting/dialects keep the existing keys, the same scope as #1050.

"""Read the header of the CSV result file as the labels of columns with the same name.

pandas gives a column that it renames, such as ``x.1``, the dtype of the first
column with the name when it has none of its own. When PyAthena builds the
dtypes, the columns are read under their labels, so that each keeps its own type.

Args:
read_csv_kwargs: The options for ``pandas.read_csv()``, updated in place.
csv_engine: The CSV engine that reads the file.
"""
import pandas as pd

names = [d[0] for d in self.description or []]
if len(set(names)) == len(names):
return
if not self._kwargs.keys() <= self._LABELED_HEADER_READ_CSV_OPTIONS:
# Other options, such as dtype, names, sep, comment, or encoding, keep the
# header row as pandas reads it.
return
# pandas renames the names in a header row, and copies their dtypes, even with
# names given, so the header row is skipped instead.
read_csv_kwargs["names"] = self._resolve_csv_column_names(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed): Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, -s read-only, session 01a10527-6805-7ce3-a2dd-fe2a929ab285. Base 8c9c201ac53a06c2f54d7f282ec4de23b0fe449f, head 936a6924722f375c75891302b6d87640ae6bf5be, run in a clean detached snapshot that was still unchanged afterwards. The prompt contained the literal diff and the intended behavior, and no PR text or prior findings. This was a static review: no tests, builds, or network.
Coverage: pandas header replacement for the C and python engines, multiline names, chunking, usecols/index_col/other overrides, BinaryCSVReader, storage_options, empty and header-only files; the Arrow positional path, .txt, binary NULL handling, and the unique-name path; tests and conventions. Result: FINDINGS (2). It found no Arrow or unique-name-path defect.

Finding 1 (P2), this line in 936a692: the replacement names came from _get_column_names() with default options, which ignore caller options such as escapechar. With engine="python", escapechar="\\" and SELECT CAST('01:02:03' AS TIME) AS "a\b", 1 AS "a\b", the dtype and date keys use the labels ab/ab.1, but the names were a\b/a\b.1, giving Missing column provided to 'parse_dates': 'ab'.
Verified: reproduced offline with pandas 3.0.6. Repaired in e1567d1: the names are resolved with _resolve_csv_column_names(names, read_csv_kwargs, pd.read_csv), the same resolution _get_csv_column_labels() uses for the keys.

names, read_csv_kwargs, pd.read_csv
)[0]
read_csv_kwargs["header"] = None
if csv_engine == "python":
# The python engine skips lines, which a name with a newline spans.
header = StringIO(newline="")
csv.writer(header, quoting=csv.QUOTE_ALL, lineterminator="").writerow(names)
read_csv_kwargs["skiprows"] = len(StringIO(header.getvalue(), newline="").readlines())
else:
# The C engine skips parsed rows.
read_csv_kwargs["skiprows"] = 1

def _configure_binary_csv_read(
self, read_csv_kwargs: dict[str, Any], labels: list[Any] | None
) -> set[int]:
Expand Down
19 changes: 14 additions & 5 deletions tests/pyathena/arrow/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,14 @@

class TestArrowCursor:
def test_binary_null_vs_empty(self, arrow_cursor):
query = """SELECT * FROM (VALUES
# The text column has the name of the binary column, and keeps its own NULL.
query = """SELECT id, value, label, text_value AS value FROM (VALUES
(1, CAST(NULL AS VARBINARY), 'null', CAST(NULL AS VARCHAR)),
(2, X'', 'empty', ''),
(3, X'00ff275c25', 'comma, quote" and' || chr(10) || 'newline', 'NULL')
) AS t(id, value, label, text_value) ORDER BY id"""
arrow_cursor.execute(query)
assert arrow_cursor.as_arrow().column("value").to_pylist() == [None, "", "00 ff 27 5c 25"]
assert arrow_cursor.as_arrow().column(1).to_pylist() == [None, "", "00 ff 27 5c 25"]
rows = arrow_cursor.fetchall()
assert [row[:3] for row in rows] == [
(1, None, "null"),
Expand Down Expand Up @@ -1071,21 +1072,29 @@ def test_fetch_all_rows(self, arrow_cursor):
)
def test_duplicate_column_names(self, arrow_cursor):
arrow_cursor.execute(
"SELECT 1 AS x, 2 AS x, 'a' AS y, json_parse('[1]') AS j, 'b' AS j, "
"CAST('12:34:56' AS TIME) AS t, CAST('01:02:03' AS TIME) AS t"
"SELECT 1 AS x, 2 AS x, CAST('01:02:03' AS TIME) AS x, 'a' AS y, "
"json_parse('[1]') AS j, 'b' AS j, CAST('12:34:56' AS TIME) AS t, "
"CAST('01:02:03' AS TIME) AS t"
)
assert arrow_cursor.fetchall() == [
(
1,
2,
datetime(2017, 1, 1, 1, 2, 3).time(),
"a",
[1],
"b",
datetime(2017, 1, 1, 12, 34, 56).time(),
datetime(2017, 1, 1, 1, 2, 3).time(),
)
]
assert arrow_cursor.as_arrow().column_names == ["x", "x", "y", "j", "j", "t", "t"]
assert arrow_cursor.as_arrow().column_names == ["x", "x", "x", "y", "j", "j", "t", "t"]

def test_duplicate_column_names_with_newline(self, arrow_cursor):
"""Columns with the same name keep their own types, with a newline in the name."""
arrow_cursor.execute('SELECT 1 AS "a\nb", CAST(\'01:02:03\' AS TIME) AS "a\nb"')
assert arrow_cursor.fetchall() == [(1, datetime(2017, 1, 1, 1, 2, 3).time())]
assert arrow_cursor.as_arrow().column_names == ["a\nb", "a\nb"]

@pytest.mark.parametrize(
"execute_kwargs", [{}, {"connect_timeout": 3.0, "request_timeout": 4.0}]
Expand Down
45 changes: 43 additions & 2 deletions tests/pyathena/pandas/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1743,13 +1743,14 @@ def test_fetch_all_rows(self, pandas_cursor):
)
def test_duplicate_column_names(self, pandas_cursor):
pandas_cursor.execute(
"SELECT 1 AS x, 'a' AS x, 'b' AS y, json_parse('[1]') AS j, json_parse('[2]') AS j, "
"CAST('12:34:56' AS TIME) AS t, 2 AS t"
"SELECT 1 AS x, 'a' AS x, CAST('01:02:03' AS TIME) AS x, 'b' AS y, "
"json_parse('[1]') AS j, json_parse('[2]') AS j, CAST('12:34:56' AS TIME) AS t, 2 AS t"
)
assert pandas_cursor.fetchall() == [
(
1,
"a",
datetime(2017, 1, 1, 1, 2, 3).time(),
"b",
[1],
[2],
Expand All @@ -1760,13 +1761,53 @@ def test_duplicate_column_names(self, pandas_cursor):
assert pandas_cursor.as_pandas().columns.tolist() == [
"x",
"x.1",
"x.2",
"y",
"j",
"j.1",
"t",
"t.1",
]

@pytest.mark.parametrize("engine", ["c", "python"])
def test_duplicate_column_names_with_newline(self, pandas_cursor, engine):
"""Columns with the same name keep their own types, with a newline in the name."""
pandas_cursor.execute(
'SELECT 1 AS "a\nb", CAST(\'01:02:03\' AS TIME) AS "a\nb", '
"INTERVAL '2' DAY AS \"a\nb\"",
engine=engine,
)
assert pandas_cursor.fetchall() == [
(1, datetime(2017, 1, 1, 1, 2, 3).time(), "2 00:00:00.000")
]
assert pandas_cursor.as_pandas().columns.tolist() == ["a\nb", "a\nb.1", "a\nb.2"]

@pytest.mark.parametrize(
("query", "execute_kwargs", "expected_rows", "expected_columns"),
[
(
'SELECT CAST(\'01:02:03\' AS TIME) AS "a\\b", 1 AS "a\\b"',
{"engine": "python", "escapechar": "\\"},
[(datetime(2017, 1, 1, 1, 2, 3).time(), 1)],
["ab", "ab.1"],
),
(
"SELECT 1 AS x, 2 AS x WHERE false",
{"engine": "python", "sep": None},
[],
["x", "x.1"],
),
],
ids=["escapechar", "detected_delimiter"],
)
def test_duplicate_column_names_header_options(
self, pandas_cursor, query, execute_kwargs, expected_rows, expected_columns
):
"""Options that change how pandas reads the header still apply to the column labels."""
pandas_cursor.execute(query, **execute_kwargs)
assert pandas_cursor.fetchall() == expected_rows
assert pandas_cursor.as_pandas().columns.tolist() == expected_columns

@pytest.mark.parametrize(
("execute_kwargs", "expected_row", "expected_columns"),
[
Expand Down
Loading