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
3 changes: 2 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -134,4 +134,5 @@ Versions are derived from git tags via `hatch-vcs` — never edit `pyathena/_ver

### Google-style Docstrings

Use Google-style docstrings for public methods. See existing code for examples.
Follow the [docstring rules](docs/contributing.md#write-docstrings), which `just lint` checks with ruff's pydocstyle rules for `pyathena/`.
Mark overrides with `pyathena.util.override` instead of repeating the base method's docstring.
17 changes: 17 additions & 0 deletions docs/contributing.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,23 @@ External-fork pull requests must not be run in the project's AWS integration CI.
Do not submit an untested change expecting a maintainer to approve an AWS CI run to validate it.
Checks that need no AWS access may still run, but their success does not establish integration coverage.

## Write docstrings

Code under `pyathena/` uses [Google-style docstrings](https://google.github.io/styleguide/pyguide.html#38-comments-and-docstrings).
`just lint` checks them with ruff's pydocstyle rules.

- Modules, packages, public classes, and public functions and methods have a docstring.
Describe arguments, return values, and raised exceptions in `Args:`, `Returns:`, and `Raises:` sections.
- `__init__` describes the constructor arguments in its `Args:` section.
- A property getter has a one-line docstring; its setter needs none.
- Magic methods such as `__enter__` and `__iter__` need no docstring.
- A method that overrides a base class method is decorated with `override` from `pyathena.util`.
It needs a docstring only when its behavior differs from the base method; otherwise the API reference shows the base method's docstring.
mypy reports a missing decorator, except on an unannotated property.
- Overrides of fsspec methods are not decorated, because fsspec has no type information, so they need a docstring.
- New or changed private functions and methods use the same style.
ruff does not require them to have a docstring, but checks the docstrings they have.

## Open a pull request

Open a draft pull request with the repository's template completed.
Expand Down
1 change: 1 addition & 0 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -320,6 +320,7 @@ Result columns from a passthrough query come back with the source system's types
PyAthena provides a callback mechanism that allows you to get immediate access to the query ID
as soon as the `start_query_execution` API call is made, before waiting for query completion.
This is useful for monitoring, logging, or cancelling long-running queries from another thread.
When `cache_size` finds a reusable query, no new query starts, and the callback receives the reused query's ID.

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, evidence), full pass: CLEAN after PR-body updates

Base 2fcbcdbe, head 358d15542471bdbc3cadb1678ba7d5c04c1397ae.

Claims checked:

  • Counts: an AST comparison with master finds 511 definitions that gained a docstring (78 modules/packages, 63 __init__, 279 property getters, 91 classes/functions/methods). That is the 508 D findings plus the three D101 classes. The PR body now states it that way (the helpers' self-reported totals were higher).
  • No code change: the AST without docstrings equals master's for all 83 modules. The only removed statement is the PIE790 pass.
  • Eleven previously empty __init__.py files now carry the 2026 header (checked with git show origin/master:<f> | wc -c). scripts/check_license_headers.py passes.
  • on_start_query_execution: all ten cursors call _call_on_start_query_execution right after _execute(), and _execute() returns a cached ID from _find_previous_query_id without StartQueryExecution. The new wording, and the sentence added to docs/usage.md (anchored here), match that. The same sentence is not added to docs/sqlalchemy.md, because the SQLAlchemy dialects never pass cache_size.
  • kill_on_interrupt: checked against BaseCursor._poll/_start_execution and AioBaseCursor after Re-raise and stop interrupted queries in the SQL cursors, sharing the Spark handling #853, and against docs/usage.md "Query cancellation on interrupt" and docs/aio.md "Task cancellation". The PR body bullet was stale (it named only two docstrings) and is updated.
  • Rendered docs on this head: 145 Sphinx warnings, none new against master's 151. ruff --select D417,D415,D101 without ignore-decorators passes.
  • Commit messages: the claims in Check every docstring..., Say when on_start_query_execution... and Align the new docstrings... match the diff.
  • Existing callers and AWS operator: only docstrings, Markdown and lint configuration change. Runtime behavior is unchanged (AST check).

Deferred, as agreed with the maintainer: pre-existing docstring errors go to a separate PR, and suspected code bugs will be verified and then filed as issues.


The `on_start_query_execution` callback can be configured at both the connection level and
the execute level. When both are set, both callbacks will be invoked.
Expand Down
8 changes: 7 additions & 1 deletion pyathena/__init__.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""DB API 2.0 interface to Amazon Athena: ``connect()``, ``aio_connect()``, and type objects."""

from __future__ import annotations

import datetime
Expand Down Expand Up @@ -30,7 +32,7 @@


class DBAPITypeObject(frozenset[str]):
"""Type Objects and Constructors
"""A DB API type object that compares equal to each of its Athena type names.

https://www.python.org/dev/peps/pep-0249/#type-objects-and-constructors
"""
Expand Down Expand Up @@ -89,6 +91,8 @@ def connect(*args, **kwargs) -> Connection[Any]:
SQL queries.

Args:
*args: Positional arguments passed to the Connection constructor, in the
order of its parameters (``s3_staging_dir``, ``region_name``, ...).
s3_staging_dir: S3 location to store query results. Required if not
using workgroups or if the workgroup doesn't have a result location.
Pass an empty string to explicitly disable S3 staging and skip
Expand Down Expand Up @@ -144,6 +148,8 @@ async def aio_connect(*args, **kwargs) -> AioConnection:
and API calls, keeping the event loop free.

Args:
*args: Forwarded to ``AioConnection.create()``, which accepts keyword
arguments only.
**kwargs: Arguments forwarded to ``AioConnection.create()``.
See :func:`connect` for the full list of supported arguments.

Expand Down
8 changes: 8 additions & 0 deletions pyathena/aio/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
# Copyright 2026 The PyAthena authors
#
# Licensed under the MIT License.
# See LICENSE or https://opensource.org/licenses/MIT.
#
# SPDX-License-Identifier: MIT

"""Native asyncio connections and cursors for Amazon Athena."""
8 changes: 8 additions & 0 deletions pyathena/aio/arrow/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
# Copyright 2026 The PyAthena authors
#
# Licensed under the MIT License.
# See LICENSE or https://opensource.org/licenses/MIT.
#
# SPDX-License-Identifier: MIT

"""Native asyncio cursor that returns Athena query results as Apache Arrow tables."""
28 changes: 27 additions & 1 deletion pyathena/aio/arrow/cursor.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""Native asyncio cursor that returns Athena query results as Apache Arrow tables."""

from __future__ import annotations

import asyncio
Expand Down Expand Up @@ -54,6 +56,28 @@ def __init__(
request_timeout: float | None = None,
**kwargs,
) -> None:
"""Initialize an AioArrowCursor.

Args:
s3_staging_dir: S3 location for query results.
schema_name: Default schema name.
catalog_name: Default catalog name.
work_group: Athena workgroup name.
poll_interval: Query status polling interval in seconds.
encryption_option: S3 encryption option for query results.
kms_key: KMS key for encrypting query results.
kill_on_interrupt: Cancel the query when the task is cancelled while
``execute()`` starts or waits for the query.
unload: Whether to wrap queries in ``UNLOAD`` and read the Parquet output.
result_reuse_enable: Whether to enable Athena query result reuse.
result_reuse_minutes: Maximum age of a reused query result in minutes.
connect_timeout: Connection timeout in seconds of the pyarrow S3 filesystem
that reads the results. If None, the pyarrow default is used.
request_timeout: Request timeout in seconds of the pyarrow S3 filesystem
that reads the results. If None, the pyarrow default is used.
**kwargs: Other cursor arguments, such as ``connection`` and ``arraysize``,
passed to the parent ``__init__``.
"""
super().__init__(
s3_staging_dir=s3_staging_dir,
schema_name=schema_name,
Expand Down Expand Up @@ -111,7 +135,9 @@ async 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').
on_start_query_execution: Callback called when query starts.
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.
Expand Down
2 changes: 2 additions & 0 deletions pyathena/aio/common.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""Asyncio base cursor and the fetch mixin shared by the asyncio SQL cursors."""

from __future__ import annotations

import asyncio
Expand Down
8 changes: 8 additions & 0 deletions pyathena/aio/connection.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
#
# SPDX-License-Identifier: MIT

"""Asyncio-aware connection to Amazon Athena."""

from __future__ import annotations

import asyncio
Expand All @@ -31,6 +33,12 @@ class AioConnection(Connection[AioCursor]):
"""

def __init__(self, **kwargs: Any) -> None:
"""Initialize the connection with ``AioCursor`` as the default cursor class.

Args:
**kwargs: Arguments forwarded to ``Connection.__init__``. If they do not
include ``cursor_class``, it is set to ``AioCursor``.
"""
if "cursor_class" not in kwargs:
kwargs["cursor_class"] = AioCursor
super().__init__(**kwargs)
Expand Down
33 changes: 31 additions & 2 deletions pyathena/aio/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
#
# SPDX-License-Identifier: MIT

"""Native asyncio cursors that return rows as tuples or dictionaries."""

from __future__ import annotations

import logging
Expand Down Expand Up @@ -50,6 +52,24 @@ def __init__(
result_reuse_minutes: int = CursorIterator.DEFAULT_RESULT_REUSE_MINUTES,
**kwargs,
) -> None:
"""Initialize an AioCursor.

Args:
s3_staging_dir: S3 location for query results.
schema_name: Default schema name.
catalog_name: Default catalog name.
work_group: Athena workgroup name.
poll_interval: Query status polling interval in seconds.
encryption_option: S3 encryption option (SSE_S3, SSE_KMS, CSE_KMS).
kms_key: KMS key for encryption.
kill_on_interrupt: Cancel the query when the task is cancelled while
``execute()`` starts or waits for the query.
result_reuse_enable: Enable Athena query result reuse.
result_reuse_minutes: Maximum age in minutes of a reused result.
**kwargs: Arguments forwarded to ``WithResultSet.__init__`` and
``AioBaseCursor.__init__``, such as ``arraysize``, ``connection``,
``converter``, ``formatter``, and ``retry_config``.
"""
super().__init__(
s3_staging_dir=s3_staging_dir,
schema_name=schema_name,
Expand Down Expand Up @@ -104,12 +124,14 @@ async def execute(
parameters: Query parameters (optional).
work_group: Athena workgroup to use (optional).
s3_staging_dir: S3 location for query results (optional).
cache_size: Query result cache size (optional).
cache_size: Number of queries to check for result caching (optional).
cache_expiration_time: Cache expiration time in seconds (optional).
result_reuse_enable: Enable result reuse (optional).
result_reuse_minutes: Result reuse duration in minutes (optional).
paramstyle: Parameter style to use (optional).
on_start_query_execution: Callback called when query starts.
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.
Expand Down Expand Up @@ -225,6 +247,13 @@ class AioDictCursor(AioCursor):
"""

def __init__(self, **kwargs) -> None:
"""Initialize an AioDictCursor.

Args:
**kwargs: Arguments forwarded to ``AioCursor.__init__``. If they include
``dict_type``, it is also assigned to the class attribute
``AthenaAioDictResultSet.dict_type``, the type used to build each row.
"""
super().__init__(**kwargs)
self._result_set_class = AthenaAioDictResultSet
if "dict_type" in kwargs:
Expand Down
8 changes: 8 additions & 0 deletions pyathena/aio/pandas/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
# Copyright 2026 The PyAthena authors
#
# Licensed under the MIT License.
# See LICENSE or https://opensource.org/licenses/MIT.
#
# SPDX-License-Identifier: MIT

"""Native asyncio cursor that returns Athena query results as pandas DataFrames."""
32 changes: 31 additions & 1 deletion pyathena/aio/pandas/cursor.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""Native asyncio cursor that returns Athena query results as pandas DataFrames."""

from __future__ import annotations

import asyncio
Expand Down Expand Up @@ -63,6 +65,32 @@ def __init__(
auto_optimize_chunksize: bool = False,
**kwargs,
) -> None:
"""Initialize an AioPandasCursor.

Args:
s3_staging_dir: S3 location for query results.
schema_name: Default schema name.
catalog_name: Default catalog name.
work_group: Athena workgroup name.
poll_interval: Query status polling interval in seconds.
encryption_option: S3 encryption option for query results.
kms_key: KMS key for encrypting query results.
kill_on_interrupt: Cancel the query when the task is cancelled while
``execute()`` starts or waits for the query.
unload: Whether to wrap queries in ``UNLOAD`` and read the Parquet output.
engine: Parsing engine (``auto``, ``c``, ``python``, or ``pyarrow``).
chunksize: Number of rows per DataFrame chunk when reading CSV results. If set,
it takes precedence over ``auto_optimize_chunksize``.
block_size: Default block size of the S3 filesystem that reads the results.
cache_type: Default cache type of the S3 filesystem that reads the results.
max_workers: Maximum number of workers of the S3 filesystem.
result_reuse_enable: Whether to enable Athena query result reuse.
result_reuse_minutes: Maximum age of a reused query result in minutes.
auto_optimize_chunksize: Whether to choose a chunk size from the size of the
CSV result file when ``chunksize`` is None.
**kwargs: Other cursor arguments, such as ``connection`` and ``arraysize``,
passed to the parent ``__init__``.
"""
super().__init__(
s3_staging_dir=s3_staging_dir,
schema_name=schema_name,
Expand Down Expand Up @@ -130,7 +158,9 @@ async def execute(
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).
on_start_query_execution: Callback called when query starts.
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.
Expand Down
8 changes: 8 additions & 0 deletions pyathena/aio/polars/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
# Copyright 2026 The PyAthena authors
#
# Licensed under the MIT License.
# See LICENSE or https://opensource.org/licenses/MIT.
#
# SPDX-License-Identifier: MIT

"""Native asyncio cursor that returns Athena query results as Polars DataFrames."""
29 changes: 28 additions & 1 deletion pyathena/aio/polars/cursor.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""Native asyncio cursor that returns Athena query results as Polars DataFrames."""

from __future__ import annotations

import asyncio
Expand Down Expand Up @@ -58,6 +60,29 @@ def __init__(
chunksize: int | None = None,
**kwargs,
) -> None:
"""Initialize an AioPolarsCursor.

Args:
s3_staging_dir: S3 location for query results.
schema_name: Default schema name.
catalog_name: Default catalog name.
work_group: Athena workgroup name.
poll_interval: Query status polling interval in seconds.
encryption_option: S3 encryption option for query results.
kms_key: KMS key for encrypting query results.
kill_on_interrupt: Cancel the query when the task is cancelled while
``execute()`` starts or waits for the query.
unload: Whether to wrap queries in ``UNLOAD`` and read the Parquet output.
result_reuse_enable: Whether to enable Athena query result reuse.
result_reuse_minutes: Maximum age of a reused query result in minutes.
block_size: Default block size of the S3 filesystem that reads the results.
cache_type: Default cache type of the S3 filesystem that reads the results.
max_workers: Maximum number of workers of the S3 filesystem.
chunksize: Number of rows per chunk. If set, result files in S3 are read
lazily in chunks of this size.
**kwargs: Other cursor arguments, such as ``connection`` and ``arraysize``,
passed to the parent ``__init__``.
"""
super().__init__(
s3_staging_dir=s3_staging_dir,
schema_name=schema_name,
Expand Down Expand Up @@ -117,7 +142,9 @@ async 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').
on_start_query_execution: Callback called when query starts.
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.
Expand Down
17 changes: 17 additions & 0 deletions pyathena/aio/result_set.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
"""Asyncio result sets that fetch Athena query results with ``GetQueryResults``."""

from __future__ import annotations

import logging
Expand Down Expand Up @@ -39,6 +41,21 @@ def __init__(
retry_config: RetryConfig,
result_set_type_hints: dict[str | int, str] | None = None,
) -> None:
"""Initialize the result set without fetching rows; ``create()`` fetches the first page.

Args:
connection: The connection that ran the query.
converter: The converter for result values.
query_execution: The query execution whose results to read.
arraysize: The number of rows per ``GetQueryResults`` page and the default
``fetchmany()`` size.
retry_config: The retry configuration for API calls.
result_set_type_hints: Athena type signatures for complex-type columns,
keyed by column name (case-insensitive) or zero-based column index.

Raises:
ProgrammingError: If ``query_execution`` is not given.
"""
super().__init__(
connection=connection,
converter=converter,
Expand Down
8 changes: 8 additions & 0 deletions pyathena/aio/s3fs/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
# Copyright 2026 The PyAthena authors
#
# Licensed under the MIT License.
# See LICENSE or https://opensource.org/licenses/MIT.
#
# SPDX-License-Identifier: MIT

"""Native asyncio cursor that reads Athena CSV query results through ``AioS3FileSystem``."""
Loading
Loading