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:
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 with session_id is 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 CompletionDateTime of the calculation, or None if there is none.

property connection: Connection[Any]

The connection that created this cursor.

property description: str | None

The Description of the calculation, or None if there is none.

property dpu_execution_in_millis: int | None

The DpuExecutionInMillis statistic 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 EngineConfiguration of 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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 progress: str | None

The Progress statistic of the calculation, or None if there is none.

property result_s3_uri: str | None

The ResultS3Uri of the calculation, or None if there is none.

property result_type: str | None

The ResultType of the calculation, or None if there is none.

property session_id: str

The ID of the Spark session that this cursor runs calculations in.

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: str | None

The State of the calculation, or None if there is none.

property state_change_reason: str | None

The StateChangeReason of the calculation, or None if there is none.

property std_error_s3_uri: str | None

The StdErrorS3Uri of the calculation, or None if there is none.

property std_out_s3_uri: str | None

The StdOutS3Uri of the calculation, or None if there is none.

property submission_date_time: datetime | None

The SubmissionDateTime of the calculation, or None if there is none.

property working_directory: str | None

The WorkingDirectory of 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 with session_id is left running.

  • **kwargs – Arguments passed to SparkBaseCursor.

Raises:
  • ValueError – If max_workers is 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_uri is 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_uri is 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 StopCalculationExecution request 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 EngineConfiguration of 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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 session_id: str

The ID of the Spark session that this cursor runs calculations in.

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 with session_id is 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 session_id: str

The ID of the Spark session that this cursor runs calculations in.

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 EngineConfiguration of 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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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 AwsDataCatalog and S3 Tables catalogs, a throttled request is answered from the AWS Glue Data Catalog; see glue_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.

__init__()[source]

Initialize the mixin, which keeps no state of its own.

abstract property calculation_execution: AthenaCalculationExecution | None

The calculation execution that the other properties read, or None.

abstract property session_id: str

The ID of the Spark session that runs the calculations.

abstract property calculation_id: str | None

The ID of the current calculation, or None.

property description: str | None

The Description of the calculation, or None if there is none.

property working_directory: str | None

The WorkingDirectory of the calculation, or None if there is none.

property state: str | None

The State of the calculation, or None if there is none.

property state_change_reason: str | None

The StateChangeReason of the calculation, or None if there is none.

property submission_date_time: datetime | None

The SubmissionDateTime of the calculation, or None if there is none.

property completion_date_time: datetime | None

The CompletionDateTime of the calculation, or None if there is none.

property dpu_execution_in_millis: int | None

The DpuExecutionInMillis statistic of the calculation, or None if there is none.

property progress: str | None

The Progress statistic of the calculation, or None if there is none.

property std_out_s3_uri: str | None

The StdOutS3Uri of the calculation, or None if there is none.

property std_error_s3_uri: str | None

The StdErrorS3Uri of the calculation, or None if there is none.

property result_s3_uri: str | None

The ResultS3Uri of the calculation, or None if there is none.

property result_type: str | None

The ResultType of the calculation, or None if there is none.