Native Asyncio Cursors

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. 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?

PyAthena has two families of async cursors:

AsyncCursor

AioCursor

Concurrency model

concurrent.futures.ThreadPoolExecutor

Native asyncio (await / async for)

Event loop

Blocks a thread per query

Non-blocking

Connection

connect() (sync)

aio_connect() (async)

execute() returns

(query_id, Future)

Awaitable cursor (self)

Fetch methods

Sync (via Future.result())

await cursor.fetchone() for streaming cursors

Iteration

for row in result_set

async for row in cursor

Context manager

with conn.cursor() as cursor

async with conn.cursor() as cursor

Best for

Adding concurrency to sync code

Async frameworks (FastAPI, aiohttp, etc.)

Choose AioCursor when your application already uses asyncio (e.g., web frameworks, async pipelines). Choose AsyncCursor when you want simple parallel query execution from synchronous code.

Connection

Use the aio_connect() function to create an async connection. It returns an AioConnection that produces AioCursor instances by default.

from pyathena import aio_connect

conn = await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                      region_name="us-west-2")

The connection supports the async context manager protocol:

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    cursor = conn.cursor()
    await cursor.execute("SELECT 1")
    print(await cursor.fetchone())

AioCursor

AioCursor is a native asyncio cursor that uses await for query execution and result fetching. It follows the DB API 2.0 interface adapted for async usage.

from pyathena import aio_connect
from pyathena.aio.cursor import AioCursor

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    cursor = conn.cursor()
    await cursor.execute("SELECT * FROM many_rows")
    print(await cursor.fetchone())
    print(await cursor.fetchmany(10))
    print(await cursor.fetchall())

The cursor supports the async with context manager:

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    async with conn.cursor() as cursor:
        await cursor.execute("SELECT * FROM many_rows")
        rows = await cursor.fetchall()

You can iterate over results with async for:

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    async with conn.cursor() as cursor:
        await cursor.execute("SELECT * FROM many_rows")
        async for row in cursor:
            print(row)

Execution information of the query can also be retrieved:

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    async with conn.cursor() as cursor:
        await cursor.execute("SELECT * FROM many_rows")
        print(cursor.state)
        print(cursor.state_change_reason)
        print(cursor.completion_date_time)
        print(cursor.submission_date_time)
        print(cursor.data_scanned_in_bytes)
        print(cursor.engine_execution_time_in_millis)
        print(cursor.query_queue_time_in_millis)
        print(cursor.total_execution_time_in_millis)
        print(cursor.query_planning_time_in_millis)
        print(cursor.service_processing_time_in_millis)
        print(cursor.output_location)

To cancel a running query:

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    async with conn.cursor() as cursor:
        await cursor.execute("SELECT * FROM many_rows")
        await cursor.cancel()

Task cancellation

With kill_on_interrupt enabled, which is the default, cancelling the task while execute() waits for the query requests cancellation of the query, waits until it reaches a terminal state, and then raises asyncio.CancelledError. Cancellation is a best-effort request, so the query can still end as SUCCEEDED or FAILED. The query_id property keeps the ID of the cancelled query. If the cancellation request or that wait fails, asyncio.CancelledError is raised with the error as its cause. Cancelling the task while execute() is still starting the query first waits for the start request to finish, and then cancels the query it started in the same way. If the task is cancelled before execute() begins the request, the request is never sent. Cancelling the task again during the cancellation request or these waits raises asyncio.CancelledError immediately, and the query can keep running. With kill_on_interrupt=False, asyncio.CancelledError is raised immediately, and a query that has already started keeps running.

With kill_on_interrupt enabled, a timeout from asyncio.wait_for() while execute() starts or waits for the query therefore cancels it and raises asyncio.TimeoutError. If the timeout expires while execute() is still looking up a cached result, no query is started and query_id is None.

import asyncio

from pyathena import aio_connect

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    async with conn.cursor() as cursor:
        try:
            await asyncio.wait_for(cursor.execute("SELECT * FROM many_rows"), timeout=60)
        except asyncio.TimeoutError:
            print(f"Query timed out: {cursor.query_id}")

AioDictCursor

AioDictCursor is an AioCursor that returns rows as dictionaries with column names as keys.

from pyathena import aio_connect
from pyathena.aio.cursor import AioDictCursor

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    cursor = conn.cursor(AioDictCursor)
    await cursor.execute("SELECT * FROM many_rows LIMIT 10")
    async for row in cursor:
        print(row["a"])

If you want to change the dictionary type (e.g., use OrderedDict):

from collections import OrderedDict
from pyathena import aio_connect
from pyathena.aio.cursor import AioDictCursor

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    cursor = conn.cursor(AioDictCursor, dict_type=OrderedDict)
    await cursor.execute("SELECT * FROM many_rows LIMIT 10")
    async for row in cursor:
        print(row)

Specialized Aio Cursors

Native asyncio versions are available for all cursor types:

Cursor

Module

Result format

AioPandasCursor

pyathena.aio.pandas.cursor

pandas DataFrame

AioArrowCursor

pyathena.aio.arrow.cursor

pyarrow Table

AioPolarsCursor

pyathena.aio.polars.cursor

polars DataFrame

AioS3FSCursor

pyathena.aio.s3fs.cursor

Row tuples (lightweight)

AioSparkCursor

pyathena.aio.spark.cursor

PySpark execution

Fetch behavior

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().

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 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. With a chunk size chosen by auto_optimize_chunksize, as_pandas() reads every remaining chunk.

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

AioS3FileSystem

AioS3FileSystem is a native asyncio filesystem interface for Amazon S3, built on fsspec’s AsyncFileSystem. It provides the same functionality as S3FileSystem but uses asyncio.gather with asyncio.to_thread for parallel operations instead of ThreadPoolExecutor.

Why AioS3FileSystem?

The synchronous S3FileSystem uses ThreadPoolExecutor for parallel S3 operations (batch deletes, multipart uploads, range reads). When used from within an asyncio application via AioS3FSCursor, this creates a thread-in-thread pattern: the cursor wraps calls in asyncio.to_thread(), and inside that thread S3FileSystem spawns additional threads via ThreadPoolExecutor.

AioS3FileSystem eliminates this inefficiency by dispatching all parallel operations through the asyncio event loop.

S3FileSystem

AioS3FileSystem

Parallelism

ThreadPoolExecutor

asyncio.gather + asyncio.to_thread

File handles

S3File with thread pool

AioS3File with S3AioExecutor

Bulk delete

Thread pool per batch

asyncio.gather per batch

Multipart copy

Thread pool per part

asyncio.gather per part

Best for

Synchronous applications

Async frameworks (FastAPI, aiohttp, etc.)

Executor strategy

S3FileSystem and S3File use a pluggable executor abstraction (S3Executor) for parallel operations. Two implementations are provided:

  • S3ThreadPoolExecutor — wraps ThreadPoolExecutor (default for sync usage)

  • 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 dispatched through the event loop with asyncio.to_thread() instead of a separate ThreadPoolExecutor per file. At most max_workers of them run at once. An instance created with asynchronous=True has no event loop of its own, so its file handles use S3ThreadPoolExecutor.

Usage with AioS3FSCursor

AioS3FSCursor automatically uses AioS3FileSystem internally. No additional configuration is needed:

from pyathena import aio_connect
from pyathena.aio.s3fs.cursor import AioS3FSCursor

async with await aio_connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
                          region_name="us-west-2") as conn:
    cursor = conn.cursor(AioS3FSCursor)
    await cursor.execute("SELECT * FROM many_rows")
    async for row in cursor:
        print(row)

Standalone usage

AioS3FileSystem can also be used directly for S3 operations:

from pyathena.filesystem.s3_async import AioS3FileSystem

# Async context
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:

from pyathena.filesystem.s3_async import AioS3FileSystem

fs = AioS3FileSystem()
files = fs.ls("s3://my-bucket/data/")