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
41 changes: 31 additions & 10 deletions docs/aio.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@

PyAthena provides native asyncio cursor implementations under `pyathena.aio`.
These cursors use `asyncio.sleep` for polling and `asyncio.to_thread` for boto3 calls,
keeping the event loop free without relying on thread pools for concurrency.
keeping the event loop free. Concurrency comes from asyncio tasks rather than a
`ThreadPoolExecutor` owned by the cursor; the boto3 calls run on the event loop's default executor.

## Why native asyncio?

Expand Down Expand Up @@ -209,20 +210,33 @@ Native asyncio versions are available for all cursor types:

### Fetch behavior

All aio cursors use `await` for fetch operations. The S3 download (CSV or Parquet)
happens inside `execute()`, wrapped in `asyncio.to_thread()`. Fetch methods are also
wrapped in `asyncio.to_thread()` to ensure the event loop is never blocked — this is
especially important when `chunksize` is set, as fetch calls trigger lazy S3 reads.
All aio cursors use `await` for fetch operations, so fetching does not block the event loop.

- `AioCursor` and `AioDictCursor` page through `GetQueryResults` as rows are fetched.
- `AioPandasCursor`, `AioArrowCursor`, and `AioPolarsCursor` download the result file (CSV or
Parquet) inside `execute()`, wrapped in `asyncio.to_thread()`. When CSV results are read in
chunks (`chunksize` for pandas and Polars, or a chunk size chosen by `auto_optimize_chunksize`
for pandas), or with `chunksize` on Polars UNLOAD results, fetch calls read S3 lazily instead.
The fetch methods are also wrapped in `asyncio.to_thread()`.
- `AioS3FSCursor` streams rows from the result file in S3 as they are fetched.

With managed query result storage, where the query has no S3 output location, the pandas, Arrow,
Polars, and S3FS cursors instead read every row through `GetQueryResults` inside `execute()`.

```python
await cursor.execute("SELECT * FROM many_rows")
row = await cursor.fetchone()
rows = await cursor.fetchall()

await cursor.execute("SELECT * FROM many_rows")
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.
The `as_pandas()`, `as_arrow()`, and `as_polars()` convenience methods are synchronous.
When `execute()` has loaded the whole result, they return that data.
When the result is read in chunks, they read S3 on the calling thread and block the event loop.
With `chunksize` on CSV results, `as_pandas()` returns an iterator that reads each chunk as it is iterated, and
`as_polars()` reads every remaining chunk.

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

Expand Down Expand Up @@ -263,8 +277,8 @@ parallel operations. Two implementations are provided:
- `S3AioExecutor` — dispatches work via `asyncio.run_coroutine_threadsafe` + `asyncio.to_thread`

`AioS3FileSystem` automatically uses `S3AioExecutor` for file handles, so multipart
uploads and parallel range reads are executed on the event loop without spawning
additional threads.
uploads and parallel range reads are dispatched through the event loop with
`asyncio.to_thread()` instead of a separate `ThreadPoolExecutor` per file.

### Usage with AioS3FSCursor

Expand Down Expand Up @@ -296,7 +310,14 @@ fs = AioS3FileSystem(asynchronous=True)
files = await fs._ls("s3://my-bucket/data/")
data = await fs._cat_file("s3://my-bucket/data/file.csv")
await fs._rm("s3://my-bucket/data/old/", recursive=True)
```

fsspec also generates synchronous wrappers such as `ls()`.
Call them on an instance created without `asynchronous=True`; each call blocks the caller until it completes:

```python
from pyathena.filesystem.s3_async import AioS3FileSystem

# Sync wrappers are auto-generated by fsspec
fs = AioS3FileSystem()
files = fs.ls("s3://my-bucket/data/")
```
27 changes: 15 additions & 12 deletions docs/arrow.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,15 +160,15 @@ class CustomArrowTypeConverter(Converter):
},
)

def convert(self, type_, value):
converter = self.get(type_)
return converter(value)
def convert(self, type_, value, type_hint=None):
converter = self.get(type_)
return converter(value)
```

`types` is used to explicitly specify the Arrow type when reading CSV files.
`mappings` is used as a conversion method when fetching data from a cursor object.

Then you simply specify an instance of this class in the convertes argument when creating a connection or cursor.
Then you simply specify an instance of this class in the `converter` argument when creating a connection or cursor.

```python
from pyathena import connect
Expand All @@ -189,14 +189,16 @@ cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",

If the unload option is enabled, the Parquet file itself has a schema, so the conversion is done to the Arrow type according to that schema,
and the `types` setting of the Converter class is not used.
The `mappings` are still applied to rows returned by the fetch methods.

(arrow-unload-options)=

### Unload options

ArrowCursor supports the unload option. When this option is enabled,
queries with SELECT statements are automatically converted to unload statements and executed to Athena,
and the results are output in Parquet format (Snappy compressed) to `s3_staging_dir`.
queries that start with SELECT or WITH are automatically converted to unload statements and executed to Athena,
and the results are output in Parquet format (Snappy compressed) under `{s3_staging_dir}unload/<date>/<uuid>/`.
This option requires `s3_staging_dir`; without it, `execute()` raises `ProgrammingError`.
The cursor reads the output Parquet file directly.

The output of query results with the unload statement is faster than normal query execution.
Expand Down Expand Up @@ -290,8 +292,9 @@ cursor = connect(
).cursor(ArrowCursor, connect_timeout=10.0, request_timeout=30.0)
```

The timeout parameters accept float values in seconds and apply to all S3 operations performed by the cursor,
including HeadObject and GetObject operations when retrieving query results.
The timeout parameters accept float values in seconds and apply to the pyarrow S3 filesystem that reads the result CSV and Parquet files.
The other S3 requests, the HeadObject request for the result file size and the GetObject request for the UNLOAD data manifest,
use the connection's boto3 client and its botocore configuration.

(async-arrow-cursor)=

Expand Down Expand Up @@ -365,7 +368,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
```

The return value of the [future object](https://docs.python.org/3/library/concurrent.futures.html#future-objects) is an `AthenaArrowResultSet` object.
This object has an interface similar to `AthenaResultSetObject`.
This object has an interface similar to `AthenaResultSet`.

```python
from pyathena import connect
Expand Down Expand Up @@ -429,7 +432,7 @@ print(table.schema)
print(table.shape)
```

As with AsyncArrowCursor, you need a query ID to cancel a query.
As with AsyncCursor, you need a query ID to cancel a query.

```python
from pyathena import connect
Expand All @@ -443,7 +446,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
cursor.cancel(query_id)
```

As with AsyncArrowCursor, the UNLOAD option is also available.
As with ArrowCursor, the UNLOAD option is also available.

```python
from pyathena import connect
Expand All @@ -459,7 +462,7 @@ cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",

```python
from pyathena import connect
from pyathena.arrow.cursor import AsyncArrowCursor
from pyathena.arrow.async_cursor import AsyncArrowCursor

cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
Expand Down
6 changes: 3 additions & 3 deletions docs/cursor.md
Original file line number Diff line number Diff line change
Expand Up @@ -216,10 +216,10 @@ NOTE: The cancel method of the [future object](https://docs.python.org/3/library

## AsyncDictCursor

AsyncDIctCursor is an AsyncCursor that can retrieve the query execution result
AsyncDictCursor is an AsyncCursor that can retrieve the query execution result
as a dictionary type with column names and values.

You can use the DictCursor by specifying the `cursor_class`
You can use the AsyncDictCursor by specifying the `cursor_class`
with the connect method or connection object.

```python
Expand Down Expand Up @@ -262,7 +262,7 @@ The basic usage is the same as the AsyncCursor.

```python
from pyathena.connection import Connection
from pyathena.cursor import DictCursor
from pyathena.async_cursor import AsyncDictCursor

cursor = Connection(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2").cursor(AsyncDictCursor)
Expand Down
12 changes: 8 additions & 4 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ PyAthena ships its own [fsspec](https://filesystem-spec.readthedocs.io/en/latest
filesystem implementation for Amazon S3 (`S3FileSystem`), built on boto3, with an API
surface compatible with [s3fs](https://github.com/fsspec/s3fs) for users migrating from it.

The filesystem is used internally by the pandas/polars result sets to read query results
The filesystem is used internally by the pandas, Polars, and S3FS result sets to read query results
from S3, and can also be used independently for S3 file operations.

## fsspec registration
Expand Down Expand Up @@ -48,7 +48,8 @@ s3fs-compatible credential arguments:
from pyathena import connect
from pyathena.filesystem.s3 import S3FileSystem

fs = S3FileSystem(connect(region_name="us-west-2"))
fs = S3FileSystem(connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2"))

# Or with direct credentials (s3fs-compatible arguments).
fs = S3FileSystem(key="YOUR_ACCESS_KEY", secret="YOUR_SECRET_KEY")
Expand Down Expand Up @@ -139,7 +140,10 @@ reading. Explicit versions can always be read with the `?versionId=` suffix or t
`version_id` argument.

```python
fs = S3FileSystem(connect(region_name="us-west-2"), version_aware=True)
fs = S3FileSystem(
connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", region_name="us-west-2"),
version_aware=True,
)

with fs.open("s3://YOUR_S3_BUCKET/path/to/object", "rb") as f:
data = f.read() # Pinned to the version observed at open time.
Expand All @@ -164,7 +168,7 @@ create or delete a bucket. Pass the opt-in flags to enable them:

```python
fs = S3FileSystem(
connect(region_name="us-west-2"),
connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", region_name="us-west-2"),
allow_bucket_creation=True,
allow_bucket_deletion=True,
)
Expand Down
2 changes: 1 addition & 1 deletion docs/introduction.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ PyAthena provides comprehensive support for Amazon Athena's data types and featu

**Additional Features:**

- **Connection Management**: Efficient connection pooling and configuration
- **Connection Management**: Flexible connection configuration (credentials, role assumption, workgroups)
- **Result Caching**: Athena query result reuse capabilities
- **Error Handling**: Comprehensive exception handling and recovery
- **S3 Integration**: Direct S3 data access and staging support
Expand Down
15 changes: 9 additions & 6 deletions docs/pandas.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ cursor.execute("SHOW PARTITIONS YOUR_TABLE")
print(cursor.fetchall())
```

`to_sql` writes the data as Parquet with pyarrow, so it requires `pip install PyAthena[Pandas,Arrow]`.
Conversion to Parquet and upload to S3 use [ThreadPoolExecutor](https://docs.python.org/3/library/concurrent.futures.html#threadpoolexecutor) by default.
It is also possible to use [ProcessPoolExecutor](https://docs.python.org/3/library/concurrent.futures.html#processpoolexecutor).

Expand Down Expand Up @@ -292,7 +293,7 @@ class CustomPandasTypeConverter(Converter):

Specify the combination of converter functions in the mappings argument and the dtypes combination in the types argument.

Then you simply specify an instance of this class in the convertes argument when creating a connection or cursor.
Then you simply specify an instance of this class in the `converter` argument when creating a connection or cursor.

```python
from pyathena import connect
Expand Down Expand Up @@ -379,7 +380,8 @@ awsathena+pandas://:@athena.{region_name}.amazonaws.com:443/{schema_name}?s3_sta
```

When this option is used, the object returned by the as_pandas method is a `PandasDataFrameIterator` object.
This object has exactly the same interface as the `TextFileReader` object and can be handled in the same way.
This object is an iterator of DataFrames, like the pandas `TextFileReader` object.
It also provides `get_chunk()`, `as_pandas()`, `iterrows()`, and `close()`, and can be used as a context manager.

```python
from pyathena import connect
Expand Down Expand Up @@ -557,6 +559,7 @@ For example, `dtype` replaces the whole column type mapping, and `parse_dates` r
### Unload options

PandasCursor also supports the unload option, as does {ref}`arrow-cursor`.
Reading the unloaded Parquet files requires pyarrow (`pip install PyAthena[Pandas,Arrow]`), and `engine` must be `auto` or `pyarrow`.

See {ref}`arrow-unload-options` for more information.

Expand Down Expand Up @@ -662,7 +665,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
```

The return value of the [future object](https://docs.python.org/3/library/concurrent.futures.html#future-objects) is an `AthenaPandasResultSet` object.
This object has an interface similar to `AthenaResultSetObject`.
This object has an interface similar to `AthenaResultSet`.

```python
from pyathena import connect
Expand Down Expand Up @@ -735,7 +738,7 @@ result_set = future.result()
print(type(result_set.fetchone()[0])) # <class 'pandas._libs.tslibs.timestamps.Timestamp'>
```

As with AsyncPandasCursor, you need a query ID to cancel a query.
As with AsyncCursor, you need a query ID to cancel a query.

```python
from pyathena import connect
Expand All @@ -749,7 +752,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
cursor.cancel(query_id)
```

As with AsyncPandasCursor, the unload option is also available.
As with PandasCursor, the unload option is also available.

```python
from pyathena import connect
Expand All @@ -765,7 +768,7 @@ cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",

```python
from pyathena import connect
from pyathena.pandas.cursor import AsyncPandasCursor
from pyathena.pandas.async_cursor import AsyncPandasCursor

cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
Expand Down
36 changes: 11 additions & 25 deletions docs/polars.md
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,7 @@ for row in cursor:
in chunks, which is more efficient for batch processing:

```python
import polars as pl
from pyathena import connect
from pyathena.polars.cursor import PolarsCursor

Expand Down Expand Up @@ -335,23 +336,8 @@ for chunk in cursor.iter_chunks():
process_chunk(chunk)
```

When the chunksize option is used, the object returned by the `as_polars` method is a `PolarsDataFrameIterator` object.
This object provides the same chunked iteration interface and can be used in the same way:

```python
from pyathena import connect
from pyathena.polars.cursor import PolarsCursor

cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
cursor_class=PolarsCursor).cursor(chunksize=50_000)
df_iter = cursor.execute("SELECT * FROM many_rows").as_polars()
for df in df_iter:
print(df.describe())
print(df.head())
```

The `PolarsDataFrameIterator` also has an `as_polars()` method that collects all chunks into a single DataFrame:
With the chunksize option, the `as_polars` method still returns a single `polars.DataFrame`.
It reads every chunk and concatenates them, so the whole result is loaded into memory:

```python
from pyathena import connect
Expand All @@ -360,11 +346,11 @@ from pyathena.polars.cursor import PolarsCursor
cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
cursor_class=PolarsCursor).cursor(chunksize=50_000)
df_iter = cursor.execute("SELECT * FROM many_rows").as_polars()
df = df_iter.as_polars() # Collect all chunks into a single DataFrame
df = cursor.execute("SELECT * FROM many_rows").as_polars() # All chunks in a single DataFrame
```

This is equivalent to using [polars.concat](https://docs.pola.rs/api/python/stable/reference/api/polars.concat.html):
Apart from returning an empty DataFrame when there are no chunks, this is equivalent to using
[polars.concat](https://docs.pola.rs/api/python/stable/reference/api/polars.concat.html) on the chunks from `iter_chunks()`:

```python
import polars as pl
Expand All @@ -374,8 +360,8 @@ from pyathena.polars.cursor import PolarsCursor
cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
cursor_class=PolarsCursor).cursor(chunksize=50_000)
df_iter = cursor.execute("SELECT * FROM many_rows").as_polars()
df = pl.concat(list(df_iter))
cursor.execute("SELECT * FROM many_rows")
df = pl.concat(list(cursor.iter_chunks()))
```

(async-polars-cursor)=
Expand Down Expand Up @@ -450,7 +436,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
```

The return value of the [future object](https://docs.python.org/3/library/concurrent.futures.html#future-objects) is an `AthenaPolarsResultSet` object.
This object has an interface similar to `AthenaResultSetObject`.
This object has an interface similar to `AthenaResultSet`.

```python
from pyathena import connect
Expand Down Expand Up @@ -508,7 +494,7 @@ print(df.describe())
print(df.head())
```

As with AsyncPolarsCursor, you need a query ID to cancel a query.
As with AsyncCursor, you need a query ID to cancel a query.

```python
from pyathena import connect
Expand All @@ -522,7 +508,7 @@ query_id, future = cursor.execute("SELECT * FROM many_rows")
cursor.cancel(query_id)
```

As with AsyncPolarsCursor, the unload option is also available.
As with PolarsCursor, the unload option is also available.

```python
from pyathena import connect
Expand Down
Loading
Loading