airflow.providers.amazon.aws.hooks.athena

This module contains AWS Athena hook.

Attributes

MULTI_LINE_QUERY_LOG_PREFIX

Classes

AthenaHook

Interact with Amazon Athena.

Functions

query_params_to_string(params)

Module Contents

airflow.providers.amazon.aws.hooks.athena.MULTI_LINE_QUERY_LOG_PREFIX = Multiline-String[source]
Show Value
"""
          """
airflow.providers.amazon.aws.hooks.athena.query_params_to_string(params)[source]
class airflow.providers.amazon.aws.hooks.athena.AthenaHook(*args, log_query=True, **kwargs)[source]

Bases: airflow.providers.amazon.aws.hooks.base_aws.AwsBaseHook

Interact with Amazon Athena.

Provide thick wrapper around boto3.client("athena").

Parameters:

log_query (bool) – Whether to log athena query and other execution params when it’s executed. Defaults to True.

Additional arguments (such as aws_conn_id) may be specified and are passed down to the underlying AwsBaseHook.

INTERMEDIATE_STATES = ('QUEUED', 'RUNNING')[source]
FAILURE_STATES = ('FAILED', 'CANCELLED')[source]
SUCCESS_STATES = ('SUCCEEDED',)[source]
TERMINAL_STATES = ('SUCCEEDED', 'FAILED', 'CANCELLED')[source]
SPARK_FAILURE_STATES = ('FAILED', 'CANCELED')[source]
SPARK_SUCCESS_STATES = ('COMPLETED',)[source]
SPARK_TERMINAL_STATES = ('COMPLETED', 'FAILED', 'CANCELED')[source]
log_query = True[source]
run_query(query, query_context, result_configuration, client_request_token=None, workgroup='primary')[source]

Run a Trino/Presto query on Athena with provided config.

Parameters:
  • query (str) – Trino/Presto query to run.

  • query_context (dict[str, str]) – Context in which query need to be run.

  • result_configuration (dict[str, Any]) – Dict with path to store results in and config related to encryption.

  • client_request_token (str | None) – Unique token created by user to avoid multiple executions of same query.

  • workgroup (str) – Athena workgroup name, when not specified, will be 'primary'.

Returns:

Submitted query execution ID.

Return type:

str

get_query_info(query_execution_id, use_cache=False)[source]

Get information about a single execution of a query.

Parameters:
  • query_execution_id (str) – Id of submitted athena query

  • use_cache (bool) – If True, use execution information cache

check_query_status(query_execution_id, use_cache=False)[source]

Fetch the state of a submitted query.

Parameters:

query_execution_id (str) – Id of submitted athena query

Returns:

One of valid query states, or None if the response is malformed.

Return type:

str | None

get_state_change_reason(query_execution_id, use_cache=False)[source]

Fetch the reason for a state change (e.g. error message). Returns None or reason string.

Parameters:

query_execution_id (str) – Id of submitted athena query

get_query_results(query_execution_id, next_token_id=None, max_results=1000)[source]

Fetch submitted query results.

Parameters:
  • query_execution_id (str) – Id of submitted athena query

  • next_token_id (str | None) – The token that specifies where to start pagination.

  • max_results (int) – The maximum number of results (rows) to return in this request.

Returns:

None if the query is in intermediate, failed, or cancelled state. Otherwise a dict of query outputs.

Return type:

dict | None

get_query_results_paginator(query_execution_id, max_items=None, page_size=None, starting_token=None)[source]

Fetch submitted Athena query results.

Parameters:
  • query_execution_id (str) – Id of submitted athena query

  • max_items (int | None) – The total number of items to return.

  • page_size (int | None) – The size of each page.

  • starting_token (str | None) – A token to specify where to start paginating.

Returns:

None if the query is in intermediate, failed, or cancelled state. Otherwise a paginator to iterate through pages of results.

Return type:

botocore.paginate.PageIterator | None

Call :meth`.build_full_result()` on the returned paginator to get all results at once.

poll_query_status(query_execution_id, max_polling_attempts=None, sleep_time=None)[source]

Poll the state of a submitted query until it reaches final state.

Parameters:
  • query_execution_id (str) – ID of submitted athena query

  • max_polling_attempts (int | None) – Number of times to poll for query state before function exits

  • sleep_time (int | None) – Time (in seconds) to wait between two consecutive query status checks.

Returns:

One of the final states

Return type:

str | None

get_output_location(query_execution_id)[source]

Get the output location of the query results in S3 URI format.

Parameters:

query_execution_id (str) – Id of submitted athena query

stop_query(query_execution_id)[source]

Cancel the submitted query.

Parameters:

query_execution_id (str) – Id of submitted athena query

start_spark_calculation(*, session_id, code_block, description=None, client_request_token=None)[source]

Start an Athena Spark calculation execution.

Parameters:
  • session_id (str) – The Athena session ID.

  • code_block (str) – Spark code to execute, typically notebook-like code.

  • description (str | None) – Optional description of the calculation. Defaults to None.

  • client_request_token (str | None) – Optional idempotency token. Defaults to None.

Returns:

CalculationExecutionId

Return type:

str

get_spark_calculation_info(calculation_execution_id, use_cache=False)[source]

Get information about a single Athena Spark calculation execution.

Parameters:
  • calculation_execution_id (str) – CalculationExecutionId returned by start_spark_calculation.

  • use_cache (bool) – If True, use execution information cache. Defaults to False.

Returns:

Calculation execution response.

Return type:

dict[str, Any]

check_spark_calculation_status(calculation_execution_id, use_cache=False)[source]

Fetch the state of a submitted Athena Spark calculation execution.

Parameters:
  • calculation_execution_id (str) – CalculationExecutionId returned by start_spark_calculation.

  • use_cache (bool) – If True, use execution information cache. Defaults to False.

Returns:

One of valid calculation states, or None if the response is malformed.

Return type:

str | None

get_spark_calculation_state_change_reason(calculation_execution_id, use_cache=False)[source]

Fetch the reason for an Athena Spark calculation state change, such as an error message.

Parameters:
  • calculation_execution_id (str) – CalculationExecutionId returned by start_spark_calculation.

  • use_cache (bool) – If True, use execution information cache. Defaults to False.

Returns:

State change reason string, or None.

Return type:

str | None

poll_spark_calculation_status(calculation_execution_id, waiter_delay=30, waiter_max_attempts=120)[source]

Poll an Athena Spark calculation until it reaches a terminal state.

Parameters:
  • calculation_execution_id (str) – ID of the submitted calculation.

  • waiter_delay (int) – Seconds to wait between status checks. Defaults to 30.

  • waiter_max_attempts (int) – Maximum number of status checks. Defaults to 120.

Returns:

The latest calculation state, or None if the response is malformed.

Return type:

str | None

stop_spark_calculation(calculation_execution_id)[source]

Cancel the submitted Athena Spark calculation execution.

Parameters:

calculation_execution_id (str) – CalculationExecutionId returned by start_spark_calculation.

Returns:

Response from stop_calculation_execution.

Return type:

dict[str, Any]

Was this entry helpful?