Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
c09ac2f
Backport #800: Find an existing table in to_sql whatever the name's case
laughingman7743 Oct 6, 2026
b18e533
Backport #817: Assert cache hits in the test work group instead of pr…
laughingman7743 Oct 6, 2026
3800144
Backport #831: Stop Spark session readiness polling on failure states
laughingman7743 Oct 6, 2026
d81b11c
Backport #830: Shut down the AsyncSparkCursor executor when session t…
laughingman7743 Oct 6, 2026
550503c
Backport #832: Terminate a newly started Spark session when cursor se…
laughingman7743 Oct 6, 2026
b82fdb9
Backport #899: Reject synchronous iteration of the asyncio cursors an…
laughingman7743 Oct 6, 2026
3a14b71
Backport #900 (decimal256 only): Map Arrow decimal256 to Athena decimal
laughingman7743 Oct 6, 2026
df27eb1
Backport #928: Fix stale S3FileSystem listing and object caches
laughingman7743 Oct 6, 2026
eb74aa8
Backport #929: Keep the existing object when appending within a large…
laughingman7743 Oct 6, 2026
0d05a7f
Backport #947: Merge a short trailing block into the previous multipa…
laughingman7743 Oct 6, 2026
0fd65a7
Backport #955: Keep multipart copy parts within the S3 part size limits
laughingman7743 Oct 6, 2026
d534549
Backport #959: Skip the result cache for qmark queries with parameters
laughingman7743 Oct 6, 2026
63eb967
Backport #950: Return the whole result from as_pandas() when auto chu…
laughingman7743 Oct 6, 2026
0afb784
Backport #995: Reflect the field and value types of top-level STRUCT …
laughingman7743 Oct 6, 2026
3d36fbc
Backport #1023: Do not convert ArrowCursor fallback values twice
laughingman7743 Oct 6, 2026
e05d33b
Backport #1029: Stop a closed PolarsDataFrameIterator for every reade…
laughingman7743 Oct 6, 2026
896c287
Backport #1031: Keep NULL rows of single-column CSV results in ArrowC…
laughingman7743 Oct 6, 2026
a3cba8d
Backport #1033 (#854 only): Keep numeric UTC offsets of TIMESTAMP WIT…
laughingman7743 Oct 6, 2026
af6f8fb
Backport #1045: Read multi-line CSV values that cross a block in Arro…
laughingman7743 Oct 6, 2026
683658f
Backport #1075: Use the C engine for Pandas DDL text results
laughingman7743 Oct 6, 2026
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
1 change: 1 addition & 0 deletions docs/aio.md
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ df = cursor.as_pandas() # In-memory conversion, no await needed

The `as_pandas()`, `as_arrow()`, and `as_polars()` convenience methods operate on
already-loaded data and remain synchronous.
With a chunk size chosen by `auto_optimize_chunksize`, `as_pandas()` reads every remaining 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.

Round 2 deferral (docs reader):


See each cursor's documentation page for detailed usage examples.

Expand Down
6 changes: 6 additions & 0 deletions docs/pandas.md
Original file line number Diff line number Diff line change
Expand Up @@ -475,6 +475,10 @@ for chunk in cursor.iter_chunks():
process_chunk(chunk)
```

Without an explicit `chunksize`, `as_pandas()` returns a single DataFrame even when a chunk size was chosen automatically.
It reads every chunk and joins them, so the whole result is held in memory.
Use `iter_chunks()`, `fetchone()`, or `fetchmany()` to read a large result one chunk at a time.

**Priority of chunksize settings:**

1. **Explicit chunksize** (highest priority): Always respected
Expand Down Expand Up @@ -551,6 +555,8 @@ Common performance options:
- `dtype`: Explicit column data types
- `parse_dates`: Columns to parse as dates

With `engine="pyarrow"`, tab-separated `.txt` results from DDL statements such as `SHOW TABLES`, `SHOW COLUMNS`, and `DESCRIBE` use the C engine to preserve leading zeros, exponent notation, and padding in string values.

### Unload options

PandasCursor also supports the unload option, as does {ref}`arrow-cursor`.
Expand Down
8 changes: 8 additions & 0 deletions docs/sqlalchemy.md
Original file line number Diff line number Diff line change
Expand Up @@ -719,6 +719,10 @@ That includes top-level columns, fields of a STRUCT, STRUCT values inside MAP, a
Integer fields, and integer MAP keys and values, use `INT` in that DDL.
`CAST` and other SQL expressions keep `ROW(...)`, `MAP(...)`, and `ARRAY(...)`, and spell integers as `INTEGER`.

Reflected STRUCT and ROW columns use `AthenaStruct` with their field names and types.
A field type that the dialect does not recognize is reflected as `NullType` with a warning.
Selecting such a column returns the value from the cursor, as described under Data format support below; SQLAlchemy does not convert it to the reflected field types.

#### Querying STRUCT data

PyAthena automatically converts STRUCT data between different formats:
Expand Down Expand Up @@ -868,6 +872,10 @@ CREATE TABLE products (
`CREATE TABLE` renders integer MAP keys and values as `INT`.
`CAST` still spells those integers as `INTEGER`.

Reflected MAP columns use `AthenaMap` with their key and value types.
A key or value type that the dialect does not recognize is reflected as `NullType` with a warning.
Selecting such a column returns the value from the cursor, as described under Data format support below; SQLAlchemy does not convert it to the reflected key and value types.

#### Querying MAP data

PyAthena automatically converts MAP data between different formats:
Expand Down
3 changes: 3 additions & 0 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -221,6 +221,8 @@ See the [Athena documentation](https://docs.aws.amazon.com/athena/latest/ug/reus

You can attempt to re-use the results from a previously executed query to help save time and money in the cases where your underlying data isn't changing.
Set the `cache_size` or `cache_expiration_time` parameter of `cursor.execute()` to a number larger than 0 to enable caching.
`cache_size` is the number of the most recent query executions in the work group to search, including executions by other clients of the same work group.
In a busy work group, a previous execution may no longer be among them.

```python
from pyathena import connect
Expand Down Expand Up @@ -256,6 +258,7 @@ cursor.execute("SELECT * FROM one_row", cache_size=100, cache_expiration_time=36

Results will only be re-used if the query strings match *exactly*,
and the query was a DML statement (the assumption being that you always want to re-run queries like `CREATE TABLE` and `DROP TABLE`).
The cache is not used for a `qmark` query with parameters.

The S3 staging directory is not checked, so it's possible that the location of the results is not in your provided `s3_staging_dir`.

Expand Down
27 changes: 20 additions & 7 deletions pyathena/aio/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
import logging
import sys
from datetime import datetime, timedelta, timezone
from typing import Any, cast
from typing import Any, NoReturn, cast

from pyathena.aio.util import async_retry_api_call
from pyathena.common import BaseCursor, CursorIterator
Expand Down Expand Up @@ -60,12 +60,16 @@ async def _execute( # type: ignore[override]
result_reuse_minutes=options.result_reuse_minutes,
execution_parameters=execution_parameters,
)
query_id = await self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
query_id = None
# Athena does not return the ExecutionParameters of earlier executions,
# so the cache cannot tell which parameters an execution ran with (#941).
if not request.get("ExecutionParameters"):
query_id = await self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
if query_id is None:
try:
response = await async_retry_api_call(
Expand Down Expand Up @@ -376,6 +380,7 @@ class WithAsyncFetch(AioBaseCursor, CursorIterator, WithResultSet):
``rownumber``, ``rowcount``), lifecycle methods (``close``, ``executemany``,
``cancel``), default sync fetch (for cursors whose result sets load all
data eagerly in ``__init__``), and the async iteration protocol.
Synchronous iteration raises ``TypeError``.

Subclasses override ``execute()`` and optionally ``__init__`` and
format-specific helpers.
Expand Down Expand Up @@ -504,6 +509,14 @@ def fetchall(
result_set = cast(AthenaResultSet, self.result_set)
return result_set.fetchall()

def __iter__(self) -> NoReturn:

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.

Round 1 note (#899):

"""Reject synchronous iteration; use ``async for`` instead.

Raises:
TypeError: Always, because the fetch methods are coroutines.
"""
raise TypeError(f"'{type(self).__name__}' object is not iterable; use 'async for' instead.")

def __aiter__(self):
return self

Expand Down
12 changes: 11 additions & 1 deletion pyathena/aio/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
from typing import (
TYPE_CHECKING,
Any,
NoReturn,
cast,
)

Expand All @@ -25,7 +26,8 @@ class AthenaAioResultSet(AthenaResultSet):

Skips the synchronous ``_pre_fetch`` by passing ``_pre_fetch=False`` to
the parent ``__init__`` and provides an ``async create()`` classmethod
factory instead.
factory instead. Synchronous iteration raises ``TypeError``; use
``async for`` instead.
"""

def __init__(
Expand Down Expand Up @@ -191,6 +193,14 @@ async def fetchall( # type: ignore[override]
break
return rows

def __iter__(self) -> NoReturn:
"""Reject synchronous iteration; use ``async for`` instead.

Raises:
TypeError: Always, because the fetch methods are coroutines.
"""
raise TypeError(f"'{type(self).__name__}' object is not iterable; use 'async for' instead.")

def __aiter__(self):
return self

Expand Down
23 changes: 17 additions & 6 deletions pyathena/arrow/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,9 @@ def __init__(
import pyarrow as pa

self._table = pa.Table.from_pydict({})
# The fetch methods convert only the values read from a result file.
# GetQueryResults values are already converted.
self._convert_rows = bool(self.output_location)
self._batches = iter(self._table.to_batches(arraysize))

def __s3_file_system(self):
Expand Down Expand Up @@ -215,11 +218,15 @@ def _fetch(self) -> None:
return
else:
dict_rows = rows.to_pydict()
column_names = dict_rows.keys()
processed_rows = [
tuple(self.converters[k](v) for k, v in zip(column_names, row, strict=False))
for row in zip(*dict_rows.values(), strict=False)
]
if self._convert_rows:
converters = self.converters
column_names = dict_rows.keys()
processed_rows = [
tuple(converters[k](v) for k, v in zip(column_names, row, strict=False))
for row in zip(*dict_rows.values(), strict=False)
]
else:
processed_rows = list(zip(*dict_rows.values(), strict=False))

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.1-sol, codex exec -s read-only, effort high. Session 01a11134-c83c-7aa2-96a0-26d06ed51cb0. Static review only: no builds, tests, edits or network.

  • Snapshot: detached worktree at head 683658f24dff765e0cc6df1bffcfebdd9f7797a8, base f09cf4269052b7dfa4157f10a5f37a5cbfa6148a. The prompt contained the diff range, the mapping from commit to master merge commit, the 3.x dependency floors and the required scope. It contained no PR text, commit messages or prior findings.
  • The snapshot and the PR worktree were unchanged afterwards.
  • Coverage reported:
    • all 20 commits against their master patches, including the three partial ports;
    • sync, threaded and asyncio cursors; pandas, Arrow, Polars and S3FS result sets;
    • Spark lifecycle;
    • S3 caches, append, multipart upload/copy and failure paths;
    • SQLAlchemy reflection and compiler, dependency bounds, tests and docs.
  • Verdict: FINDINGS, with two P2 findings: one introduced (this thread) and one pre-existing (next thread).

Finding (introduced, P2):

  • With managed result storage (no output location), the GetQueryResults fallback no longer applies the result set's converter mapping in the fetch methods.
  • So a custom mapping such as DefaultArrowTypeConverter().set("varchar", str.upper) stops applying: fetchall() returns ("hello",) instead of ("HELLO",). The same holds on master since Do not convert ArrowCursor fallback values twice #1023.

Author verification: confirmed, deferred.

self._rows.extend(processed_rows)

def fetchone(
Expand Down Expand Up @@ -297,9 +304,13 @@ def _read_csv(self) -> Table:
parse_opts = csv.ParseOptions(
delimiter=",",
quote_char='"',
ignore_empty_lines=not binary_columns,
# Athena writes a single-column row with a NULL value as an empty line.
ignore_empty_lines=False,
double_quote=True,
escape_char=False,
# A quoted value can contain a newline, so the reader must not split
# blocks inside quotes.
newlines_in_values=True,
)
else:
return pa.Table.from_pydict({})
Expand Down
2 changes: 1 addition & 1 deletion pyathena/arrow/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ def get_athena_type(type_: DataType) -> tuple[str, int, int]:
return "date", 0, 0
if type_.id == types.Type_TIMESTAMP: # 18
return "timestamp", 3, 0
if type_.id in [types.Type_DECIMAL128, types.Decimal256Type]: # 23, 24
if type_.id in [types.Type_DECIMAL128, types.Type_DECIMAL256]: # 23, 24

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 2: claims, callers and operations: FINDINGS, repaired (wording in commit messages and the PR body only; no code change)

Scope: git diff f09cf4269052b7dfa4157f10a5f37a5cbfa6148a..683658f24dff765e0cc6df1bffcfebdd9f7797a8, the PR body, and the 20 commit messages.

Claims checked against evidence:

type_ = cast(types.Decimal128Type, type_)
return "decimal", type_.precision, type_.scale
if type_.id in [
Expand Down
16 changes: 10 additions & 6 deletions pyathena/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -750,12 +750,16 @@ def _execute(
result_reuse_minutes=options.result_reuse_minutes,
execution_parameters=execution_parameters,
)
query_id = self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
query_id = None
# Athena does not return the ExecutionParameters of earlier executions,
# so the cache cannot tell which parameters an execution ran with (#941).
if not request.get("ExecutionParameters"):
query_id = self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
if query_id is None:
try:
query_id = retry_api_call(
Expand Down
26 changes: 24 additions & 2 deletions pyathena/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from abc import ABCMeta, abstractmethod
from collections.abc import Callable
from copy import deepcopy
from datetime import date, datetime, time
from datetime import date, datetime, time, timedelta, timezone
from decimal import Decimal
from typing import Any, ClassVar

Expand Down Expand Up @@ -40,11 +40,33 @@ def _to_datetime(varchar_value: str | None) -> datetime | None:
return datetime.strptime(varchar_value, "%Y-%m-%d %H:%M:%S.%f")


_UTC_OFFSET_PATTERN: re.Pattern[str] = re.compile(r"([+-])(\d{2}):(\d{2})")


def _parse_utc_offset(value: str) -> timezone | 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 1: implementation behavior: CLEAN (no actionable findings)

Scope: git diff f09cf4269052b7dfa4157f10a5f37a5cbfa6148a..5bf0490e450e56ca9a1be3840e94d86234beae4f, the 20 commits, 36 files.

Covered:

Recorded limitations (3.x consequences of the agreed scope, not defects of this PR) are in the threads below.

"""Parse a ``+HH:MM`` or ``-HH:MM`` UTC offset.

Args:
value: The text to parse.

Returns:
The fixed-offset time zone, or None if the text is not an offset.
"""
match = _UTC_OFFSET_PATTERN.fullmatch(value)
if not match:
return None
sign, hours, minutes = match.groups()
offset = timedelta(hours=int(hours), minutes=int(minutes))
return timezone(-offset if sign == "-" else offset)


def _to_datetime_with_tz(varchar_value: str | None) -> datetime | None:
if varchar_value is None:
return None
datetime_, _, tz = varchar_value.rpartition(" ")
return datetime.strptime(datetime_, "%Y-%m-%d %H:%M:%S.%f").replace(tzinfo=gettz(tz))
return datetime.strptime(datetime_, "%Y-%m-%d %H:%M:%S.%f").replace(
tzinfo=_parse_utc_offset(tz) or gettz(tz)
)


def _to_time(varchar_value: str | None) -> time | None:
Expand Down
Loading
Loading