Spark Integration¶
This section covers Apache Spark-specific cursors and execution management.
Spark Cursors¶
- class pyathena.spark.cursor.SparkCursor(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, terminate_session_on_close: bool | None = None, **kwargs)[source]¶
Cursor for executing PySpark code on Amazon Athena for Apache Spark.
This cursor allows you to execute PySpark code directly on Athena’s managed Spark environment. It’s designed for big data processing, ETL operations, and machine learning workloads that require Spark’s distributed computing capabilities.
The cursor manages Spark sessions automatically and provides an interface similar to other PyAthena cursors but optimized for Spark calculations rather than SQL queries.
- session_id¶
The Athena Spark session ID.
- description¶
The description of the current calculation.
- calculation_id¶
ID of the current calculation being executed.
Example
>>> from pyathena.spark.cursor import SparkCursor >>> cursor = connection.cursor(SparkCursor) >>> >>> # Execute PySpark code >>> spark_code = ''' ... df = spark.read.table("my_database.my_table") ... result = df.groupBy("category").count() ... result.show() ... ''' >>> cursor.execute(spark_code) >>> output = cursor.get_std_out()
# Configure Spark session >>> cursor = connection.cursor( … SparkCursor, … engine_configuration={ … ‘CoordinatorDpuSize’: 1, … ‘MaxConcurrentDpus’: 20, … ‘DefaultExecutorDpuSize’: 1 … } … )
Note
Requires an Athena workgroup configured for Spark calculations. Spark sessions have associated costs and idle timeout settings.
- property calculation_execution: AthenaCalculationExecution | None¶
The calculation execution that the other properties read, or None.
- get_std_out() str | None[source]¶
Get the standard output from the Spark calculation execution.
Retrieves and returns the contents of the standard output generated during the Spark calculation execution, if available.
- Returns:
The standard output as a string, or None if no output is available or the calculation has not been executed.
- get_std_error() str | None[source]¶
Get the standard error from the Spark calculation execution.
Retrieves and returns the contents of the standard error generated during the Spark calculation execution, if available. This is useful for debugging failed or problematic Spark operations.
- Returns:
The standard error as a string, or None if no error output is available or the calculation has not been executed.
- execute(operation: str, parameters: dict[str, Any] | list[str] | None = None, session_id: str | None = None, description: str | None = None, client_request_token: str | None = None, work_group: str | None = None, **kwargs) SparkCursor[source]¶
Execute a SQL query.
- Parameters:
operation – SQL query string.
parameters – Query parameters.
**kwargs – Execution options defined by the cursor implementation.
- cancel() None[source]¶
Stop the calculation that
execute()last started.- Raises:
ProgrammingError – If no calculation ID is set.
OperationalError – If the
StopCalculationExecutionrequest fails.
- LIST_DATABASES_MAX_RESULTS = 50¶
- LIST_QUERY_EXECUTIONS_MAX_RESULTS = 50¶
- LIST_TABLE_METADATA_MAX_RESULTS = 50¶
- __init__(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, terminate_session_on_close: bool | None = None, **kwargs) None¶
Initialize the cursor and start or attach to a Spark session.
If waiting for a newly started session fails, that session is terminated regardless of
terminate_session_on_close; a supplied session is not.- Parameters:
session_id – ID of an existing session to use. If omitted, a new session is started.
description – Description of a new session.
engine_configuration – Engine configuration of a new session. Defaults to
get_default_engine_configuration().notebook_version – Notebook version of a new session.
session_idle_timeout_minutes – Idle timeout of a new session in minutes.
terminate_session_on_close – Whether
close()terminates the session. If None, only a session started by this cursor is terminated; a session supplied withsession_idis left running.**kwargs – Arguments passed to
BaseCursor.
- Raises:
OperationalError – If the supplied session does not exist, or the session cannot be started or does not become idle.
- property calculation_id: str | None¶
The ID of the calculation tracked by this cursor, or None if there is none.
- close() None¶
Close the cursor, terminating its Spark session if configured to.
See
terminate_session_on_close. After a successful termination, further calls do not terminate the session again; after a failed one, calling this method again retries it.- Raises:
OperationalError – If terminating the session fails.
- property completion_date_time: datetime | None¶
The
CompletionDateTimeof the calculation, or None if there is none.
- property connection: Connection[Any]¶
The connection that created this cursor.
- property dpu_execution_in_millis: int | None¶
The
DpuExecutionInMillisstatistic of the calculation, or None if there is none.
- executemany(operation: str, seq_of_parameters: list[dict[str, Any] | list[str] | None], **kwargs) None¶
Execute a SQL query once for each set of parameters.
- Parameters:
operation – SQL query string.
seq_of_parameters – Sequence of parameter sets.
**kwargs – Execution options defined by the cursor implementation.
- static get_default_converter(unload: bool = False) DefaultTypeConverter | Any¶
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.
- static get_default_engine_configuration() dict[str, Any]¶
Return the engine configuration used when none is given.
- Returns:
a coordinator DPU size of 1, at most 2 concurrent DPUs, and a default executor DPU size of 1.
- Return type:
The
EngineConfigurationof a new session
- 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.
- 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
StateChangeReasonof the calculation, or None if there is none.
- property std_error_s3_uri: str | None¶
The
StdErrorS3Uriof the calculation, or None if there is none.
- class pyathena.spark.async_cursor.AsyncSparkCursor(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, max_workers: int = 20, terminate_session_on_close: bool | None = None, **kwargs)[source]¶
Asynchronous cursor for executing PySpark code on Amazon Athena for Apache Spark.
This cursor provides asynchronous execution of PySpark code on Athena’s managed Spark environment. It’s designed for non-blocking big data processing, ETL operations, and machine learning workloads that require Spark’s distributed computing capabilities without blocking the main thread.
- Features:
Asynchronous PySpark code execution with concurrent futures
Non-blocking query submission and result polling
Managed Spark sessions with configurable resources
Access to standard output and error streams asynchronously
Automatic session lifecycle management
Thread pool executor for concurrent operations
- session_id¶
The Athena Spark session ID.
Example
>>> from pyathena.spark.async_cursor import AsyncSparkCursor >>> >>> cursor = connection.cursor( ... AsyncSparkCursor, ... engine_configuration={ ... 'CoordinatorDpuSize': 1, ... 'MaxConcurrentDpus': 20 ... } ... ) >>> >>> # Execute PySpark code asynchronously >>> spark_code = ''' ... df = spark.read.table("my_database.my_table") ... result = df.groupBy("category").count() ... result.show() ... ''' >>> calculation_id, future = cursor.execute(spark_code) >>> >>> # Get result when ready >>> calc_execution = future.result() >>> stdout_future = cursor.get_std_out(calc_execution) >>> if stdout_future: ... output = stdout_future.result() ... print(output)
Note
Requires an Athena workgroup configured for Spark calculations. Spark sessions have associated costs and idle timeout settings. The cursor manages a thread pool for asynchronous operations.
- __init__(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, max_workers: int = 20, terminate_session_on_close: bool | None = None, **kwargs)[source]¶
Initialize the cursor and start or attach to a Spark session.
- Parameters:
session_id – ID of an existing session to use. If omitted, a new session is started.
description – Description of a new session.
engine_configuration – Engine configuration of a new session.
notebook_version – Notebook version of a new session.
session_idle_timeout_minutes – Idle timeout of a new session in minutes.
max_workers – Maximum number of threads for asynchronous operations.
terminate_session_on_close – Whether
close()terminates the session. If None, only a session started by this cursor is terminated; a session supplied withsession_idis left running.**kwargs – Arguments passed to
SparkBaseCursor.
- Raises:
ValueError – If
max_workersis not greater than 0.OperationalError – If the supplied session does not exist, or the session cannot be started or does not become idle.
- close(wait: bool = False) None[source]¶
Close the cursor, then shut down the executor.
The session is terminated as described in
SparkBaseCursor.close(). The executor is shut down even if terminating the session fails.- Parameters:
wait – Whether to wait for submitted futures to finish before returning or raising.
- Raises:
OperationalError – If terminating the session fails.
- calculation_execution(query_id: str) Future[AthenaCalculationExecution][source]¶
Get calculation execution details asynchronously.
- Parameters:
query_id – The calculation execution ID.
- Returns:
Future object containing the
AthenaCalculationExecution.
- get_std_out(calculation_execution: AthenaCalculationExecution) Future[str] | None[source]¶
Read the standard output of a calculation from S3 asynchronously.
- Parameters:
calculation_execution – The calculation execution whose
std_out_s3_uriis read.- Returns:
Future object containing the output text with leading and trailing whitespace removed, or None if the calculation has no
std_out_s3_uri.
- get_std_error(calculation_execution: AthenaCalculationExecution) Future[str] | None[source]¶
Read the standard error output of a calculation from S3 asynchronously.
- Parameters:
calculation_execution – The calculation execution whose
std_error_s3_uriis read.- Returns:
Future object containing the error output text with leading and trailing whitespace removed, or None if the calculation has no
std_error_s3_uri.
- poll(query_id: str) Future[AthenaCalculationExecution][source]¶
Wait for a calculation to reach a terminal state asynchronously.
- Parameters:
query_id – The calculation execution ID.
- Returns:
Future object containing the calculation execution in a terminal state.
- execute(operation: str, parameters: dict[str, Any] | list[str] | None = None, session_id: str | None = None, description: str | None = None, client_request_token: str | None = None, work_group: str | None = None, **kwargs) tuple[str, Future[AthenaQueryExecution | AthenaCalculationExecution]][source]¶
Execute a SQL query.
- Parameters:
operation – SQL query string.
parameters – Query parameters.
**kwargs – Execution options defined by the cursor implementation.
- cancel(query_id: str) Future[None][source]¶
Stop a calculation execution asynchronously.
- Parameters:
query_id – The calculation execution ID.
- Returns:
Future object that completes when the
StopCalculationExecutionrequest has been sent.
- LIST_DATABASES_MAX_RESULTS = 50¶
- LIST_QUERY_EXECUTIONS_MAX_RESULTS = 50¶
- LIST_TABLE_METADATA_MAX_RESULTS = 50¶
- property calculation_id: str | None¶
The ID of the calculation tracked by this cursor, or None if there is none.
- property connection: Connection[Any]¶
The connection that created this cursor.
- executemany(operation: str, seq_of_parameters: list[dict[str, Any] | list[str] | None], **kwargs) None¶
Execute a SQL query once for each set of parameters.
- Parameters:
operation – SQL query string.
seq_of_parameters – Sequence of parameter sets.
**kwargs – Execution options defined by the cursor implementation.
- static get_default_converter(unload: bool = False) DefaultTypeConverter | Any¶
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.
- static get_default_engine_configuration() dict[str, Any]¶
Return the engine configuration used when none is given.
- Returns:
a coordinator DPU size of 1, at most 2 concurrent DPUs, and a default executor DPU size of 1.
- Return type:
The
EngineConfigurationof a new session
- 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.
- 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.
Spark Base Classes¶
- class pyathena.spark.common.SparkBaseCursor(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, terminate_session_on_close: bool | None = None, **kwargs)[source]¶
Abstract base class for Spark-enabled cursor implementations.
This class provides the foundational functionality for executing PySpark code on Amazon Athena for Apache Spark. It manages Spark sessions, handles calculation execution lifecycle, and provides utilities for reading results from S3.
- Features:
Automatic Spark session management and lifecycle
Configurable engine resources (DPU allocation)
Session idle timeout and automatic cleanup
Standard output and error stream access via S3
Calculation execution status monitoring
Session validation and error handling
- session_id¶
The Athena Spark session identifier.
- calculation_id¶
ID of the current calculation being executed.
Note
This is an abstract base class used by concrete Spark cursor implementations like SparkCursor and AsyncSparkCursor. It should not be instantiated directly.
- __init__(session_id: str | None = None, description: str | None = None, engine_configuration: dict[str, Any] | None = None, notebook_version: str | None = None, session_idle_timeout_minutes: int | None = None, terminate_session_on_close: bool | None = None, **kwargs) None[source]¶
Initialize the cursor and start or attach to a Spark session.
If waiting for a newly started session fails, that session is terminated regardless of
terminate_session_on_close; a supplied session is not.- Parameters:
session_id – ID of an existing session to use. If omitted, a new session is started.
description – Description of a new session.
engine_configuration – Engine configuration of a new session. Defaults to
get_default_engine_configuration().notebook_version – Notebook version of a new session.
session_idle_timeout_minutes – Idle timeout of a new session in minutes.
terminate_session_on_close – Whether
close()terminates the session. If None, only a session started by this cursor is terminated; a session supplied withsession_idis left running.**kwargs – Arguments passed to
BaseCursor.
- Raises:
OperationalError – If the supplied session does not exist, or the session cannot be started or does not become idle.
- property calculation_id: str | None¶
The ID of the calculation tracked by this cursor, or None if there is none.
- static get_default_engine_configuration() dict[str, Any][source]¶
Return the engine configuration used when none is given.
- Returns:
a coordinator DPU size of 1, at most 2 concurrent DPUs, and a default executor DPU size of 1.
- Return type:
The
EngineConfigurationof a new session
- close() None[source]¶
Close the cursor, terminating its Spark session if configured to.
See
terminate_session_on_close. After a successful termination, further calls do not terminate the session again; after a failed one, calling this method again retries it.- Raises:
OperationalError – If terminating the session fails.
- executemany(operation: str, seq_of_parameters: list[dict[str, Any] | list[str] | None], **kwargs) None[source]¶
Execute a SQL query once for each set of parameters.
- Parameters:
operation – SQL query string.
seq_of_parameters – Sequence of parameter sets.
**kwargs – Execution options defined by the cursor implementation.
- LIST_DATABASES_MAX_RESULTS = 50¶
- LIST_QUERY_EXECUTIONS_MAX_RESULTS = 50¶
- LIST_TABLE_METADATA_MAX_RESULTS = 50¶
- property connection: Connection[Any]¶
The connection that created this cursor.
- abstractmethod execute(operation: str, parameters: dict[str, Any] | list[str] | None = None, **kwargs)¶
Execute a SQL query.
- Parameters:
operation – SQL query string.
parameters – Query parameters.
**kwargs – Execution options defined by the cursor implementation.
- static get_default_converter(unload: bool = False) DefaultTypeConverter | Any¶
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.
- 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.
- 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.
- class pyathena.spark.common.WithCalculationExecution[source]¶
Mixin class providing access to Spark calculation execution properties.
This mixin provides property accessors for calculation execution metadata and status information. It’s designed to be mixed with cursor classes that execute Spark calculations on Athena.
- Properties:
description: Human-readable description of the calculation
working_directory: S3 path where calculation files are stored
state: Current execution state (COMPLETED, FAILED, etc.)
state_change_reason: Explanation for state changes
submission_date_time: When the calculation was submitted
completion_date_time: When the calculation completed
dpu_execution_in_millis: DPU execution time in milliseconds
progress: Current execution progress information
std_out_s3_uri: S3 URI for standard output
std_error_s3_uri: S3 URI for standard error
result_s3_uri: S3 URI for calculation results
result_type: Type of result produced by the calculation
Note
This class requires that the implementing class provides calculation_execution, session_id, and calculation_id properties.
- abstract property calculation_execution: AthenaCalculationExecution | None¶
The calculation execution that the other properties read, or None.
- property working_directory: str | None¶
The
WorkingDirectoryof the calculation, or None if there is none.
- property state_change_reason: str | None¶
The
StateChangeReasonof the calculation, or None if there is none.
- property submission_date_time: datetime | None¶
The
SubmissionDateTimeof the calculation, or None if there is none.
- property completion_date_time: datetime | None¶
The
CompletionDateTimeof the calculation, or None if there is none.
- property dpu_execution_in_millis: int | None¶
The
DpuExecutionInMillisstatistic of the calculation, or None if there is none.