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
5 changes: 2 additions & 3 deletions pyathena/aio/arrow/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -138,9 +138,8 @@ async def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
3 changes: 2 additions & 1 deletion pyathena/aio/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -815,7 +815,8 @@ async def fetchmany(
blocking the event loop.

Args:
size: Maximum number of rows to fetch. Defaults to arraysize.
size: Maximum number of rows to fetch. If None or not positive,
``arraysize`` is used.

Returns:
List of tuples representing the fetched rows.
Expand Down
19 changes: 10 additions & 9 deletions pyathena/aio/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ class AioCursor(WithAsyncFetch):
calls, keeping the event loop free.

Example:
>>> async with AioConnection.create(...) as conn:
>>> async with await AioConnection.create(...) as conn:
... async with conn.cursor() as cursor:
... await cursor.execute("SELECT * FROM my_table")
... rows = await cursor.fetchall()
Expand Down Expand Up @@ -133,9 +133,8 @@ async def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down Expand Up @@ -188,7 +187,8 @@ async def fetchone(
"""Fetch the next row of a query result set.

Returns:
A tuple representing the next row, or None if no more rows.
The next row (a tuple, or a dict for ``AioDictCursor``), or None if
no more rows.

Raises:
ProgrammingError: If called before executing a query that
Expand All @@ -204,10 +204,11 @@ async def fetchmany(self, size: int | None = None) -> list[Any | dict[Any, Any |
"""Fetch multiple rows from a query result set.

Args:
size: Maximum number of rows to fetch. If None, uses arraysize.
size: Maximum number of rows to fetch. If None or not positive,
``arraysize`` is used.

Returns:
List of tuples representing the fetched rows.
The fetched rows.

Raises:
ProgrammingError: If called before executing a query that
Expand All @@ -225,7 +226,7 @@ async def fetchall(
"""Fetch all remaining rows from a query result set.

Returns:
List of tuples representing all remaining rows in the result set.
The remaining rows.

Raises:
ProgrammingError: If called before executing a query that
Expand All @@ -241,7 +242,7 @@ class AioDictCursor(AioCursor):
"""Native asyncio cursor that returns rows as dictionaries.

Example:
>>> async with AioConnection.create(...) as conn:
>>> async with await AioConnection.create(...) as conn:
... cursor = conn.cursor(AioDictCursor)
... await cursor.execute("SELECT id, name FROM users")
... row = await cursor.fetchone()
Expand Down
5 changes: 2 additions & 3 deletions pyathena/aio/pandas/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,9 +161,8 @@ async def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
5 changes: 2 additions & 3 deletions pyathena/aio/polars/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,9 +145,8 @@ async def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
15 changes: 8 additions & 7 deletions pyathena/aio/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -87,8 +87,8 @@ async def create(
query_execution: Query execution metadata.
arraysize: Number of rows to fetch per request.
retry_config: Retry configuration for API calls.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
**kwargs: Additional arguments passed to the constructor of ``cls``,
such as ``dict_type`` for ``AthenaAioDictResultSet``.

Expand Down Expand Up @@ -179,7 +179,8 @@ async def fetchone( # type: ignore[override]
page is exhausted and more pages are available.

Returns:
A tuple representing the next row, or None if no more rows.
The next row (a tuple, or a dict for ``AthenaAioDictResultSet``), or
None if no more rows.
"""
if not self._rows and self._next_token:
await self._async_fetch()
Expand All @@ -197,11 +198,11 @@ async def fetchmany( # type: ignore[override]
"""Fetch multiple rows from the result set.

Args:
size: Maximum number of rows to fetch. If None, uses arraysize.
size: Maximum number of rows to fetch. If None or not positive,
``arraysize`` is used.

Returns:
List of row tuples. May contain fewer rows than requested if
fewer are available.
The rows, fewer than ``size`` when the result is exhausted.
"""
if not size or size <= 0:
size = self._arraysize
Expand All @@ -221,7 +222,7 @@ async def fetchall( # type: ignore[override]
"""Fetch all remaining rows from the result set.

Returns:
List of all remaining row tuples.
The remaining rows.
"""
rows = []
while True:
Expand Down
5 changes: 2 additions & 3 deletions pyathena/aio/s3fs/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -141,9 +141,8 @@ async def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
5 changes: 2 additions & 3 deletions pyathena/arrow/async_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,9 +206,8 @@ def execute(
result_reuse_enable: Enable Athena result reuse for this query.
result_reuse_minutes: Minutes to reuse cached results.
paramstyle: Parameter style ('qmark' or 'pyformat').
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
9 changes: 4 additions & 5 deletions pyathena/arrow/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,13 +45,13 @@ class ArrowCursor(WithFetch):
>>> from pyathena.arrow.cursor import ArrowCursor
>>> cursor = connection.cursor(ArrowCursor)
>>> cursor.execute("SELECT * FROM large_table")
>>> table = cursor.fetchall() # Returns pyarrow.Table
>>> table = cursor.as_arrow() # Returns pyarrow.Table
>>> df = table.to_pandas() # Convert to pandas if needed

# High-performance UNLOAD for large datasets
>>> cursor = connection.cursor(ArrowCursor, unload=True)
>>> cursor.execute("SELECT * FROM huge_table")
>>> table = cursor.fetchall() # Faster Parquet-based result
>>> table = cursor.as_arrow() # Faster Parquet-based result
"""

def __init__(
Expand Down Expand Up @@ -164,9 +164,8 @@ def execute(
on_start_query_execution: Callback invoked with the query ID before ``execute()``
waits for the query: after the ``StartQueryExecution`` call, or after a
reusable query ID is found through ``cache_size``.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down
2 changes: 1 addition & 1 deletion pyathena/arrow/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ class AthenaArrowResultSet(AthenaResultSet):
>>> cursor.execute("SELECT * FROM large_table")
>>>
>>> # Get Arrow Table
>>> table = cursor.fetchall()
>>> table = cursor.as_arrow()
>>>
>>> # Convert to pandas if needed
>>> df = table.to_pandas()
Expand Down
29 changes: 14 additions & 15 deletions pyathena/async_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,16 @@ class AsyncCursor(BaseCursor):
provides methods to check query status and retrieve results when ready.

Attributes:
description: Sequence of column descriptions for the last query.
rowcount: Number of rows affected by the last query (-1 for SELECT queries).
arraysize: Default number of rows to fetch with fetchmany().
max_workers: Maximum number of worker threads for concurrent execution.
arraysize: Default number of rows that fetchmany() returns on the result

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 9b74767cf65e0ef513338fb2bb2da678cdfbcd2a, head 763004580effd74f444ca6a0045f25eb987f9fa7 (full diff, 31 files, docstrings only).

I traced each changed claim into the code it describes:

  • AsyncCursor: execute() returns (query_id, future) (async_cursor.py:219). arraysize reaches the result sets through _collect_result_set. description(query_id) returns a Future.
  • Fetch wording:
    • fetchmany() falls back to arraysize for non-positive sizes in the base result set (result_set.py:516), in the aio result set (aio/result_set.py:202), and through WithFetch/WithAsyncFetch delegation.
    • Dict rows come from AthenaDictResultSet / AthenaAioDictResultSet.
  • Pandas retry_config: it is used for GetQueryResults and for _get_content_length/_read_data_manifest. CSV and Parquet reads go through S3FileSystem(connection=...), which takes connection.retry_config (filesystem/s3.py:182).
  • Polars converter: dtypes feed schema_overrides in both read_csv and scan_csv (polars/result_set.py:440,625). The conversion mapping covers date/time/varbinary/json only.
  • S3StorageClass BUCKET/DIRECTORY are assigned in filesystem/s3.py:284,319,378,581.
  • ExecuteOptions is a frozen dataclass, so dataclasses.replace can clear a field.
  • S3AioExecutor uses asyncio.to_thread. The sync filesystem creates one executor per opened file (_open → _create_executor).
  • Type compiler: FLOAT/REAL render FLOAT (compiler.py:104-110), CAST uses REAL (:852-856), and string types render CHAR(n)/VARCHAR(n) or STRING (:180-198).
  • Attributes sweep: a script compared each class's Attributes: with its public attributes. Only the TypeNode dataclass fields remain flagged, and they exist.

No code or test changes. Checks run: just lint and sphinx-build (no new warnings).

sets this cursor creates.

Example:
>>> cursor = connection.cursor(AsyncCursor)
>>>
>>> # Execute multiple queries concurrently
>>> future1 = cursor.execute("SELECT COUNT(*) FROM table1")
>>> future2 = cursor.execute("SELECT COUNT(*) FROM table2")
>>> future3 = cursor.execute("SELECT COUNT(*) FROM table3")
>>> query_id1, future1 = cursor.execute("SELECT COUNT(*) FROM table1")
>>> query_id2, future2 = cursor.execute("SELECT COUNT(*) FROM table2")
>>> query_id3, future3 = cursor.execute("SELECT COUNT(*) FROM table3")
>>>
>>> # Check if queries are done and get results
>>> if future1.done():
Expand All @@ -50,8 +48,10 @@ class AsyncCursor(BaseCursor):
>>> results = [f.result().fetchall() for f in [future1, future2, future3]]

Note:
Each execute() call returns a Future object that can be used to
check completion status and retrieve results.
Each execute() call returns a ``(query_id, future)`` tuple. The future
resolves to the result set, which provides the column descriptions and
the fetch methods. ``description(query_id)`` also returns the column
descriptions as a Future.
"""

def __init__(
Expand Down Expand Up @@ -235,9 +235,8 @@ def execute(
result_reuse_enable: Enable result reuse for identical queries (optional).
result_reuse_minutes: Result reuse duration in minutes (optional).
paramstyle: Parameter style to use (optional).
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
keyword arguments take precedence over ``options`` fields.
Expand Down Expand Up @@ -331,9 +330,9 @@ class AsyncDictCursor(AsyncCursor):

Example:
>>> cursor = connection.cursor(AsyncDictCursor)
>>> future = cursor.execute("SELECT id, name, email FROM users")
>>> result_cursor = future.result()
>>> row = result_cursor.fetchone()
>>> query_id, future = cursor.execute("SELECT id, name, email FROM users")
>>> result_set = future.result()
>>> row = result_set.fetchone()
>>> print(f"User: {row['name']} ({row['email']})")
"""

Expand Down
1 change: 0 additions & 1 deletion pyathena/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -501,7 +501,6 @@ class Converter(metaclass=ABCMeta):

Attributes:
mappings: Dictionary mapping Athena type names to conversion functions.
default: Default conversion function for unmapped types.
types: Optional dictionary mapping type names to Python type objects.
"""

Expand Down
8 changes: 4 additions & 4 deletions pyathena/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ class Cursor(WithFetch):

Example:
>>> cursor = connection.cursor()
>>> cursor.execute("SELECT name, age FROM users WHERE age > %s", (18,))
>>> cursor.execute("SELECT name, age FROM users WHERE age > %(age)s", {"age": 18})
>>> while True:
... row = cursor.fetchone()
... if not row:
Expand Down Expand Up @@ -141,9 +141,9 @@ def execute(
Function signature: (query_id: str) -> None
This allows early access to query_id for
monitoring/cancellation.
result_set_type_hints: Optional dictionary mapping column names to
Athena DDL type signatures for precise type conversion within
complex types. For example:
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.
For example:
``{"tags": "array(varchar)", "metadata": "map(varchar, integer)"}``
options: Shared execution options as an
:class:`~pyathena.options.ExecuteOptions` instance. Individual
Expand Down
4 changes: 0 additions & 4 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -66,10 +66,6 @@ class S3FileSystem(AbstractFileSystem):
(e.g., ``404`` -> ``FileNotFoundError``, ``403`` -> ``PermissionError``)

Attributes:
session: The boto3 session used for S3 operations.
client: The S3 client for direct API calls.
config: Boto3 configuration for the client.
retry_config: Configuration for retry behavior on failed operations.
allow_bucket_creation: Whether mkdir/makedirs may create buckets.
Defaults to False.
allow_bucket_deletion: Whether rmdir may delete buckets.
Expand Down
12 changes: 7 additions & 5 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,8 +46,8 @@ class AioS3FileSystem(AsyncFileSystem):
This avoids diamond inheritance issues and keeps all boto3 logic in one place.

File handles created by ``_open`` use ``S3AioExecutor`` so that parallel
operations (range reads, multipart uploads) are dispatched via the event loop
instead of spawning additional threads.
operations (range reads, multipart uploads) are dispatched through the event
loop with ``asyncio.to_thread`` instead of a ``ThreadPoolExecutor`` per file.

Attributes:
_sync_fs: The internal synchronous S3FileSystem instance.
Expand All @@ -59,8 +59,9 @@ class AioS3FileSystem(AsyncFileSystem):
>>> # Use in async context
>>> files = await fs._ls('s3://my-bucket/data/')
>>>
>>> # Sync wrappers also available (auto-generated by fsspec)
>>> files = fs.ls('s3://my-bucket/data/')
>>> # Sync wrappers (auto-generated by fsspec) need an instance created
>>> # without asynchronous=True; they block the caller until done
>>> files = AioS3FileSystem().ls('s3://my-bucket/data/')
"""

# https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html
Expand Down Expand Up @@ -587,5 +588,6 @@ class AioS3File(S3File):
``isinstance`` checks and to document the async execution model.
All parallel operations (range reads, multipart uploads) are dispatched
through the ``S3Executor`` interface — the ``S3AioExecutor``
provided by ``AioS3FileSystem`` uses the event loop instead of threads.
provided by ``AioS3FileSystem`` dispatches them through the event loop with
``asyncio.to_thread`` instead of a ``ThreadPoolExecutor`` per file.
"""
2 changes: 2 additions & 0 deletions pyathena/filesystem/s3_object.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ class S3StorageClass:
- DEEP_ARCHIVE: Lowest cost archive storage
- GLACIER_IR: Archive with faster retrieval than standard Glacier
- OUTPOSTS: Storage on AWS Outposts
- BUCKET: Pseudo storage class PyAthena assigns to bucket entries
- DIRECTORY: Pseudo storage class PyAthena assigns to directory entries

See Also:
AWS S3 storage classes documentation:
Expand Down
Loading
Loading