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
10 changes: 4 additions & 6 deletions pyathena/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
TypeNode,
TypeSignatureParser,
_split_array_items,
_split_native_array_items,
)
from pyathena.util import override, strtobool

Expand Down Expand Up @@ -234,7 +235,7 @@ def _to_array(varchar_value: str | None) -> list[Any] | str | None:
# If JSON parsing fails, fall back to basic parsing for simple cases
pass

inner = varchar_value[1:-1].strip()
inner = varchar_value[1:-1]
if not inner:
return []

Expand Down Expand Up @@ -378,13 +379,10 @@ def _parse_array_native(inner: str) -> list[Any] | None:
"""
result = []

# Smart split by comma - respect brace groupings
items = _split_array_items(inner)
# Split as Athena joins the items, respecting brace groupings
items = _split_native_array_items(inner)

for item in items:
if not item:
continue

# Handle struct (ROW) values in format {a, b, c} or {key=value, ...}
if item.strip().startswith("{") and item.strip().endswith("}"):
# This is a struct value - parse it as a struct
Expand Down
11 changes: 7 additions & 4 deletions pyathena/pandas/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -157,14 +157,14 @@ def get_chunk(self, size: int | None = None) -> DataFrame:
size: Number of rows to retrieve. If None, returns entire chunk.

Returns:
DataFrame chunk.
DataFrame chunk, with date truncation applied as in iteration.
"""
from pandas.io.parsers import TextFileReader

try:
if isinstance(self._reader, TextFileReader):
return self._reader.get_chunk(size)
return next(self._reader)
return self._trunc_date(self._reader.get_chunk(size))

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): no regressions; 1 pre-existing finding → folded in

Reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort high, sandbox read-only, session 01a102c7-ffb5-7902-8830-9fe54dca52a8. Static review of ec5323ea..a3566039. Snapshot: detached worktree at a3566039fd28bc6d5f5e0b90e550f4768e3726dd. The prompt contained the literal diff and the intended behavior only. Afterwards, the snapshot and the PR worktree were clean at a3566039.

Covered (reviewer): typed/untyped arrays; scalar, row/map, and nested-array elements; NULLs, whitespace, commas, quotes, = and brackets; fallback contracts and callers; pandas whole and chunked reads, all-NULL columns, fetch methods, as_pandas(), engine selection, and test compatibility. No regressions. By source comparison, the first five new native-array cases and the long-prefix JSON case fail on the base.

Finding, pre-existing P2, pyathena/pandas/result_set.py:166: PandasDataFrameIterator.get_chunk() (documented in docs/pandas.md) returned the reader's chunk without _trunc_date. Verified live on this branch with chunksize=2: get_chunk() gave [Timestamp('2026-10-04 12:34:56'), NaT] for a TIME column, while iteration gave [time(12, 34, 56), None].

Repaired in bab2687: get_chunk() applies self._trunc_date to both reader kinds, as __next__ does. The non-chunked iterators use _no_trunc_date, so the whole-result path is unchanged. New TestPandasCursor::test_get_chunk_time passes, and fails with the previous get_chunk(). tests/pyathena/pandas and tests/pyathena/aio/pandas: 309 passed; just lint passed. Self-review of the repair (both rounds): no findings. The release note was extended to cover get_chunk().

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 (relayed): CLEAN for a3566039..bab2687a6a644dc4240d758aa88860905d6a6085.

Reviewer: Codex CLI 0.160.0, gpt-6-astra, effort high, read-only, session 01a102d3-db57-7892-8c52-dd2169ae7b3f. Static review; afterwards, the snapshot was clean at bab2687a. Covered: TextFileReader chunks and single-DataFrame iterators with either callback (one conversion, as in iteration); the explicit-chunksize, auto-optimized, and whole-result paths (_no_trunc_date prevents double truncation); UNLOAD results; and close, exhaustion, and size handling (unchanged; conversion errors get the same cleanup as iteration). The test fails with the previous implementation.

return self._trunc_date(next(self._reader))
except BaseException:
self.close()
raise
Expand Down Expand Up @@ -518,7 +518,10 @@ def parse_dates(self) -> list[Any | None]:

def _trunc_date(self, df: DataFrame) -> DataFrame:
if self._time_columns:
truncated = df.loc[:, self._time_columns].apply(lambda r: r.dt.time)
# A NULL is None, as with the GetQueryResults fallback and the other types.
truncated = df.loc[:, self._time_columns].apply(
lambda r: r.dt.time.astype(object).where(r.notna(), None)
)
for time_col in self._time_columns:
df.isetitem(df.columns.get_loc(time_col), truncated[time_col])
return df
Expand Down
50 changes: 43 additions & 7 deletions pyathena/parser.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,40 @@ def _split_array_items(inner: str) -> list[str]:
return items


def _split_native_array_items(inner: str) -> list[str]:

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 (implementation behavior): CLEAN

Base ec5323ea30e3fc2da1aca536d9cbdf8b51c7b3e2, head a3566039fd28bc6d5f5e0b90e550f4768e3726dd.

Covered: _split_native_array_items, typed _convert_typed_array (JSON-first condition and native branch), untyped _to_array/_parse_array_native, the unchanged map/struct splitting (_split_array_items), and pandas _trunc_date (non-chunked, chunked, and all-NULL columns).

  • Separator: Athena joins array items with ", " (measured). A top-level , followed by a space separates items; brace/bracket depth keeps row/map items ({a=1, b=}) intact. Not stripping the interior keeps [x, ] and [ , x].
  • Safety fallbacks stay as they were: nested arrays in native text return None (typed) or the original string (untyped), and items with []=" make the untyped parser return the original string.
  • Typed array(json): JSON-first only for the json element type. Other types keep the prefix check (json.loads('[1.50]') gives 1.5, which would change array(varchar)). Invalid JSON still falls back to the native branch.
  • pandas: .dt.time of an all-NULL column stays datetime64 (measured with pandas 3), so astype(object) comes before where(..., None). Both a mixed and an all-NULL column give None. parse_dates is unchanged, so pyarrow engine selection is unchanged.
  • Regression coverage: 6 of the 8 new unit cases, both test_fetch_complex_values cases, and pandas[default] test_fetch_all_rows fail with the source reverted.

"""Split the items of an array in Athena's native format.

Athena joins the items with ``", "``, so only a top-level comma followed by a space
separates items. Other commas, leading and trailing spaces, and empty items belong
to the items. Brace and bracket groupings are respected.

Args:
inner: Interior content of the array without brackets, not stripped.

Returns:
List of item strings.
"""
items: list[str] = []
current: list[str] = []
depth = 0
index = 0
while index < len(inner):
char = inner[index]
if char in "{[":
depth += 1
elif char in "}]":
depth -= 1
elif char == "," and depth == 0 and inner.startswith(" ", index + 1):
items.append("".join(current))
current = []
index += 2
continue
current.append(char)
index += 1
items.append("".join(current))
return items


@dataclass
class TypeNode:
"""Parsed representation of an Athena DDL type signature.
Expand Down Expand Up @@ -321,9 +355,14 @@ def _convert_typed_array(self, value: str, type_node: TypeNode) -> list[Any] | N

element_type = type_node.children[0] if type_node.children else TypeNode("varchar")

# Try JSON first (only if content looks like JSON)
# Try JSON first if the elements are JSON, whose values Athena renders as JSON
# text, or if the content looks like JSON
inner_preview = value[1:10] if len(value) > 10 else value[1:-1]
if '"' in inner_preview or value.startswith(("[{", "[null", "[[")):
if (
element_type.type_name == "json"

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): CLEAN

Base ec5323ea30e3fc2da1aca536d9cbdf8b51c7b3e2, head a3566039fd28bc6d5f5e0b90e550f4768e3726dd. Full pass over the PR body, the commit message, the docstrings, the comments, and #1041.

  • The measured Athena text table (master abc99b0) is from a live run with raw converters on Cursor (GetQueryResults) and S3FSCursor (CSV result file); the two give identical text.
  • "Maps and rows already kept empty strings": same run ({=v, k=} → {'': 'v', 'k': ''}, {a=, b=x} → {'a': '', 'b': 'x'}).
  • The ambiguity list matches the measurements: [] for ARRAY[''], [null, NULL] for ARRAY['null', 'NULL'], and "[ ]" parsing as an empty JSON array in the untyped path (unit test removed for that reason).
  • The release notes call out the kept leading/trailing spaces and the pandas NULL TIME change, including as_pandas().
  • No AWS request changes and no public signature changes; _split_native_array_items is new and private.

or '"' in inner_preview
or value.startswith(("[{", "[null", "[["))
):
try:
parsed = json.loads(value)
if isinstance(parsed, list):
Expand All @@ -337,19 +376,16 @@ def _convert_typed_array(self, value: str, type_node: TypeNode) -> list[Any] | N
pass

# Native format
inner = value[1:-1].strip()
inner = value[1:-1]
if not inner:
return []

if "[" in inner:
return None # Nested arrays not supported in native format

items = _split_array_items(inner)
items = _split_native_array_items(inner)
result: list[Any] = []
for item in items:
item = item.strip()
if not item:
continue
if item.startswith("{") and item.endswith("}"):
if element_type.type_name in ("row", "struct"):
result.append(self._convert_typed_struct(item, element_type))
Expand Down
11 changes: 11 additions & 0 deletions tests/pyathena/pandas/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1589,6 +1589,17 @@ def test_pandas_cursor_iter_chunks_consistency(self, pandas_cursor):
for chunk1, chunk2 in zip(chunks_via_method, chunks_via_direct, strict=False):
pd.testing.assert_frame_equal(chunk1, chunk2)

@pytest.mark.parametrize(
"pandas_cursor", [{"cursor_kwargs": {"chunksize": 2}}], indirect=["pandas_cursor"]
)
def test_get_chunk_time(self, pandas_cursor):
"""get_chunk() converts time columns as iteration does."""
pandas_cursor.execute(
"SELECT * FROM (VALUES (1, CAST('12:34:56' AS TIME)), (2, NULL)) AS t(i, v) ORDER BY i"
)
chunk = pandas_cursor.as_pandas().get_chunk()
assert chunk["v"].tolist() == [datetime(2017, 1, 1, 12, 34, 56).time(), None]

@pytest.mark.parametrize(
"pandas_cursor",
[
Expand Down
30 changes: 30 additions & 0 deletions tests/pyathena/test_converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -726,3 +726,33 @@ def test_to_datetime_with_tz_offsets_and_zone_names(input_value, expected):
if expected is not None:
assert result.utcoffset() == expected.utcoffset()
assert result.tzinfo is not None


@pytest.mark.parametrize(
("input_value", "expected"),
[
("[x, , y]", ["x", "", "y"]),
("[x, ]", ["x", ""]),
("[x, , y]", ["x", " ", "y"]),
("[ , x]", [" ", "x"]),
("[a,b, c]", ["a,b", "c"]),
("[x, null]", ["x", None]),
("[{a=1, b=}, {a=, b=2}]", [{"a": "1", "b": ""}, {"a": "", "b": "2"}]),
],
)
def test_native_array_items(input_value, expected):
"""Native arrays are split as Athena joins them, keeping empty items."""
assert _to_array(input_value) == expected
assert (
DefaultTypeConverter().convert("array", input_value, type_hint="array(varchar)") == expected
)


def test_typed_json_array_starting_with_non_string():
"""A typed array(json) value is parsed as JSON whatever its first element is."""
converter = DefaultTypeConverter()
assert converter.convert("array", '[1234567890, "a,b"]', type_hint="array(json)") == [
1234567890,
"a,b",
]
assert converter.convert("array", "[1, 2.5, true]", type_hint="array(json)") == [1, 2.5, True]
36 changes: 36 additions & 0 deletions tests/pyathena/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1693,6 +1693,42 @@ def test_fetch_all_rows(self, cursor):
cursor.execute(CONVERTED_VALUES_QUERY)
assert cursor.fetchall() == [CONVERTED_VALUES_ROW]

@pytest.mark.parametrize(
"cursor",
[
pytest.param({}, id="default"),
pytest.param(
{"work_group": ENV.managed_work_group, "s3_staging_dir": ""},
id="managed",
marks=pytest.mark.skipif(
not ENV.managed_work_group,
reason="AWS_ATHENA_MANAGED_WORKGROUP not set",
),
),
],
indirect=["cursor"],
)
def test_fetch_complex_values(self, cursor):
"""Native arrays keep empty, space-only, and comma items; array(json) keeps JSON."""
query = """
SELECT
ARRAY['x', '', 'y'] AS col_empty
,ARRAY['x', ' ', 'a,b'] AS col_space_comma
,ARRAY[json_parse('1234567890'), json_parse('"a,b"')] AS col_json
"""
cursor.execute(query)
expected = (["x", "", "y"], ["x", " ", "a,b"], [1234567890, "a,b"])
assert cursor.fetchall() == [expected]
cursor.execute(
query,
result_set_type_hints={
"col_empty": "array(varchar)",
"col_space_comma": "array(varchar)",
"col_json": "array(json)",
},
)
assert cursor.fetchall() == [expected]

@staticmethod
def _metadata_view(metadata):
return (
Expand Down
6 changes: 4 additions & 2 deletions tests/pyathena/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -110,8 +110,8 @@ def unreachable_glue(connection):


# TIME values of several precisions, with and without a time zone, TIMESTAMP WITH TIME
# ZONE values with UTC offsets and a zone name, a NULL JSON value, and the row that
# every cursor should fetch for them.
# ZONE values with UTC offsets and a zone name, NULL JSON and TIME values, and the row
# that every cursor should fetch for them.
CONVERTED_VALUES_QUERY = """
SELECT
1 AS col
Expand All @@ -125,6 +125,7 @@ def unreachable_glue(connection):
,TIMESTAMP '2024-02-29 23:59:58.123 -08:00' AS col_timestamp_tz_negative
,TIMESTAMP '2024-02-29 23:59:58.123 America/New_York' AS col_timestamp_tz_name
,CAST(NULL AS TIMESTAMP WITH TIME ZONE) AS col_timestamp_tz_null
,CAST(NULL AS TIME) AS col_time_null
"""
CONVERTED_VALUES_ROW = (
1,
Expand All @@ -138,6 +139,7 @@ def unreachable_glue(connection):
datetime(2024, 2, 29, 23, 59, 58, 123000, tzinfo=timezone(-timedelta(hours=8))),
datetime(2024, 2, 29, 23, 59, 58, 123000, tzinfo=gettz("America/New_York")),
None,
None,
)


Expand Down
Loading