From 3840d11b241cca19b501ac63bdefaf40ba835a31 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 12:48:29 +0900 Subject: [PATCH 1/3] Correct docstrings that contradict the implementation Fix the docstring claims verified in #927: examples that use a nonexistent attribute, treat AsyncCursor.execute() as returning a future, omit await on AioConnection.create(), call fetchall() for DataFrames or Spark output, iterate a cursor for chunks, or pass a tuple to the pyformat style. Also remove Attributes entries that are not public attributes, describe result_set_type_hints keys and the fetchmany() fallback as implemented, and correct the wrap_unload output, the supported formatter types, the pandas arraysize/retry_config/kwargs descriptions, the Polars converter scope, the S3 storage class list, and the type compiler mapping. Co-Authored-By: Claude Opus 5.5 --- pyathena/aio/arrow/cursor.py | 5 ++--- pyathena/aio/cursor.py | 19 ++++++++++--------- pyathena/aio/pandas/cursor.py | 5 ++--- pyathena/aio/polars/cursor.py | 5 ++--- pyathena/aio/result_set.py | 15 ++++++++------- pyathena/aio/s3fs/cursor.py | 5 ++--- pyathena/arrow/async_cursor.py | 5 ++--- pyathena/arrow/cursor.py | 9 ++++----- pyathena/arrow/result_set.py | 2 +- pyathena/async_cursor.py | 29 ++++++++++++++--------------- pyathena/converter.py | 1 - pyathena/cursor.py | 8 ++++---- pyathena/filesystem/s3.py | 4 ---- pyathena/filesystem/s3_async.py | 9 +++++---- pyathena/filesystem/s3_object.py | 2 ++ pyathena/formatter.py | 13 ++++++++----- pyathena/model.py | 10 +++++----- pyathena/options.py | 11 ++++++----- pyathena/pandas/async_cursor.py | 12 +++++------- pyathena/pandas/cursor.py | 28 +++++++++++++++++----------- pyathena/pandas/result_set.py | 21 ++++++++++++--------- pyathena/polars/async_cursor.py | 5 ++--- pyathena/polars/converter.py | 9 +++++---- pyathena/polars/cursor.py | 5 ++--- pyathena/polars/result_set.py | 4 ++-- pyathena/result_set.py | 10 ++++++---- pyathena/s3fs/async_cursor.py | 10 ++++------ pyathena/s3fs/cursor.py | 5 ++--- pyathena/spark/async_cursor.py | 2 -- pyathena/spark/common.py | 1 - pyathena/spark/cursor.py | 5 ++--- pyathena/sqlalchemy/compiler.py | 6 ++++-- 32 files changed, 140 insertions(+), 140 deletions(-) diff --git a/pyathena/aio/arrow/cursor.py b/pyathena/aio/arrow/cursor.py index f1b4aa71d..b8c915515 100644 --- a/pyathena/aio/arrow/cursor.py +++ b/pyathena/aio/arrow/cursor.py @@ -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. diff --git a/pyathena/aio/cursor.py b/pyathena/aio/cursor.py index 00a57be0b..037addb33 100644 --- a/pyathena/aio/cursor.py +++ b/pyathena/aio/cursor.py @@ -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() @@ -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. @@ -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 @@ -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 @@ -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 @@ -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() diff --git a/pyathena/aio/pandas/cursor.py b/pyathena/aio/pandas/cursor.py index 967ec3e4f..43350c1f2 100644 --- a/pyathena/aio/pandas/cursor.py +++ b/pyathena/aio/pandas/cursor.py @@ -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. diff --git a/pyathena/aio/polars/cursor.py b/pyathena/aio/polars/cursor.py index 8338c9d61..8f584da6f 100644 --- a/pyathena/aio/polars/cursor.py +++ b/pyathena/aio/polars/cursor.py @@ -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. diff --git a/pyathena/aio/result_set.py b/pyathena/aio/result_set.py index 59299f101..27642580a 100644 --- a/pyathena/aio/result_set.py +++ b/pyathena/aio/result_set.py @@ -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``. @@ -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() @@ -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 @@ -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: diff --git a/pyathena/aio/s3fs/cursor.py b/pyathena/aio/s3fs/cursor.py index 91a176170..172348761 100644 --- a/pyathena/aio/s3fs/cursor.py +++ b/pyathena/aio/s3fs/cursor.py @@ -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. diff --git a/pyathena/arrow/async_cursor.py b/pyathena/arrow/async_cursor.py index a9f5d6766..873a19090 100644 --- a/pyathena/arrow/async_cursor.py +++ b/pyathena/arrow/async_cursor.py @@ -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. diff --git a/pyathena/arrow/cursor.py b/pyathena/arrow/cursor.py index ded70ee77..e98890604 100644 --- a/pyathena/arrow/cursor.py +++ b/pyathena/arrow/cursor.py @@ -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__( @@ -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. diff --git a/pyathena/arrow/result_set.py b/pyathena/arrow/result_set.py index de9d4c66d..ea11a419f 100644 --- a/pyathena/arrow/result_set.py +++ b/pyathena/arrow/result_set.py @@ -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() diff --git a/pyathena/async_cursor.py b/pyathena/async_cursor.py index 2963e5f40..076cee37c 100644 --- a/pyathena/async_cursor.py +++ b/pyathena/async_cursor.py @@ -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 + 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(): @@ -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__( @@ -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. @@ -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']})") """ diff --git a/pyathena/converter.py b/pyathena/converter.py index 560549430..3ea1577a2 100644 --- a/pyathena/converter.py +++ b/pyathena/converter.py @@ -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. """ diff --git a/pyathena/cursor.py b/pyathena/cursor.py index 099f8e8eb..d29ff5105 100644 --- a/pyathena/cursor.py +++ b/pyathena/cursor.py @@ -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: @@ -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 diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 940a2649e..e0633e25e 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -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. diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index 3e8279728..8dcd2ab0f 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -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. @@ -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 and a caller outside a running event loop + >>> files = AioS3FileSystem().ls('s3://my-bucket/data/') """ # https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html diff --git a/pyathena/filesystem/s3_object.py b/pyathena/filesystem/s3_object.py index c28e45dd9..dc6ffbdfc 100644 --- a/pyathena/filesystem/s3_object.py +++ b/pyathena/filesystem/s3_object.py @@ -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: diff --git a/pyathena/formatter.py b/pyathena/formatter.py index 632c8e676..d453c19a5 100644 --- a/pyathena/formatter.py +++ b/pyathena/formatter.py @@ -41,7 +41,6 @@ class Formatter(metaclass=ABCMeta): Attributes: mappings: Dictionary mapping Python types to formatting functions. - default: Default formatting function for unmapped types. """ def __init__( @@ -139,8 +138,9 @@ def wrap_unload( for large datasets and preserves data types more accurately. Args: - operation: SQL query to wrap. Must be a SELECT or WITH statement. - s3_staging_dir: Base S3 directory for storing UNLOAD results. + operation: SQL query to wrap. Only SELECT and WITH statements are wrapped. + s3_staging_dir: Base S3 directory for storing UNLOAD results, ending + with ``/``. format_: Output file format. Defaults to Parquet for optimal performance. compression: Compression algorithm. Defaults to Snappy for balanced compression ratio and speed. @@ -159,12 +159,15 @@ def wrap_unload( UNLOAD ( SELECT * FROM sales WHERE year = 2023 ) - TO 's3://my-bucket/results/unload/20231215/uuid//' + TO 's3://my-bucket/results/unload/20231215//' WITH ( format = 'PARQUET', compression = 'SNAPPY' ) + Raises: + ProgrammingError: If the query is None, empty, or only whitespace. + Note: Only SELECT and WITH statements are wrapped. Other statement types are returned unchanged with location=None. @@ -413,7 +416,7 @@ class DefaultParameterFormatter(Formatter): - Strings: Properly escaped and quoted - Binary data: bytes, bytearray, memoryview as hexadecimal literals - Numbers: int, float, Decimal - - Dates and times: date, datetime, time + - Dates and timestamps: date (DATE literal), datetime (TIMESTAMP literal) - Booleans: Converted to SQL boolean literals - Sequences: list, tuple, set (for IN clauses) diff --git a/pyathena/model.py b/pyathena/model.py index d76872703..283903d04 100644 --- a/pyathena/model.py +++ b/pyathena/model.py @@ -34,9 +34,9 @@ class AthenaQueryExecution: - UTILITY: Utility statements (SHOW, DESCRIBE, EXPLAIN) Example: - >>> # Typically accessed through cursor execution - >>> cursor.execute("SELECT COUNT(*) FROM my_table") - >>> query_execution = cursor._last_query_execution # Internal access + >>> # AsyncCursor returns the query execution through a Future + >>> query_id, future = cursor.execute("SELECT COUNT(*) FROM my_table") + >>> query_execution = cursor.query_execution(query_id).result() >>> print(f"Query ID: {query_execution.query_id}") >>> print(f"State: {query_execution.state}") >>> print(f"Data scanned: {query_execution.data_scanned_in_bytes} bytes") @@ -472,8 +472,8 @@ class AthenaCalculationExecution(AthenaCalculationExecutionStatus): and timing information. See Also: - AWS Athena CalculationExecution API reference: - https://docs.aws.amazon.com/athena/latest/APIReference/API_CalculationSummary.html + AWS Athena GetCalculationExecution API reference: + https://docs.aws.amazon.com/athena/latest/APIReference/API_GetCalculationExecution.html """ def __init__(self, response: dict[str, Any]) -> None: diff --git a/pyathena/options.py b/pyathena/options.py index 33d95bf9e..04a383521 100644 --- a/pyathena/options.py +++ b/pyathena/options.py @@ -35,8 +35,9 @@ class ExecuteOptions: >>> cursor.execute("SELECT ...", options=options, work_group="adhoc") Passing None for an individual keyword argument is treated as "not - provided" and leaves the corresponding ``options`` field unchanged; to - reset a field, use :meth:`merge` or construct a new instance. + provided" and leaves the corresponding ``options`` field unchanged. :meth:`merge` + also ignores None, so to clear a field, use :func:`dataclasses.replace` or + construct a new instance. Attributes: work_group: Athena workgroup to use for this query. Overrides the @@ -62,9 +63,9 @@ class ExecuteOptions: synchronous and aio cursors; ``AsyncCursor``-based cursors return the query ID directly through their execution model and do not invoke it. - result_set_type_hints: Mapping of column names (or indices) 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)"}`` """ diff --git a/pyathena/pandas/async_cursor.py b/pyathena/pandas/async_cursor.py index bd2279ea6..2328dc18a 100644 --- a/pyathena/pandas/async_cursor.py +++ b/pyathena/pandas/async_cursor.py @@ -40,9 +40,8 @@ class AsyncPandasCursor(AsyncCursor): - Memory optimization through configurable chunking Attributes: - arraysize: Number of rows to fetch per batch. - engine: Parsing engine ('auto', 'c', 'python', 'pyarrow'). - chunksize: Number of rows per chunk for large datasets. + arraysize: Default number of rows that fetchmany() returns on the result + sets this cursor creates. Example: >>> from pyathena.pandas.async_cursor import AsyncPandasCursor @@ -55,7 +54,7 @@ class AsyncPandasCursor(AsyncCursor): >>> df = result_set.as_pandas() >>> >>> # Or iterate through chunks for large datasets - >>> for chunk_df in result_set: + >>> for chunk_df in result_set.iter_chunks(): ... process_chunk(chunk_df) Note: @@ -207,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. keep_default_na: Whether to keep default pandas NA values. na_values: Additional values to treat as NA. quoting: CSV quoting behavior (pandas csv.QUOTE_* constants). diff --git a/pyathena/pandas/cursor.py b/pyathena/pandas/cursor.py index 81738492a..97bf773f9 100644 --- a/pyathena/pandas/cursor.py +++ b/pyathena/pandas/cursor.py @@ -33,8 +33,9 @@ class PandasCursor(WithFetch): """Cursor for handling pandas DataFrame results from Athena queries. This cursor returns query results as pandas DataFrames with memory-efficient - processing through chunking support and automatic chunksize optimization - for large result sets. It's ideal for data analysis and data science workflows. + processing through chunking support and optional automatic chunksize + optimization for large result sets. It's ideal for data analysis and data + science workflows. The cursor supports both regular CSV-based results and high-performance UNLOAD operations that return results in Parquet format, which is significantly @@ -44,24 +45,24 @@ class PandasCursor(WithFetch): 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(). - chunksize: Number of rows per chunk when iterating through results. Example: >>> from pyathena.pandas.cursor import PandasCursor >>> cursor = connection.cursor(PandasCursor) >>> cursor.execute("SELECT * FROM sales_data WHERE year = 2023") - >>> df = cursor.fetchall() # Returns pandas DataFrame + >>> df = cursor.as_pandas() # Returns pandas DataFrame >>> print(df.describe()) # Memory-efficient iteration for large datasets + >>> cursor = connection.cursor(PandasCursor, chunksize=50_000) >>> cursor.execute("SELECT * FROM huge_table") - >>> for chunk_df in cursor: + >>> for chunk_df in cursor.iter_chunks(): ... process_chunk(chunk_df) # Process data in chunks # High-performance UNLOAD for large datasets >>> cursor = connection.cursor(PandasCursor, unload=True) >>> cursor.execute("SELECT * FROM big_table") - >>> df = cursor.fetchall() # Faster Parquet-based result + >>> df = cursor.as_pandas() # Faster Parquet-based result """ def __init__( @@ -108,7 +109,10 @@ def __init__( auto_optimize_chunksize: Enable automatic chunksize determination for large files. Only effective when chunksize is None. Default: False (no automatic chunking). - **kwargs: Additional arguments passed to pandas.read_csv. + **kwargs: Arguments forwarded to ``WithResultSet.__init__`` and + ``BaseCursor.__init__``, such as ``arraysize``, ``connection``, + ``converter``, ``formatter``, and ``retry_config``. Pass pandas + ``read_csv``/``read_parquet`` options to ``execute()`` instead. """ super().__init__( s3_staging_dir=s3_staging_dir, @@ -183,9 +187,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. @@ -197,7 +200,7 @@ def execute( Example: >>> cursor.execute("SELECT * FROM sales WHERE year = %(year)s", ... {"year": 2023}) - >>> df = cursor.fetchall() # Returns pandas DataFrame + >>> df = cursor.as_pandas() # Returns pandas DataFrame """ self._reset_state() options = ExecuteOptions.resolve( @@ -253,6 +256,9 @@ def as_pandas(self) -> DataFrame | PandasDataFrameIterator: Returns: DataFrame when chunksize is None, PandasDataFrameIterator when chunksize is set. + + Raises: + ProgrammingError: If no result set is available. """ if not self.has_result_set: raise ProgrammingError("No result set.") diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index f7f488d1d..4ee90040b 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -190,11 +190,11 @@ class AthenaPandasResultSet(AthenaResultSet): This result set handles CSV and Parquet result files from S3, converting them to pandas DataFrames with configurable chunking for memory-efficient processing. - It automatically optimizes chunk sizes based on file size and provides iterative - processing capabilities for large datasets. + With ``auto_optimize_chunksize=True``, it chooses a chunk size based on file + size, and it provides iterative processing capabilities for large datasets. Features: - - Automatic chunk size optimization based on file size + - Optional chunk size optimization based on file size - Support for both CSV and Parquet result formats - Memory-efficient iterative processing - Automatic date/time parsing for pandas compatibility @@ -211,10 +211,10 @@ class AthenaPandasResultSet(AthenaResultSet): >>> cursor.execute("SELECT * FROM large_table") >>> >>> # Get full DataFrame - >>> df = cursor.fetchall() + >>> df = cursor.as_pandas() >>> >>> # Or iterate through chunks for memory efficiency - >>> for chunk_df in cursor: + >>> for chunk_df in cursor.iter_chunks(): ... process_chunk(chunk_df) Note: @@ -266,8 +266,11 @@ def __init__( connection: Database connection instance. converter: Data type converter for Athena types to pandas types. query_execution: Query execution metadata from Athena. - arraysize: Number of rows to fetch in each batch (not used for pandas processing). - retry_config: Retry configuration for S3 operations. + arraysize: Default number of rows that ``fetchmany()`` returns. + retry_config: Retry configuration for the ``GetQueryResults`` calls and + for the HeadObject and GetObject calls this result set makes. Result + files are read through the connection's S3 filesystem, which uses the + connection's retry configuration. keep_default_na: pandas option for handling NA values. na_values: Additional values to recognize as NA. quoting: CSV quoting behavior. @@ -281,8 +284,8 @@ def __init__( max_workers: Maximum worker threads for parallel operations. auto_optimize_chunksize: Enable automatic chunksize determination for large files when chunksize is None. - 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 pandas.read_csv/read_parquet. """ super().__init__( diff --git a/pyathena/polars/async_cursor.py b/pyathena/polars/async_cursor.py index 9e42ca78d..d119e0d60 100644 --- a/pyathena/polars/async_cursor.py +++ b/pyathena/polars/async_cursor.py @@ -223,9 +223,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. diff --git a/pyathena/polars/converter.py b/pyathena/polars/converter.py index 9cf87718a..d0b2caf6f 100644 --- a/pyathena/polars/converter.py +++ b/pyathena/polars/converter.py @@ -42,10 +42,11 @@ class DefaultPolarsTypeConverter(Converter): optimized type conversion for Polars DataFrames. The converter focuses on: - - Converting date/time types to appropriate Python objects - - Handling decimal and binary types - - Preserving JSON and complex types - - Maintaining high performance for columnar operations + - Mapping Athena types to Polars dtypes when reading CSV results, + including ``decimal`` as ``pl.Decimal(precision, scale)`` + - Converting date, time, and varbinary values, and parsing JSON values, + in rows returned by fetchone(), fetchmany(), and fetchall() + - Keeping array, map, and row values as strings Example: >>> from pyathena.polars.converter import DefaultPolarsTypeConverter diff --git a/pyathena/polars/cursor.py b/pyathena/polars/cursor.py index 53288bf91..095d1ce7a 100644 --- a/pyathena/polars/cursor.py +++ b/pyathena/polars/cursor.py @@ -183,9 +183,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. diff --git a/pyathena/polars/result_set.py b/pyathena/polars/result_set.py index 237b4ac6f..b6a2979bb 100644 --- a/pyathena/polars/result_set.py +++ b/pyathena/polars/result_set.py @@ -232,8 +232,8 @@ def __init__( chunksize: Number of rows per chunk for memory-efficient processing. If specified, data is loaded lazily in chunks for all data access methods including fetchone(), fetchmany(), and iter_chunks(). - 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 Polars read functions. """ super().__init__( diff --git a/pyathena/result_set.py b/pyathena/result_set.py index adbb36c96..410c471b3 100644 --- a/pyathena/result_set.py +++ b/pyathena/result_set.py @@ -1265,7 +1265,8 @@ def fetchone( """Fetch the next row of the result set. Returns: - A tuple representing the next row, or None if no more rows. + The next row (a tuple, or a dict for dict cursors), or None if no + more rows. Raises: ProgrammingError: If no result set is available. @@ -1282,10 +1283,11 @@ def fetchmany( """Fetch multiple rows from the result set. 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. + The fetched rows. Raises: ProgrammingError: If no result set is available. @@ -1302,7 +1304,7 @@ def fetchall( """Fetch all remaining rows from the result set. Returns: - List of tuples representing all remaining rows. + The remaining rows. Raises: ProgrammingError: If no result set is available. diff --git a/pyathena/s3fs/async_cursor.py b/pyathena/s3fs/async_cursor.py index 98bac4fc3..f2ec8a7dd 100644 --- a/pyathena/s3fs/async_cursor.py +++ b/pyathena/s3fs/async_cursor.py @@ -156,9 +156,8 @@ def _collect_result_set( Args: query_id: The Athena query execution ID. - 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. kwargs: Additional keyword arguments for result set. Returns: @@ -210,9 +209,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. diff --git a/pyathena/s3fs/cursor.py b/pyathena/s3fs/cursor.py index 839b8c128..e78b51ac2 100644 --- a/pyathena/s3fs/cursor.py +++ b/pyathena/s3fs/cursor.py @@ -159,9 +159,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. diff --git a/pyathena/spark/async_cursor.py b/pyathena/spark/async_cursor.py index 84f35ff42..1e34e617d 100644 --- a/pyathena/spark/async_cursor.py +++ b/pyathena/spark/async_cursor.py @@ -39,9 +39,7 @@ class AsyncSparkCursor(SparkBaseCursor): - Thread pool executor for concurrent operations Attributes: - max_workers: Maximum number of worker threads for async operations. session_id: The Athena Spark session ID. - engine_configuration: Spark engine configuration settings. Example: >>> from pyathena.spark.async_cursor import AsyncSparkCursor diff --git a/pyathena/spark/common.py b/pyathena/spark/common.py index b7c3948f6..0b9f2ce31 100644 --- a/pyathena/spark/common.py +++ b/pyathena/spark/common.py @@ -51,7 +51,6 @@ class SparkBaseCursor(BaseCursor, metaclass=ABCMeta): Attributes: session_id: The Athena Spark session identifier. calculation_id: ID of the current calculation being executed. - engine_configuration: DPU and resource configuration for Spark. Note: This is an abstract base class used by concrete Spark cursor implementations diff --git a/pyathena/spark/cursor.py b/pyathena/spark/cursor.py index 3376454de..87fe0ef00 100644 --- a/pyathena/spark/cursor.py +++ b/pyathena/spark/cursor.py @@ -34,8 +34,7 @@ class SparkCursor(SparkBaseCursor, WithCalculationExecution): Attributes: session_id: The Athena Spark session ID. - description: Optional description for the Spark session. - engine_configuration: Spark engine configuration settings. + description: The description of the current calculation. calculation_id: ID of the current calculation being executed. Example: @@ -49,7 +48,7 @@ class SparkCursor(SparkBaseCursor, WithCalculationExecution): ... result.show() ... ''' >>> cursor.execute(spark_code) - >>> result = cursor.fetchall() + >>> output = cursor.get_std_out() # Configure Spark session >>> cursor = connection.cursor( diff --git a/pyathena/sqlalchemy/compiler.py b/pyathena/sqlalchemy/compiler.py index 2e254a745..a68ad60ea 100644 --- a/pyathena/sqlalchemy/compiler.py +++ b/pyathena/sqlalchemy/compiler.py @@ -85,8 +85,10 @@ class AthenaTypeCompiler(GenericTypeCompiler): SQLAlchemy's portable types and Athena's specific type syntax. Athena has specific requirements for type names that differ from standard - SQL. For example, FLOAT maps to REAL in CAST expressions, and various - string types (TEXT, NCHAR, NVARCHAR) all map to STRING. + SQL. For example, FLOAT and REAL render as FLOAT here, while the statement + compiler renders them as REAL in CAST expressions. TEXT, and CHAR, NCHAR, + VARCHAR, or NVARCHAR without a length, render as STRING; with a length, + they render as CHAR(n) or VARCHAR(n). The compiler also supports Athena-specific complex types: - STRUCT/ROW: Nested record types with named fields From b445a93f859781c37aabb0eba6ebccbd1ef6e1ee Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 12:56:51 +0900 Subject: [PATCH 2/3] Describe the fetchmany() size fallback in the aio fetch mixin Co-Authored-By: Claude Opus 5.5 --- pyathena/aio/common.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pyathena/aio/common.py b/pyathena/aio/common.py index 7e581ab7f..5b04cdec0 100644 --- a/pyathena/aio/common.py +++ b/pyathena/aio/common.py @@ -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. From a27fd581fab01d74ce3918186043995e2fb27283 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:18:17 +0900 Subject: [PATCH 3/3] Apply the independent review to the async, S3 and pandas docstrings Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3_async.py | 5 +++-- pyathena/pandas/result_set.py | 2 ++ pyathena/polars/async_cursor.py | 5 +++-- 3 files changed, 8 insertions(+), 4 deletions(-) diff --git a/pyathena/filesystem/s3_async.py b/pyathena/filesystem/s3_async.py index 8dcd2ab0f..f5abdb446 100644 --- a/pyathena/filesystem/s3_async.py +++ b/pyathena/filesystem/s3_async.py @@ -60,7 +60,7 @@ class AioS3FileSystem(AsyncFileSystem): >>> files = await fs._ls('s3://my-bucket/data/') >>> >>> # Sync wrappers (auto-generated by fsspec) need an instance created - >>> # without asynchronous=True and a caller outside a running event loop + >>> # without asynchronous=True; they block the caller until done >>> files = AioS3FileSystem().ls('s3://my-bucket/data/') """ @@ -588,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. """ diff --git a/pyathena/pandas/result_set.py b/pyathena/pandas/result_set.py index 4ee90040b..0eeab6078 100644 --- a/pyathena/pandas/result_set.py +++ b/pyathena/pandas/result_set.py @@ -214,6 +214,8 @@ class AthenaPandasResultSet(AthenaResultSet): >>> df = cursor.as_pandas() >>> >>> # Or iterate through chunks for memory efficiency + >>> cursor = connection.cursor(PandasCursor, chunksize=50_000) + >>> cursor.execute("SELECT * FROM large_table") >>> for chunk_df in cursor.iter_chunks(): ... process_chunk(chunk_df) diff --git a/pyathena/polars/async_cursor.py b/pyathena/polars/async_cursor.py index d119e0d60..cf63e6ad8 100644 --- a/pyathena/polars/async_cursor.py +++ b/pyathena/polars/async_cursor.py @@ -210,8 +210,9 @@ def execute( ) -> tuple[str, Future[AthenaPolarsResultSet | Any]]: """Execute a SQL query asynchronously and return results as Polars DataFrames. - Executes the SQL query on Amazon Athena asynchronously and returns a - future that resolves to a result set for Polars DataFrame output. + Executes the SQL query on Amazon Athena asynchronously and returns the + query ID with a future that resolves to a result set for Polars + DataFrame output. Args: operation: SQL query string to execute.