Apache Arrow Integration¶
This section covers Apache Arrow-specific cursors, result sets, and data converters.
Arrow Cursors¶
- class pyathena.arrow.cursor.ArrowCursor(s3_staging_dir: str | None = None, schema_name: str | None = None, catalog_name: str | None = None, work_group: str | None = None, poll_interval: float = 1, encryption_option: str | None = None, kms_key: str | None = None, kill_on_interrupt: bool = True, unload: bool = False, result_reuse_enable: bool = False, result_reuse_minutes: int = 60, connect_timeout: float | None = None, request_timeout: float | None = None, **kwargs)[source]¶
Cursor for handling Apache Arrow Table results from Athena queries.
This cursor returns query results as Apache Arrow Tables, which provide efficient columnar data processing and memory usage. Arrow Tables are especially useful for analytical workloads and data science applications.
The cursor supports both regular CSV-based results and high-performance UNLOAD operations that return results in Parquet format for improved performance with large datasets.
- description¶
Sequence of column descriptions for the last query.
- rowcount¶
Number of rows affected by the last query (-1 for SELECT queries).
- arraysize¶
Default number of rows to fetch with fetchmany().
Example
>>> from pyathena.arrow.cursor import ArrowCursor >>> cursor = connection.cursor(ArrowCursor) >>> cursor.execute("SELECT * FROM large_table") >>> table = cursor.as_arrow() # Returns pyarrow.Table >>> df = table.to_pandas() # Convert to pandas if needed
# High-performance UNLOAD for large datasets >>> cursor = connection.cursor(ArrowCursor, unload=True) >>> cursor.execute(“SELECT * FROM huge_table”) >>> table = cursor.as_arrow() # Faster Parquet-based result
- __init__(s3_staging_dir: str | None = None, schema_name: str | None = None, catalog_name: str | None = None, work_group: str | None = None, poll_interval: float = 1, encryption_option: str | None = None, kms_key: str | None = None, kill_on_interrupt: bool = True, unload: bool = False, result_reuse_enable: bool = False, result_reuse_minutes: int = 60, connect_timeout: float | None = None, request_timeout: float | None = None, **kwargs) None[source]¶
Initialize an ArrowCursor.
- Parameters:
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 ARN for encryption.
kill_on_interrupt – Cancel running query on keyboard interrupt.
unload – Enable UNLOAD for high-performance Parquet output.
result_reuse_enable – Enable Athena query result reuse.
result_reuse_minutes – Minutes to reuse cached results.
connect_timeout – Socket connection timeout in seconds for S3 operations. Defaults to AWS SDK default (typically 1 second) if not specified.
request_timeout – Request timeout in seconds for S3 operations. Defaults to AWS SDK default (typically 3 seconds) if not specified. Increase this value if you experience timeout errors when using role assumption with STS or have high latency to S3.
**kwargs – Additional connection parameters.
Example
>>> # Use higher timeouts for role assumption scenarios >>> cursor = connection.cursor( ... ArrowCursor, ... connect_timeout=10, ... request_timeout=30 ... )
- static get_default_converter(unload: bool = False) DefaultArrowTypeConverter | DefaultArrowUnloadTypeConverter | Any[source]¶
Get the default type converter for this cursor class.
- Parameters:
unload – Whether the converter is for UNLOAD operations. Some cursor types may return different converters for UNLOAD operations.
- Returns:
The default type converter instance for this cursor type.
- execute(operation: str, parameters: dict[str, Any] | list[str] | None = None, work_group: str | None = None, s3_staging_dir: str | None = None, cache_size: int | None = None, cache_expiration_time: int | None = None, result_reuse_enable: bool | None = None, result_reuse_minutes: int | None = None, paramstyle: str | None = None, on_start_query_execution: Callable[[str], None] | None = None, result_set_type_hints: dict[str | int, str] | None = None, *, options: ExecuteOptions | None = None, **kwargs) ArrowCursor[source]¶
Execute a SQL query and return results as Apache Arrow Tables.
Executes the SQL query on Amazon Athena and configures the result set for Apache Arrow Table output. Arrow format provides high-performance columnar data processing with efficient memory usage.
- Parameters:
operation – SQL query string to execute.
parameters – Query parameters for parameterized queries.
work_group – Athena workgroup to use for this query.
s3_staging_dir – S3 location for query results.
cache_size – Number of queries to check for result caching.
cache_expiration_time – Cache expiration time in seconds.
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 invoked with the query ID before
execute()waits for the query: after theStartQueryExecutioncall, or after a reusable query ID is found throughcache_size.result_set_type_hints – Athena type signatures for complex-type columns, keyed by column name (case-insensitive) or zero-based column index.
options – Shared execution options as an
ExecuteOptionsinstance. Individual keyword arguments take precedence overoptionsfields.**kwargs – Additional execution parameters.
- Returns:
Self reference for method chaining.
Example
>>> cursor.execute("SELECT * FROM sales WHERE year = 2023") >>> table = cursor.as_arrow() # Returns Apache Arrow Table
- as_arrow() Table[source]¶
Return query results as an Apache Arrow Table.
Converts the entire result set into an Apache Arrow Table for efficient columnar data processing. Arrow Tables provide excellent performance for analytical workloads and interoperability with other data processing frameworks.
- Returns:
Apache Arrow Table containing all query results.
- Raises:
ProgrammingError – If no query has been executed or no results are available.
Example
>>> cursor = connection.cursor(ArrowCursor) >>> cursor.execute("SELECT * FROM my_table") >>> table = cursor.as_arrow() >>> print(f"Table has {table.num_rows} rows and {table.num_columns} columns")
- as_polars() pl.DataFrame[source]¶
Return query results as a Polars DataFrame.
Converts the Apache Arrow Table to a Polars DataFrame for interoperability with the Polars data processing library.
- Returns:
Polars DataFrame containing all query results.
- Raises:
ProgrammingError – If no query has been executed or no results are available.
ImportError – If polars is not installed.
Example
>>> cursor = connection.cursor(ArrowCursor) >>> cursor.execute("SELECT * FROM my_table") >>> df = cursor.as_polars() >>> print(f"DataFrame has {df.height} rows and {df.width} columns")
- DEFAULT_RESULT_REUSE_MINUTES = 60¶
- LIST_DATABASES_MAX_RESULTS = 50¶
- LIST_QUERY_EXECUTIONS_MAX_RESULTS = 50¶
- LIST_TABLE_METADATA_MAX_RESULTS = 50¶
- property arraysize: int¶
The default number of rows per
fetchmany()call.execute()passes it to the new result set, so a change applies to the result sets of later executions. Setting it to zero or a negative value raisesProgrammingError.- Returns:
The default number of rows per
fetchmany()call.
- cancel() None¶
Cancel the currently executing query.
- Raises:
ProgrammingError – If no query is currently executing.
- property connection: Connection[Any]¶
The connection that created this cursor.
- property data_manifest_location: str | None¶
The S3 location of the data manifest that lists the files the query wrote.
- property description: list[tuple[str, str, None, None, int, int, str]] | None¶
The DB API 2.0 column descriptions of the result set, or None without one.
- property encryption_option: str | None¶
The
EncryptionOptionof the query results, such asSSE_S3orSSE_KMS.
- property engine_execution_time_in_millis: int | None¶
The time in milliseconds that the query engine took to run the query.
- property error_category: int | None¶
1 for system, 2 for user, 3 for other.
- Type:
The
ErrorCategoryof the failure
- executemany(operation: str, seq_of_parameters: list[dict[str, Any] | list[str] | None], **kwargs) None¶
Execute a SQL query multiple times with different parameters.
On success,
rowcountis the sum of the affected row counts, or -1 if any execution has an unknown count. An empty parameter list sets it to 0. On failure, it is -1; earlier executions are not rolled back. Result sets are discarded.On failure,
query_idretains the current query ID when available. If parameter iteration fails, this can identify the last successful execution.- Parameters:
operation – SQL query string to execute.
seq_of_parameters – Sequence of parameter sets, one per execution.
**kwargs – Additional keyword arguments passed to each
execute().
- property expected_bucket_owner: str | None¶
The AWS account ID expected to own the S3 bucket of the query results.
- fetchall() list[tuple[Any | None, ...] | dict[Any, Any | None]]¶
Fetch all remaining rows from the result set.
- Returns:
The remaining rows.
- Raises:
ProgrammingError – If no result set is available.
- fetchmany(size: int | None = None) list[tuple[Any | None, ...] | dict[Any, Any | None]]¶
Fetch multiple rows from the result set.
- Parameters:
size – Maximum number of rows to fetch. If None or not positive,
arraysizeis used.- Returns:
The fetched rows.
- Raises:
ProgrammingError – If no result set is available.
- fetchone() tuple[Any | None, ...] | dict[Any, Any | None] | None¶
Fetch the next row of the result set.
- Returns:
The next row (a tuple, or a dict for dict cursors), or None if no more rows.
- Raises:
ProgrammingError – If no result set is available.
- get_table_metadata(table_name: str, catalog_name: str | None = None, schema_name: str | None = None, logging_: bool = True) AthenaTableMetadata¶
Get one table’s metadata.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
table_name – The table name.
catalog_name – The catalog, or None for the cursor’s catalog.
schema_name – The database, or None for the cursor’s schema.
logging – Whether to log a failed request.
- Returns:
The table’s metadata.
- Raises:
OperationalError – If the request fails, including when the table does not exist.
- list_databases(catalog_name: str | None, max_results: int | None = None) list[AthenaDatabase]¶
List the catalog’s databases.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
catalog_name – The catalog, or None for the cursor’s catalog.
max_results – The page size of each request.
- Returns:
The catalog’s databases.
- Raises:
OperationalError – If the request fails.
- list_table_metadata(catalog_name: str | None = None, schema_name: str | None = None, expression: str | None = None, max_results: int | None = None, logging_: bool = True) list[AthenaTableMetadata]¶
List a database’s table metadata.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
catalog_name – The catalog, or None for the cursor’s catalog.
schema_name – The database, or None for the cursor’s schema.
expression – A table name pattern.
max_results – The page size of each request.
logging – Whether to log a failed request.
- Returns:
The metadata of the database’s tables.
- Raises:
OperationalError – If the request fails.
- property query_id: str | None¶
The query execution ID of the last execution.
With
cache_sizeorcache_expiration_time, this can be the ID of a previous execution whose result is reused.- Returns:
The query execution ID, or None if there is none since the last reset.
- property query_planning_time_in_millis: int | None¶
The time in milliseconds that Athena took to plan the query.
- property query_queue_time_in_millis: int | None¶
The time in milliseconds that the query waited in the queue.
- property result_reuse_enabled: bool | None¶
Whether reuse of previous query results by age is enabled for the query.
- property result_reuse_minutes: int | None¶
The maximum age in minutes of a previous query result that Athena can reuse.
- property result_set: AthenaResultSet | None¶
The result set of the last executed query.
- Returns:
The result set, or None before a query succeeds or after a reset.
- property reused_previous_result: bool | None¶
Whether Athena reused a previous query result instead of running the query.
- property rowcount: int¶
Get the number of rows affected by the last operation.
For SELECT statements, this returns -1 as per DB API 2.0 specification. For DML operations (INSERT, UPDATE, DELETE) and CTAS, this returns the number of affected rows. After a successful
executemany(), this is the sum across executions, or -1 if any count is unknown.- Returns:
The number of rows, or -1 if not applicable or unknown.
- property rownumber: int | None¶
The zero-based index of the next row in the result set.
- Returns:
The row index, or None if there is no result set or the index is unknown.
- property s3_acl_option: str | None¶
The
S3AclOptionof the query results, such asBUCKET_OWNER_FULL_CONTROL.
- property service_processing_time_in_millis: int | None¶
The time in milliseconds that Athena took to publish the query results.
- setinputsizes(sizes)¶
Accept input sizes as DB API 2.0 requires, and ignore them.
- Parameters:
sizes – Sequence of parameter types or sizes.
- setoutputsize(size, column=None)¶
Accept a column buffer size as DB API 2.0 requires, and ignore it.
- Parameters:
size – Buffer size for large columns.
column – Index of the column the size applies to, or None for all large columns.
- property state_change_reason: str | None¶
The
StateChangeReasonthat gives further detail about the state.
- class pyathena.arrow.async_cursor.AsyncArrowCursor(s3_staging_dir: str | None = None, schema_name: str | None = None, catalog_name: str | None = None, work_group: str | None = None, poll_interval: float = 1, encryption_option: str | None = None, kms_key: str | None = None, kill_on_interrupt: bool = True, max_workers: int = 20, arraysize: int = 1000, unload: bool = False, result_reuse_enable: bool = False, result_reuse_minutes: int = 60, connect_timeout: float | None = None, request_timeout: float | None = None, **kwargs)[source]¶
Asynchronous cursor that returns results in Apache Arrow format.
This cursor extends AsyncCursor to provide asynchronous query execution with results returned as Apache Arrow Tables or RecordBatches. It’s optimized for high-performance analytics workloads and interoperability with the Apache Arrow ecosystem.
- Features:
Asynchronous query execution with concurrent futures
Apache Arrow columnar data format for high performance
Memory-efficient processing of large datasets
Support for UNLOAD operations with Parquet output
Integration with pandas, Polars, and other Arrow-compatible libraries
- arraysize¶
Number of rows to fetch per batch (configurable).
Example
>>> from pyathena.arrow.async_cursor import AsyncArrowCursor >>> >>> cursor = connection.cursor(AsyncArrowCursor, unload=True) >>> query_id, future = cursor.execute("SELECT * FROM large_table") >>> >>> # Get result when ready >>> result_set = future.result() >>> arrow_table = result_set.as_arrow() >>> >>> # Convert to pandas if needed >>> df = arrow_table.to_pandas() >>> >>> # Convert to Polars if needed (requires polars) >>> polars_df = result_set.as_polars()
Note
Requires pyarrow to be installed. UNLOAD operations generate Parquet files in S3 for optimal Arrow compatibility. For Polars interoperability, polars must be installed separately.
- __init__(s3_staging_dir: str | None = None, schema_name: str | None = None, catalog_name: str | None = None, work_group: str | None = None, poll_interval: float = 1, encryption_option: str | None = None, kms_key: str | None = None, kill_on_interrupt: bool = True, max_workers: int = 20, arraysize: int = 1000, unload: bool = False, result_reuse_enable: bool = False, result_reuse_minutes: int = 60, connect_timeout: float | None = None, request_timeout: float | None = None, **kwargs) None[source]¶
Initialize an AsyncArrowCursor.
- Parameters:
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 ARN for encryption.
kill_on_interrupt – Cancel running query on keyboard interrupt.
max_workers – Maximum number of workers for concurrent execution.
arraysize – Number of rows to fetch per batch.
unload – Enable UNLOAD for high-performance Parquet output.
result_reuse_enable – Enable Athena query result reuse.
result_reuse_minutes – Minutes to reuse cached results.
connect_timeout – Socket connection timeout in seconds for S3 operations. Defaults to AWS SDK default (typically 1 second) if not specified.
request_timeout – Request timeout in seconds for S3 operations. Defaults to AWS SDK default (typically 3 seconds) if not specified. Increase this value if you experience timeout errors when using role assumption with STS or have high latency to S3.
**kwargs – Additional connection parameters.
Example
>>> # Use higher timeouts for role assumption scenarios >>> cursor = connection.cursor( ... AsyncArrowCursor, ... connect_timeout=10.0, ... request_timeout=30.0 ... )
- static get_default_converter(unload: bool = False) DefaultArrowTypeConverter | DefaultArrowUnloadTypeConverter | Any[source]¶
Get the default type converter for this cursor class.
- Parameters:
unload – Whether the converter is for UNLOAD operations. Some cursor types may return different converters for UNLOAD operations.
- Returns:
The default type converter instance for this cursor type.
- LIST_DATABASES_MAX_RESULTS = 50¶
- LIST_QUERY_EXECUTIONS_MAX_RESULTS = 50¶
- LIST_TABLE_METADATA_MAX_RESULTS = 50¶
- cancel(query_id: str) Future[None]¶
Cancel a running query asynchronously.
Submits a cancellation request for the specified query. The cancellation itself runs asynchronously in the background.
- Parameters:
query_id – The Athena query execution ID to cancel.
- Returns:
Future object that completes when the cancellation request finishes.
Example
>>> query_id, future = cursor.execute("SELECT * FROM huge_table") >>> # Later, cancel the query >>> cancel_future = cursor.cancel(query_id) >>> cancel_future.result() # Wait for cancellation to complete
- property connection: Connection[Any]¶
The connection that created this cursor.
- description(query_id: str) Future[list[tuple[str, str, None, None, int, int, str]] | None]¶
Get the column descriptions of a query’s result set asynchronously.
The future waits for the query to finish before it reads the result set.
- Parameters:
query_id – The Athena query execution ID.
- Returns:
Future object containing the DB API 2.0 column descriptions, or None.
- execute(operation: str, parameters: dict[str, Any] | list[str] | None = None, work_group: str | None = None, s3_staging_dir: str | None = None, cache_size: int | None = None, cache_expiration_time: int | None = None, result_reuse_enable: bool | None = None, result_reuse_minutes: int | None = None, paramstyle: str | None = None, result_set_type_hints: dict[str | int, str] | None = None, *, options: ExecuteOptions | None = None, **kwargs) tuple[str, Future[AthenaArrowResultSet | Any]][source]¶
Execute a SQL query asynchronously and return results as Arrow Tables.
- Parameters:
operation – SQL query string to execute.
parameters – Query parameters for parameterized queries.
work_group – Athena workgroup to use for this query.
s3_staging_dir – S3 location for query results.
cache_size – Number of queries to check for result caching.
cache_expiration_time – Cache expiration time in seconds.
result_reuse_enable – Enable Athena result reuse for this query.
result_reuse_minutes – Minutes to reuse cached results.
paramstyle – Parameter style (‘qmark’ or ‘pyformat’).
result_set_type_hints – Athena type signatures for complex-type columns, keyed by column name (case-insensitive) or zero-based column index.
options – Shared execution options as an
ExecuteOptionsinstance. Individual keyword arguments take precedence overoptionsfields.**kwargs – Additional execution parameters.
- Returns:
Tuple of (query_id, future) where future resolves to AthenaArrowResultSet.
- executemany(operation: str, seq_of_parameters: list[dict[str, Any] | list[str] | None], **kwargs) None¶
Execute multiple queries asynchronously (not supported).
This method is not supported for asynchronous cursors because managing multiple concurrent queries would be complex and resource-intensive.
- Parameters:
operation – SQL query string.
seq_of_parameters – Sequence of parameter sets.
**kwargs – Additional arguments.
- Raises:
NotSupportedError – Always raised as this operation is not supported.
Note
For bulk operations, consider using execute() with parameterized queries or batch processing patterns instead.
- get_table_metadata(table_name: str, catalog_name: str | None = None, schema_name: str | None = None, logging_: bool = True) AthenaTableMetadata¶
Get one table’s metadata.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
table_name – The table name.
catalog_name – The catalog, or None for the cursor’s catalog.
schema_name – The database, or None for the cursor’s schema.
logging – Whether to log a failed request.
- Returns:
The table’s metadata.
- Raises:
OperationalError – If the request fails, including when the table does not exist.
- list_databases(catalog_name: str | None, max_results: int | None = None) list[AthenaDatabase]¶
List the catalog’s databases.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
catalog_name – The catalog, or None for the cursor’s catalog.
max_results – The page size of each request.
- Returns:
The catalog’s databases.
- Raises:
OperationalError – If the request fails.
- list_table_metadata(catalog_name: str | None = None, schema_name: str | None = None, expression: str | None = None, max_results: int | None = None, logging_: bool = True) list[AthenaTableMetadata]¶
List a database’s table metadata.
In
AwsDataCatalogand S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; seeglue_metadata_fallback.- Parameters:
catalog_name – The catalog, or None for the cursor’s catalog.
schema_name – The database, or None for the cursor’s schema.
expression – A table name pattern.
max_results – The page size of each request.
logging – Whether to log a failed request.
- Returns:
The metadata of the database’s tables.
- Raises:
OperationalError – If the request fails.
- poll(query_id: str) Future[AthenaQueryExecution]¶
Poll for query completion asynchronously.
Waits for the query to complete (succeed, fail, or be cancelled) and returns the final execution status. This method blocks until completion but runs the polling in a background thread.
- Parameters:
query_id – The Athena query execution ID to poll.
- Returns:
Future object containing the final AthenaQueryExecution status.
Note
This method performs polling internally, so it will take time proportional to your query execution duration.
- query_execution(query_id: str) Future[AthenaQueryExecution]¶
Get query execution details asynchronously.
Retrieves the current execution status and metadata for a query. This is useful for monitoring query progress without blocking.
- Parameters:
query_id – The Athena query execution ID.
- Returns:
Future object containing AthenaQueryExecution with query details.
- setinputsizes(sizes)¶
Accept input sizes as DB API 2.0 requires, and ignore them.
- Parameters:
sizes – Sequence of parameter types or sizes.
- setoutputsize(size, column=None)¶
Accept a column buffer size as DB API 2.0 requires, and ignore it.
- Parameters:
size – Buffer size for large columns.
column – Index of the column the size applies to, or None for all large columns.
Arrow Result Set¶
- class pyathena.arrow.result_set.AthenaArrowResultSet(connection: Connection[Any], converter: Converter, query_execution: AthenaQueryExecution, arraysize: int, retry_config: RetryConfig, block_size: int | None = None, unload: bool = False, unload_location: str | None = None, connect_timeout: float | None = None, request_timeout: float | None = None, result_set_type_hints: dict[str | int, str] | None = None, **kwargs)[source]¶
Result set that provides Apache Arrow Table results with columnar optimization.
This result set handles CSV and Parquet result files from S3, converting them to Apache Arrow Tables which provide efficient columnar data processing and memory usage. It’s optimized for analytical workloads and large dataset operations.
- Features:
Efficient columnar data processing with Apache Arrow
Support for both CSV and Parquet result formats
Optimized memory usage for large datasets
Advanced timestamp parsing with multiple format support
Zero-copy operations where possible
- DEFAULT_BLOCK_SIZE¶
Default block size for Arrow operations (128MB).
Example
>>> # Used automatically by ArrowCursor >>> cursor = connection.cursor(ArrowCursor) >>> cursor.execute("SELECT * FROM large_table") >>> >>> # Get Arrow Table >>> table = cursor.as_arrow() >>> >>> # Convert to pandas if needed >>> df = table.to_pandas() >>> >>> # Or work with Arrow directly >>> print(f"Table has {table.num_rows} rows and {table.num_columns} columns")
Note
This class is used internally by ArrowCursor and typically not instantiated directly by users. Requires pyarrow to be installed.
- DEFAULT_BLOCK_SIZE = 134217728¶
- __init__(connection: Connection[Any], converter: Converter, query_execution: AthenaQueryExecution, arraysize: int, retry_config: RetryConfig, block_size: int | None = None, unload: bool = False, unload_location: str | None = None, connect_timeout: float | None = None, request_timeout: float | None = None, result_set_type_hints: dict[str | int, str] | None = None, **kwargs) None[source]¶
Initialize the result set and load the query results into an Arrow Table.
- Parameters:
connection – The connection that ran the query.
converter – The converter for result values.
query_execution – The query execution whose results to read.
arraysize – The default
fetchmany()size and the maximum number of rows per record batch that the fetch methods read from the table.retry_config – The retry configuration for API calls.
block_size – The block size in bytes for reading CSV results. If not set,
DEFAULT_BLOCK_SIZEis used.unload – Whether the query is an
UNLOADwhose Parquet output is read instead of the CSV results.unload_location – The S3 location of the
UNLOADoutput. If None, it is derived from the first file in the data manifest.connect_timeout – The connect timeout in seconds for the pyarrow S3 filesystem.
request_timeout – The request timeout in seconds for the pyarrow S3 filesystem.
result_set_type_hints – Athena type signatures for complex-type columns, keyed by column name (case-insensitive) or zero-based column index.
**kwargs – Additional keyword arguments, stored but not used.
- Raises:
ProgrammingError – If
query_executionis not given.OperationalError – If reading the query results fails.
- property timestamp_parsers: list[str]¶
The timestamp formats for reading CSV results, starting with pyarrow’s
ISO8601.
- property column_types: dict[str, type[Any]]¶
The converter’s types for the result columns it maps, keyed by column name.
- property converters: dict[str, Callable[[str | None], Any | None]]¶
The conversion functions for the result columns, keyed by column name.
- fetchone() tuple[Any | None, ...] | dict[Any, Any | None] | None[source]¶
Fetch the next row of the result.
- as_arrow() Table[source]¶
Return the query results as an Apache Arrow Table.
- Returns:
The Arrow Table that holds the query results.
- DEFAULT_RESULT_REUSE_MINUTES = 60¶
- as_polars() pl.DataFrame[source]¶
Return query results as a Polars DataFrame.
Converts the Apache Arrow Table to a Polars DataFrame for interoperability with the Polars data processing library.
- Returns:
Polars DataFrame containing all query results.
- Raises:
ImportError – If polars is not installed.
Example
>>> cursor = connection.cursor(ArrowCursor) >>> cursor.execute("SELECT * FROM my_table") >>> df = cursor.as_polars() >>> # Use with Polars operations
- property connection: Connection[Any]¶
The connection of the result set; raises
ProgrammingErrorif closed.
- property data_manifest_location: str | None¶
The S3 location of the data manifest that lists the files the query wrote.
- property description: list[tuple[str, str, None, None, int, int, str]] | None¶
The DB API 2.0 column descriptions.
None without result metadata, or for
INSERT,UPDATE,DELETE, andMERGE.
- property encryption_option: str | None¶
The
EncryptionOptionof the query results, such asSSE_S3orSSE_KMS.
- property engine_execution_time_in_millis: int | None¶
The time in milliseconds that the query engine took to run the query.
- property error_category: int | None¶
1 for system, 2 for user, 3 for other.
- Type:
The
ErrorCategoryof the failure
- property expected_bucket_owner: str | None¶
The AWS account ID expected to own the S3 bucket of the query results.
- fetchall() list[tuple[Any | None, ...] | dict[Any, Any | None]]¶
Fetch all remaining rows of the query result.
- Returns:
The remaining rows.
- fetchmany(size: int | None = None) list[tuple[Any | None, ...] | dict[Any, Any | None]]¶
Fetch the next set of rows of the query result.
- Parameters:
size – Maximum number of rows to fetch. If None or not positive,
arraysizeis used.- Returns:
The rows, fewer than
sizewhen the result is exhausted.
- property is_unload: bool¶
Check if the query is an UNLOAD statement.
- Returns:
True if the query is an UNLOAD statement, False otherwise.
- property query_planning_time_in_millis: int | None¶
The time in milliseconds that Athena took to plan the query.
- property query_queue_time_in_millis: int | None¶
The time in milliseconds that the query waited in the queue.
- property result_reuse_enabled: bool | None¶
Whether reuse of previous query results by age is enabled for the query.
- property result_reuse_minutes: int | None¶
The maximum age in minutes of a previous query result that Athena can reuse.
- property reused_previous_result: bool | None¶
Whether Athena reused a previous query result instead of running the query.
- property s3_acl_option: str | None¶
The
S3AclOptionof the query results, such asBUCKET_OWNER_FULL_CONTROL.
- property service_processing_time_in_millis: int | None¶
The time in milliseconds that Athena took to publish the query results.
- property state_change_reason: str | None¶
The
StateChangeReasonthat gives further detail about the state.
Arrow Data Converters¶
- class pyathena.arrow.converter.DefaultArrowTypeConverter[source]¶
Optimized type converter for Apache Arrow Table results.
This converter is specifically designed for the ArrowCursor and provides optimized type conversion for Apache Arrow’s columnar data format. It converts Athena data types to Python types that are efficiently handled by Apache Arrow.
- The converter focuses on:
Converting date/time types to appropriate Python objects
Handling decimal and binary types for Arrow compatibility
Preserving JSON and complex types
Maintaining high performance for columnar operations
Example
>>> from pyathena.arrow.converter import DefaultArrowTypeConverter >>> converter = DefaultArrowTypeConverter() >>> >>> # Used automatically by ArrowCursor >>> cursor = connection.cursor(ArrowCursor) >>> # converter is applied automatically to results
Note
This converter is used by default in ArrowCursor. Most users don’t need to instantiate it directly.
- __init__() None[source]¶
Initialize the converter with the default Arrow conversion functions and types.
- convert(type_: str, value: str | None, type_hint: str | None = None) Any | None[source]¶
Convert a value returned by Athena to a Python object.
- Parameters:
type – The Athena data type name.
value – The string value to convert, or None.
type_hint – Optional Athena DDL type signature of the value.
- Returns:
The converted value.
- class pyathena.arrow.converter.DefaultArrowUnloadTypeConverter[source]¶
Type converter for Arrow UNLOAD operations.
This converter is designed for use with UNLOAD queries that write results directly to Parquet files in S3. Since UNLOAD operations bypass the normal conversion process and write data in native Parquet format, this converter has minimal functionality.
Note
Used automatically when ArrowCursor is configured with unload=True. UNLOAD results are read directly as Arrow tables from Parquet files.
- convert(type_: str, value: str | None, type_hint: str | None = None) Any | None[source]¶
Convert a value returned by Athena to a Python object.
- Parameters:
type – The Athena data type name.
value – The string value to convert, or None.
type_hint – Optional Athena DDL type signature of the value.
- Returns:
The converted value.
Arrow Utilities¶
- pyathena.arrow.util.to_column_info(schema: Schema) tuple[dict[str, Any], ...][source]¶
Convert a PyArrow schema to Athena column information.
Iterates through all fields in the schema and converts each field’s type information to an Athena-compatible column metadata dictionary.
- Parameters:
schema – A PyArrow Schema object containing field definitions.
- Returns:
Name: The column name
Type: The Athena SQL type name
Precision: Numeric precision (0 for non-numeric types)
Scale: Numeric scale (0 for non-numeric types)
Nullable: Either “NULLABLE” or “NOT_NULL”
- Return type:
A tuple of dictionaries, each containing column metadata with keys
- pyathena.arrow.util.get_athena_type(type_: DataType) tuple[str, int, int][source]¶
Map a PyArrow data type to an Athena SQL type.
Converts PyArrow type identifiers to corresponding Athena SQL type names with appropriate precision and scale values. Handles all common Arrow types including numeric, string, binary, temporal, and complex types.
- Parameters:
type – A PyArrow DataType object to convert.
- Returns:
type_name: The Athena SQL type (e.g., “varchar”, “bigint”, “timestamp”)
precision: The numeric precision or max length
scale: The numeric scale (decimal places)
- Return type:
A tuple of (type_name, precision, scale) where
Note
Unknown types default to “string” with maximum varchar length. Decimal types preserve their original precision and scale.