From 8bd91260c83e1b70d0906cc64f67d4b157b0d707 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 12:55:50 +0900 Subject: [PATCH 1/8] Correct user guide claims that contradict the implementation Fix the guide findings verified in #927: examples with wrong import modules, missing imports, an abstract converter class, connect() calls without a staging directory, Core select().from_statement(), json.loads() on decoded JSON, and cancellation examples that always wait out the timeout. Describe as_polars() with chunksize, the cache match rule, on_start_query_execution support, aio fetch behavior, Arrow timeout scope, ROW value types, and Athena's JSON results as measured. Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 31 ++++++--- docs/arrow.md | 27 ++++---- docs/cursor.md | 6 +- docs/filesystem.md | 12 ++-- docs/introduction.md | 2 +- docs/pandas.md | 15 +++-- docs/polars.md | 35 +++------- docs/s3fs.md | 9 ++- docs/spark.md | 3 +- docs/sqlalchemy.md | 149 ++++++++++++++++++++----------------------- docs/usage.md | 64 +++++++++---------- 11 files changed, 175 insertions(+), 178 deletions(-) diff --git a/docs/aio.md b/docs/aio.md index 9feffb7d6..df3d10a58 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -209,10 +209,14 @@ 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()`. Their fetch methods are also + wrapped in `asyncio.to_thread()`, which matters when `chunksize` is set, as fetch calls then + read S3 lazily. +- `AioS3FSCursor` streams rows from S3 as they are fetched. ```python await cursor.execute("SELECT * FROM many_rows") @@ -221,8 +225,10 @@ rows = await cursor.fetchall() 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. +Without `chunksize`, they return data that `execute()` has already loaded. +With `chunksize`, `as_pandas()` returns an iterator that reads S3 as it is iterated, and +`as_polars()` reads every remaining chunk; both read on the calling thread and block the event loop. See each cursor's documentation page for detailed usage examples. @@ -263,8 +269,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 @@ -296,7 +302,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`, outside a running event loop: + +```python +from pyathena.filesystem.s3_async import AioS3FileSystem -# Sync wrappers are auto-generated by fsspec +fs = AioS3FileSystem() files = fs.ls("s3://my-bucket/data/") ``` diff --git a/docs/arrow.md b/docs/arrow.md index a99f1d246..9cc66e731 100644 --- a/docs/arrow.md +++ b/docs/arrow.md @@ -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 @@ -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///`. +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. @@ -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)= @@ -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 @@ -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 @@ -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 @@ -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", diff --git a/docs/cursor.md b/docs/cursor.md index 5fc907d98..a3db13a2a 100644 --- a/docs/cursor.md +++ b/docs/cursor.md @@ -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 @@ -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) diff --git a/docs/filesystem.md b/docs/filesystem.md index 24b3acc46..b19f93d43 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -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 @@ -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") @@ -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. @@ -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, ) diff --git a/docs/introduction.md b/docs/introduction.md index 580d139b2..0ecf9d151 100644 --- a/docs/introduction.md +++ b/docs/introduction.md @@ -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 diff --git a/docs/pandas.md b/docs/pandas.md index 138adb40b..e2de86ca3 100644 --- a/docs/pandas.md +++ b/docs/pandas.md @@ -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). @@ -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 @@ -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 @@ -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. @@ -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 @@ -735,7 +738,7 @@ result_set = future.result() print(type(result_set.fetchone()[0])) # ``` -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 @@ -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 @@ -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", diff --git a/docs/polars.md b/docs/polars.md index 443f6f736..e522b4168 100644 --- a/docs/polars.md +++ b/docs/polars.md @@ -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 @@ -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 @@ -360,11 +346,10 @@ 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): +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 @@ -374,8 +359,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)= @@ -450,7 +435,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 @@ -508,7 +493,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 @@ -522,7 +507,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 diff --git a/docs/s3fs.md b/docs/s3fs.md index 8e33479fb..89f6a475c 100644 --- a/docs/s3fs.md +++ b/docs/s3fs.md @@ -123,14 +123,17 @@ The following type mappings are used: | char, varchar, string | str | | date | datetime.date | | timestamp | datetime.datetime | +| timestamp with time zone | datetime.datetime (timezone-aware) | | time | datetime.time | -| binary, varbinary | bytes | -| array, map, row (struct) | Parsed as Python list/dict using JSON-like parsing | +| varbinary | bytes | +| array, map, row (struct) | Parsed into Python list/dict; nested scalar values stay strings unless `result_set_type_hints` is given, and values that cannot be parsed are returned as the original string | | json | Parsed JSON (dict or list) | If you want to customize type conversion, create a converter class like this: ```python +from typing import Any + from pyathena.s3fs.converter import DefaultS3FSTypeConverter class CustomS3FSTypeConverter(DefaultS3FSTypeConverter): @@ -312,7 +315,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 `AthenaS3FSResultSet` object. -This object has an interface similar to `AthenaResultSetObject`. +This object has an interface similar to `AthenaResultSet`. ```python from pyathena import connect diff --git a/docs/spark.md b/docs/spark.md index ed05caa3f..924d5cbde 100644 --- a/docs/spark.md +++ b/docs/spark.md @@ -305,7 +305,7 @@ Athena returns the earlier calculation for a reused token, even when the code di ## AsyncSparkCursor -AsyncSparkCursor is an AsyncCursor that can handle Spark applications. +AsyncSparkCursor is an asynchronous Spark cursor that, like AsyncCursor, returns [future objects](https://docs.python.org/3/library/concurrent.futures.html#future-objects). You can use the AsyncSparkCursor by specifying the `cursor_class` with the connect method or connection object. @@ -401,6 +401,7 @@ with conn.cursor() as cursor: ``` Standard output and standard error can be retrieved by passing this object to the cursor class. +`get_std_out()` and `get_std_error()` return None instead of a future when the calculation has no standard output or standard error location. ```python from pyathena import connect diff --git a/docs/sqlalchemy.md b/docs/sqlalchemy.md index 533639012..16e244188 100644 --- a/docs/sqlalchemy.md +++ b/docs/sqlalchemy.md @@ -328,7 +328,7 @@ bucket_count Description: The number of buckets for bucketing your data. - Value: Integer value greater than or equal to 0 + Value: Positive integer Example: @@ -339,7 +339,7 @@ bucket_count All table options can also be configured with the connection string as follows: ```text -awsathena+rest://:@athena.us-west-2.amazonaws.com:443/default?s3_staging_dir=s3%3A%2F%2Fbucket%2Fpath%2Fto%2F&location=s3%3A%2F%2Fbucket%2Fpath%2Fto%2F&file_format=parquet&compression=snappy&... +awsathena+rest://:@athena.us-west-2.amazonaws.com:443/default?s3_staging_dir=s3%3A%2F%2Fbucket%2Fpath%2Fto%2F&location=s3%3A%2F%2Fbucket%2Fpath%2Fto%2F&file_format=PARQUET&compression=SNAPPY&... ``` `serdeproperties` and `tblproperties` must be converted to strings in the `'key'='value','key'='value'` format and url encoded. @@ -394,7 +394,7 @@ partition_transform_bucket_count Only has an effect for ICEBERG tables and when partition is set to true and when the partition transform is set to 'bucket' for the column. - Value: Integer value greater than or equal to 0 + Value: Positive integer Example: @@ -409,7 +409,7 @@ partition_transform_truncate_length Only has an effect for ICEBERG tables and when partition is set to true and when the partition transform is set to 'truncate' for the column. - Value: Integer value greater than or equal to 0 + Value: Positive integer Example: @@ -556,7 +556,7 @@ engine = create_engine( with engine.connect() as connection: result = connection.execute(text("SELECT * FROM many_rows")) - # query_callback will be invoked before query execution + # query_callback was invoked when the query started, before execute() waited for it ``` ### Execution options callback @@ -585,26 +585,25 @@ with engine.connect() as connection: A practical example for managing long-running analytical queries with timeout: ```python -import time -from concurrent.futures import ThreadPoolExecutor, TimeoutError +import threading from sqlalchemy import create_engine, text def run_analytics_with_timeout(): """Run analytics query with automatic timeout and cancellation.""" - - query_info = {'query_id': None, 'connection': None} + timeout_minutes = 15 + query_info = {'query_id': None} def track_query_start(query_id): query_info['query_id'] = query_id print(f"Analytics query started: {query_id}") - def timeout_monitor(timeout_minutes): - """Cancel query after timeout period.""" - time.sleep(timeout_minutes * 60) - if query_info['query_id'] and query_info['connection']: + def cancel_on_timeout(connection): + """Cancel the query that is still running when the timer fires.""" + if query_info['query_id']: try: - # Cancel via raw connection's cursor - cursor = query_info['connection'].connection.cursor() + # Cancel through a new cursor on the same DB-API connection + cursor = connection.connection.cursor() + cursor.query_id = query_info['query_id'] cursor.cancel() print(f"Query {query_info['query_id']} cancelled after {timeout_minutes}min timeout") except Exception as e: @@ -650,38 +649,31 @@ def run_analytics_with_timeout(): ORDER BY cohort_month, month_number """) - with ThreadPoolExecutor(max_workers=1) as executor: - with engine.connect() as connection: - query_info['connection'] = connection - - # Start timeout monitor (15 minutes for complex analytics) - timeout_future = executor.submit(timeout_monitor, 15) + with engine.connect() as connection: + # Cancel the query if it is still running after the timeout + timer = threading.Timer(timeout_minutes * 60, cancel_on_timeout, args=(connection,)) + timer.start() + try: + print("Starting cohort analysis (15-minute timeout)...") + result = connection.execute(analytics_query) - try: - print("Starting cohort analysis (15-minute timeout)...") - result = connection.execute(analytics_query) + # Process results + rows = result.fetchall() + print(f"Cohort analysis completed: {len(rows)} data points") - # Process results - rows = result.fetchall() - print(f"Cohort analysis completed: {len(rows)} data points") + # Show sample results + for i, row in enumerate(rows[:5]): # First 5 rows + print(f" Cohort {row.cohort_month}: Month {row.month_number}, " + f"{row.users} users, {row.retention_rate}% retention") - # Show sample results - for i, row in enumerate(rows[:5]): # First 5 rows - print(f" Cohort {row.cohort_month}: Month {row.month_number}, " - f"{row.users} users, {row.retention_rate}% retention") + if len(rows) > 5: + print(f" ... and {len(rows) - 5} more rows") - if len(rows) > 5: - print(f" ... and {len(rows) - 5} more rows") - - except Exception as e: - print(f"Analytics query failed or was cancelled: {e}") - finally: - # Clean up - query_info['connection'] = None - try: - timeout_future.result(timeout=1) - except TimeoutError: - pass # Timeout monitor still running + except Exception as e: + print(f"Analytics query failed or was cancelled: {e}") + finally: + # Stop the timer once the query has finished + timer.cancel() # Run the analytics example run_analytics_with_timeout() @@ -727,6 +719,7 @@ The `on_start_query_execution` callback is supported by all PyAthena SQLAlchemy - `awsathena+arrow` (arrow cursor) - `awsathena+polars` (polars cursor) - `awsathena+s3fs` (S3FS cursor) +- `awsathena+aiorest`, `awsathena+aiopandas`, `awsathena+aioarrow`, `awsathena+aiopolars`, and `awsathena+aios3fs` (aio cursors) Usage with different dialects: @@ -841,34 +834,32 @@ Integer fields, and integer MAP keys and values, use `INT` in that DDL. PyAthena automatically converts STRUCT data between different formats: ```python -from sqlalchemy import create_engine, select +from sqlalchemy import text # Query STRUCT data using ROW constructor result = connection.execute( - select().from_statement( - text("SELECT ROW('John Doe', 30, 'john@example.com') as profile") - ) + text("SELECT ROW('John Doe', 30, 'john@example.com') as profile") ).fetchone() -# Access STRUCT fields as dictionary -profile = result.profile # {"0": "John Doe", "1": 30, "2": "john@example.com"} +# Access STRUCT fields as dictionary; scalar values stay strings +profile = result.profile # {"0": "John Doe", "1": "30", "2": "john@example.com"} ``` #### Named STRUCT fields -For better readability, use JSON casting to get named fields: +For better readability, cast the ROW to a named ROW type, and cast that to JSON to also get typed values: ```python # Using CAST AS JSON for named field access result = connection.execute( - select().from_statement( - text("SELECT CAST(ROW('John', 30) AS JSON) as user_data") + text( + "SELECT CAST(CAST(ROW('John', 30) AS ROW(name VARCHAR, age INTEGER)) AS JSON)" + " as user_data" ) ).fetchone() -# Parse JSON result -import json -user_data = json.loads(result.user_data) # ["John", 30] +# The JSON result is already decoded +user_data = result.user_data # {"name": "John", "age": 30} ``` #### Data format support @@ -879,7 +870,7 @@ PyAthena supports multiple STRUCT data formats: ```python # Input: "{name=John, age=30}" -# Output: {"name": "John", "age": 30} +# Output: {"name": "John", "age": "30"} ``` **JSON Format (Recommended):** @@ -893,7 +884,7 @@ PyAthena supports multiple STRUCT data formats: ```python # Input: "{Alice, 25}" -# Output: {"0": "Alice", "1": 25} +# Output: {"0": "Alice", "1": "25"} ``` #### Performance considerations @@ -939,16 +930,14 @@ PyAthena supports multiple STRUCT data formats: ```python result = cursor.execute("SELECT struct_column FROM table").fetchone() -raw_data = result[0] # "{\"name\": \"John\", \"age\": 30}" -import json -parsed_data = json.loads(raw_data) +raw_data = result[0] # "{name=John, age=30}" - Athena's native ROW text ``` **After (automatic conversion):** ```python result = cursor.execute("SELECT struct_column FROM table").fetchone() -struct_data = result[0] # {"name": "John", "age": 30} - automatically converted +struct_data = result[0] # {"name": "John", "age": "30"} - automatically converted name = struct_data['name'] # Direct access ``` @@ -990,13 +979,11 @@ CREATE TABLE products ( PyAthena automatically converts MAP data between different formats: ```python -from sqlalchemy import create_engine, select +from sqlalchemy import text # Query MAP data using MAP constructor result = connection.execute( - select().from_statement( - text("SELECT MAP(ARRAY['name', 'category'], ARRAY['Laptop', 'Electronics']) as product_info") - ) + text("SELECT MAP(ARRAY['name', 'category'], ARRAY['Laptop', 'Electronics']) as product_info") ).fetchone() # Access MAP data as dictionary @@ -1010,14 +997,11 @@ For complex MAP operations, use JSON casting: ```python # Using CAST AS JSON for complex MAP operations result = connection.execute( - select().from_statement( - text("SELECT CAST(MAP(ARRAY['price', 'rating'], ARRAY['999', '4.5']) AS JSON) as data") - ) + text("SELECT CAST(MAP(ARRAY['price', 'rating'], ARRAY['999', '4.5']) AS JSON) as data") ).fetchone() -# Parse JSON result -import json -data = json.loads(result.data) # {"price": "999", "rating": "4.5"} +# The JSON result is already decoded +data = result.data # {"price": "999", "rating": "4.5"} ``` #### Data format support @@ -1073,9 +1057,7 @@ PyAthena supports multiple MAP data formats: ```python result = cursor.execute("SELECT map_column FROM table").fetchone() -raw_data = result[0] # "{\"key1\": \"value1\", \"key2\": \"value2\"}" -import json -parsed_data = json.loads(raw_data) +raw_data = result[0] # "{key1=value1, key2=value2}" - Athena's native MAP text ``` **After (automatic conversion):** @@ -1162,7 +1144,7 @@ This creates a table definition equivalent to: ```sql CREATE TABLE orders ( - id INTEGER, + id INT, item_ids ARRAY, tags ARRAY, categories ARRAY @@ -1383,8 +1365,7 @@ print(type(result.json_col)) # Athena's JSON type support has specific limitations: -- **JSON objects are fully supported** - Objects with key-value pairs work correctly -- **Top-level JSON arrays are not supported** - Direct CAST of arrays like `[1, 2, 3]` will fail +- **JSON objects and arrays are supported** - `CAST('...' AS JSON)` accepts an object or a top-level array such as `[1, 2, 3]` - **Arrays within objects are supported** - JSON objects can contain arrays as property values - **DML only** - JSON type is supported for SELECT queries but not in CREATE TABLE statements @@ -1400,9 +1381,16 @@ result = connection.execute( ).fetchone() print(result.json_col) # {"items": [1, 2, 3]} -# Not supported: Top-level array -# This will raise InvalidRequestException -# CAST('[1, 2, 3]' AS JSON) +# Supported: Top-level array +result = connection.execute( + select( + type_coerce( + literal_column("CAST('[1, 2, 3]' AS JSON)"), + JSON + ).label("json_col") + ) +).fetchone() +print(result.json_col) # [1, 2, 3] ``` #### Best practices @@ -1410,4 +1398,3 @@ print(result.json_col) # {"items": [1, 2, 3]} 1. **Use with SELECT queries** - JSON type works best for querying existing data 2. **Handle nested structures** - Objects with nested arrays and objects are fully supported 3. **Explicit type coercion** - Use `type_coerce()` when working with literal JSON values -4. **Error handling** - Be prepared to handle `InvalidRequestException` for unsupported operations diff --git a/docs/usage.md b/docs/usage.md index f62a83575..048344ce8 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -110,6 +110,7 @@ SELECT CAST(%(param)s AS TIMESTAMP(3)) AS col_timestamp If you want to use Athena's parameterized queries, you can do so by changing the `paramstyle` to `qmark` as follows. ```python +import pyathena from pyathena import connect pyathena.paramstyle = "qmark" @@ -146,7 +147,10 @@ You can find more information about the [considerations and limitations of param The `execute()` method of every SQL cursor (`Cursor`, `AsyncCursor`, the aio cursors, and their pandas/arrow/polars/s3fs variants) accepts the same set of shared keyword arguments, such as `work_group`, `s3_staging_dir`, `cache_size`, `cache_expiration_time`, `result_reuse_enable`, -`result_reuse_minutes`, `paramstyle`, `on_start_query_execution`, and `result_set_type_hints`. +`result_reuse_minutes`, `paramstyle`, and `result_set_type_hints`. +The synchronous and aio cursors also accept `on_start_query_execution`. +`AsyncCursor` and its variants return the query ID from `execute()` instead; they do not take +`on_start_query_execution` as a keyword argument and ignore it in `options`. The Spark cursors execute calculations instead of SQL queries and do not accept these arguments. These arguments can also be passed together as an `ExecuteOptions` instance using the `options` @@ -259,7 +263,7 @@ cursor.execute("SELECT * FROM one_row", cache_size=10) # re-use earlier results print(cursor.query_id) # You should expect to see the same Query ID ``` -The unit of `expiration_time` is seconds. To use the results of queries executed up to one hour ago, specify like the following. +The unit of `cache_expiration_time` is seconds. To use the results of queries executed up to one hour ago, specify like the following. ```python from pyathena import connect @@ -280,8 +284,9 @@ cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", cursor.execute("SELECT * FROM one_row", cache_size=100, cache_expiration_time=3600) # Use the last 100 queries within 1 hour as cache. ``` -Results will only be re-used if the query strings match *exactly*, -and the query was a DML statement (the assumption being that you always want to re-run queries like `CREATE TABLE` and `DROP TABLE`). +Results will only be re-used from a succeeded DML query (the assumption being that you always want to re-run queries like `CREATE TABLE` and `DROP TABLE`) +whose query string matches *exactly* after parameters are substituted, and that ran with the same schema and catalog as the cursor. +With `unload=True` on the pandas, Arrow, and Polars cursors, each query is written to a new `UNLOAD` location, so the cache never matches. The S3 staging directory is not checked, so it's possible that the location of the results is not in your provided `s3_staging_dir`. @@ -371,20 +376,18 @@ cursor.execute( A common use case is to cancel long-running analytical queries after a timeout: ```python -import time -from concurrent.futures import ThreadPoolExecutor, TimeoutError +import threading from pyathena import connect def cancel_long_running_query(): """Example: Cancel a complex analytical query after 10 minutes.""" + timeout_minutes = 10 def track_query_start(query_id): print(f"Long-running analysis started: {query_id}") - return query_id - def monitor_and_cancel(cursor, timeout_minutes): - """Monitor query and cancel if it exceeds timeout.""" - time.sleep(timeout_minutes * 60) # Convert to seconds + def cancel_on_timeout(cursor): + """Cancel the query that is still running when the timer fires.""" try: cursor.cancel() print(f"Query cancelled after {timeout_minutes} minutes timeout") @@ -430,29 +433,24 @@ def cancel_long_running_query(): ORDER BY avg_daily_events DESC """ - # Use ThreadPoolExecutor for timeout management - with ThreadPoolExecutor(max_workers=1) as executor: - # Start timeout monitor (cancel after 10 minutes) - timeout_future = executor.submit(monitor_and_cancel, cursor, 10) - - try: - print("Starting complex analytical query (10-minute timeout)...") - cursor.execute(long_query) - - # Process results - results = cursor.fetchall() - print(f"Analysis completed successfully: {len(results)} segments found") - for row in results: - print(f" {row[0]}: {row[1]} users, {row[2]:.1f} avg events") - - except Exception as e: - print(f"Query failed or was cancelled: {e}") - finally: - # Clean up timeout monitor - try: - timeout_future.result(timeout=1) - except TimeoutError: - pass # Monitor is still running, which is fine + # Cancel the query if it is still running after the timeout + timer = threading.Timer(timeout_minutes * 60, cancel_on_timeout, args=(cursor,)) + timer.start() + try: + print("Starting complex analytical query (10-minute timeout)...") + cursor.execute(long_query) + + # Process results + results = cursor.fetchall() + print(f"Analysis completed successfully: {len(results)} segments found") + for row in results: + print(f" {row[0]}: {row[1]} users, {row[2]:.1f} avg events") + + except Exception as e: + print(f"Query failed or was cancelled: {e}") + finally: + # Stop the timer once the query has finished + timer.cancel() # Run the example cancel_long_running_query() From 4d919eec2376f6c7fc57b23bf1576147245d7a61 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 12:58:38 +0900 Subject: [PATCH 2/8] Say that chunked aio cursors read S3 in fetch calls, not in execute() Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/aio.md b/docs/aio.md index df3d10a58..1fab5a686 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -213,9 +213,9 @@ All aio cursors use `await` for fetch operations, so fetching does not block the - `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()`. Their fetch methods are also - wrapped in `asyncio.to_thread()`, which matters when `chunksize` is set, as fetch calls then - read S3 lazily. + Parquet) inside `execute()`, wrapped in `asyncio.to_thread()`. With `chunksize` (pandas and + Polars), fetch calls read S3 lazily instead. The fetch methods are also wrapped in + `asyncio.to_thread()`. - `AioS3FSCursor` streams rows from S3 as they are fetched. ```python From 855fd11e23c1593d6d0a4a01a28cd850571d93c9 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:02:42 +0900 Subject: [PATCH 3/8] Point the S3FS complex-type row to the type hints section Co-Authored-By: Claude Opus 5.5 --- docs/s3fs.md | 2 +- docs/usage.md | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/docs/s3fs.md b/docs/s3fs.md index 89f6a475c..a9cdd5dd8 100644 --- a/docs/s3fs.md +++ b/docs/s3fs.md @@ -126,7 +126,7 @@ The following type mappings are used: | timestamp with time zone | datetime.datetime (timezone-aware) | | time | datetime.time | | varbinary | bytes | -| array, map, row (struct) | Parsed into Python list/dict; nested scalar values stay strings unless `result_set_type_hints` is given, and values that cannot be parsed are returned as the original string | +| array, map, row (struct) | Parsed into Python list/dict (see {ref}`usage-type-hints` for the types of nested values); values too complex to parse are returned as the original string | | json | Parsed JSON (dict or list) | If you want to customize type conversion, create a converter class like this: diff --git a/docs/usage.md b/docs/usage.md index 048344ce8..959b3aff7 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -614,6 +614,8 @@ the synchronous cursors, the `Async*` cursors, the native-async `Aio*` cursors, cursors. For Spark cursors the callback receives the per-poll `AthenaCalculationExecutionStatus` rather than an `AthenaQueryExecution`. +(usage-type-hints)= + ## Type hints for complex types *New in version 3.30.0.* From 1cddff1e9aaa0a0f31d942d4c9da965033951d2b Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:20:40 +0900 Subject: [PATCH 4/8] Apply the independent review to the aio, usage, Spark, S3FS and Polars guides Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 20 ++++++++++++-------- docs/polars.md | 3 ++- docs/s3fs.md | 2 +- docs/spark.md | 8 ++++++-- docs/usage.md | 6 +++--- 5 files changed, 24 insertions(+), 15 deletions(-) diff --git a/docs/aio.md b/docs/aio.md index 1fab5a686..76e9a2be1 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -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? @@ -213,22 +214,25 @@ All aio cursors use `await` for fetch operations, so fetching does not block the - `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()`. With `chunksize` (pandas and - Polars), fetch calls read S3 lazily instead. The fetch methods are also wrapped in - `asyncio.to_thread()`. + Parquet) inside `execute()`, wrapped in `asyncio.to_thread()`. With `chunksize` on CSV results + (pandas and Polars) or on Polars UNLOAD results, fetch calls read S3 lazily instead. The fetch + methods are also wrapped in `asyncio.to_thread()`. - `AioS3FSCursor` streams rows from S3 as they are fetched. ```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 are synchronous. -Without `chunksize`, they return data that `execute()` has already loaded. -With `chunksize`, `as_pandas()` returns an iterator that reads S3 as it is iterated, and -`as_polars()` reads every remaining chunk; both read on the calling thread and block the event loop. +When `execute()` has loaded the whole result, they return that data. +When the result is read in chunks (`chunksize`, or `auto_optimize_chunksize` with pandas), they read +S3 on the calling thread and block the event loop: `as_pandas()` returns or reads chunks as it is +iterated, and `as_polars()` reads every remaining chunk. See each cursor's documentation page for detailed usage examples. @@ -305,7 +309,7 @@ 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`, outside a running event loop: +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 diff --git a/docs/polars.md b/docs/polars.md index e522b4168..9fcc692aa 100644 --- a/docs/polars.md +++ b/docs/polars.md @@ -349,7 +349,8 @@ cursor = connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", 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) on the chunks from `iter_chunks()`: +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 diff --git a/docs/s3fs.md b/docs/s3fs.md index a9cdd5dd8..eb1f53709 100644 --- a/docs/s3fs.md +++ b/docs/s3fs.md @@ -127,7 +127,7 @@ The following type mappings are used: | time | datetime.time | | varbinary | bytes | | array, map, row (struct) | Parsed into Python list/dict (see {ref}`usage-type-hints` for the types of nested values); values too complex to parse are returned as the original string | -| json | Parsed JSON (dict or list) | +| json | Parsed JSON value (dict, list, or scalar) | If you want to customize type conversion, create a converter class like this: diff --git a/docs/spark.md b/docs/spark.md index 924d5cbde..295954fcf 100644 --- a/docs/spark.md +++ b/docs/spark.md @@ -411,8 +411,12 @@ conn = connect(work_group="YOUR_SPARK_WORKGROUP", cursor_class=AsyncSparkCursor) with conn.cursor() as cursor: calculation_id, future = cursor.execute("""spark.sql("SELECT * FROM many_rows")""") calculation_execution = future.result() - print(cursor.get_std_out(calculation_execution).result()) - print(cursor.get_std_error(calculation_execution).result()) + std_out = cursor.get_std_out(calculation_execution) + if std_out: + print(std_out.result()) + std_error = cursor.get_std_error(calculation_execution) + if std_error: + print(std_error.result()) ``` As with AsyncCursor, you need a calculation ID to cancel a calculation. diff --git a/docs/usage.md b/docs/usage.md index 959b3aff7..643c950e5 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -285,8 +285,8 @@ cursor.execute("SELECT * FROM one_row", cache_size=100, cache_expiration_time=36 ``` Results will only be re-used from a succeeded DML query (the assumption being that you always want to re-run queries like `CREATE TABLE` and `DROP TABLE`) -whose query string matches *exactly* after parameters are substituted, and that ran with the same schema and catalog as the cursor. -With `unload=True` on the pandas, Arrow, and Polars cursors, each query is written to a new `UNLOAD` location, so the cache never matches. +whose query string (with `pyformat` parameters substituted) matches *exactly*, and that ran with the same schema and catalog as the cursor. +With `unload=True` on the pandas, Arrow, and Polars cursors, a query that is wrapped in `UNLOAD` is written to a new location each time, so the cache never matches it. The S3 staging directory is not checked, so it's possible that the location of the results is not in your provided `s3_staging_dir`. @@ -425,7 +425,7 @@ def cancel_long_running_query(): ) SELECT segment, - COUNT(DISTINCT user_id) as users, + COUNT(DISTINCT dm.user_id) as users, AVG(events) as avg_daily_events FROM daily_metrics dm JOIN user_segments us ON dm.user_id = us.user_id From 7f0b100743fbbc147a165db2be0edf7b7cadaa32 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:27:54 +0900 Subject: [PATCH 5/8] Separate explicit chunking from automatic chunking in the aio guide Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/aio.md b/docs/aio.md index 76e9a2be1..bebca1934 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -230,9 +230,9 @@ 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 (`chunksize`, or `auto_optimize_chunksize` with pandas), they read -S3 on the calling thread and block the event loop: `as_pandas()` returns or reads chunks as it is -iterated, and `as_polars()` reads every remaining chunk. +When the result is read in chunks, they read S3 on the calling thread and block the event loop. +With `chunksize`, `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. From 56e13686024a0beab9071ebe6272d28f2ba5b10d Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:31:33 +0900 Subject: [PATCH 6/8] Limit lazy pandas iteration in the aio guide to CSV results Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/aio.md b/docs/aio.md index bebca1934..de458b2a2 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -231,7 +231,7 @@ 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`, `as_pandas()` returns an iterator that reads each chunk as it is iterated, and +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. From 2e447f23edffa4e57091eade93fe6c1383deb822 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:36:51 +0900 Subject: [PATCH 7/8] Cover automatic chunking and the result file in the aio fetch bullets Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/docs/aio.md b/docs/aio.md index de458b2a2..95db40417 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -214,10 +214,11 @@ All aio cursors use `await` for fetch operations, so fetching does not block the - `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()`. With `chunksize` on CSV results - (pandas and Polars) or on Polars UNLOAD results, fetch calls read S3 lazily instead. The fetch - methods are also wrapped in `asyncio.to_thread()`. -- `AioS3FSCursor` streams rows from S3 as they are fetched. + 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. ```python await cursor.execute("SELECT * FROM many_rows") From 89c8bf4536164d02a88d8f382ac25f56a54bf88e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 13:40:03 +0900 Subject: [PATCH 8/8] Note the GetQueryResults path for managed query results in the aio guide Co-Authored-By: Claude Opus 5.5 --- docs/aio.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/docs/aio.md b/docs/aio.md index 95db40417..be806a16a 100644 --- a/docs/aio.md +++ b/docs/aio.md @@ -220,6 +220,9 @@ All aio cursors use `await` for fetch operations, so fetching does not block the 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()