airflow.providers.amazon.aws.triggers.emr

Classes

EmrAddStepsTrigger

Poll for the status of EMR steps until they reach terminal state.

EmrCreateJobFlowTrigger

Asynchronously poll the boto3 API and wait for the JobFlow to finish executing.

EmrTerminateJobFlowTrigger

Asynchronously poll the boto3 API and wait for the JobFlow to finish terminating.

EmrContainerTrigger

Poll for the status of EMR container until reaches terminal state.

EmrStepSensorTrigger

Poll for the status of EMR container until reaches terminal state.

EmrServerlessCreateApplicationTrigger

Poll an Emr Serverless application and wait for it to be created.

EmrServerlessStartApplicationTrigger

Poll an Emr Serverless application and wait for it to be started.

EmrServerlessStopApplicationTrigger

Poll an Emr Serverless application and wait for it to be stopped.

EmrServerlessJobSensorTrigger

Poll an EMR Serverless job run until it reaches a target or failure state.

EmrServerlessStartJobTrigger

Poll an Emr Serverless job run and wait for it to be completed.

EmrServerlessDeleteApplicationTrigger

Poll an Emr Serverless application and wait for it to be deleted.

EmrServerlessCancelJobsTrigger

Trigger for canceling a list of jobs in an EMR Serverless application.

Module Contents

class airflow.providers.amazon.aws.triggers.emr.EmrAddStepsTrigger(job_flow_id, step_ids, waiter_delay, waiter_max_attempts, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll for the status of EMR steps until they reach terminal state.

Parameters:
aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrCreateJobFlowTrigger(job_flow_id, aws_conn_id=None, waiter_delay=30, waiter_max_attempts=60, waiter_name='job_flow_waiting', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Asynchronously poll the boto3 API and wait for the JobFlow to finish executing.

Parameters:
aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrTerminateJobFlowTrigger(job_flow_id, aws_conn_id=None, waiter_delay=30, waiter_max_attempts=60, region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Asynchronously poll the boto3 API and wait for the JobFlow to finish terminating.

Parameters:
aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrContainerTrigger(virtual_cluster_id, job_id, aws_conn_id='aws_default', waiter_delay=30, waiter_max_attempts=sys.maxsize, cancel_on_kill=True, region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll for the status of EMR container until reaches terminal state.

Parameters:
  • virtual_cluster_id (str) – Reference Emr cluster id

  • job_id (str) – job_id to check the state

  • aws_conn_id (str | None) – Reference to AWS connection id

  • waiter_delay (int) – polling period in seconds to check for the status

  • waiter_max_attempts (int) – The maximum number of attempts to be made. Defaults to an infinite wait.

  • cancel_on_kill (bool) – If True (default), cancel the EMR container job when the user marks the deferred task failed, clears it, or mark-succeeds it. Requires apache-airflow with BaseTrigger.on_kill() support; on older versions the hook is silently inert.

  • region_name (str | None) – The AWS region where the resources to watch are.

  • verify (bool | str | None) – Whether or not to verify SSL certificates. See: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/core/session.html

  • botocore_config (dict | None) – Configuration dictionary (key-values) for botocore client. See: https://botocore.amazonaws.com/v1/documentation/api/latest/reference/config.html

aws_hook_class[source]
virtual_cluster_id[source]
job_id[source]
cancel_on_kill = True[source]
async on_kill()[source]

Cancel the EMR container job when the user acts on the deferred task.

class airflow.providers.amazon.aws.triggers.emr.EmrStepSensorTrigger(job_flow_id, step_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll for the status of EMR container until reaches terminal state.

Parameters:
aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrServerlessCreateApplicationTrigger(application_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an Emr Serverless application and wait for it to be created.

Parameters:
Waiter_delay:

polling period in seconds to check for the status

aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrServerlessStartApplicationTrigger(application_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an Emr Serverless application and wait for it to be started.

Parameters:
Waiter_delay:

polling period in seconds to check for the status

aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrServerlessStopApplicationTrigger(application_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an Emr Serverless application and wait for it to be stopped.

Parameters:
Waiter_delay:

polling period in seconds to check for the status

aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrServerlessJobSensorTrigger(application_id, job_run_id, target_states, waiter_delay=60, waiter_max_attempts=sys.maxsize, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an EMR Serverless job run until it reaches a target or failure state.

Parameters:
  • application_id (str) – The ID of the application the job is running on.

  • job_run_id (str) – The ID of the job run.

  • target_states (set[str] | frozenset[str]) – The states that indicate the sensor has succeeded.

  • waiter_delay (int) – The time in seconds to wait between polling attempts.

  • waiter_max_attempts (int) – The maximum number of attempts to be made. Defaults to an infinite wait.

  • aws_conn_id (str | None) – Reference to the AWS connection ID.

  • region_name (str | None) – The AWS region where the job is running.

  • verify (bool | str | None) – Whether to verify SSL certificates.

  • botocore_config (dict | None) – Configuration dictionary for the botocore client.

hook()[source]

Build the hook this trigger waits with.

class airflow.providers.amazon.aws.triggers.emr.EmrServerlessStartJobTrigger(application_id, job_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', cancel_on_kill=True, region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an Emr Serverless job run and wait for it to be completed.

Parameters:
aws_hook_class[source]
application_id[source]
job_id[source]
cancel_on_kill = True[source]
get_task_instance(*, session)[source]

Get the task instance for the current trigger (Airflow 2.x compatibility).

async get_task_state()[source]

Get the current state of the task instance (Airflow 3.x).

async safe_to_cancel()[source]

Whether it is safe to cancel the EMR Serverless job.

Returns True if task is NOT DEFERRED (user-initiated cancellation). Returns False if task is DEFERRED (triggerer restart - don’t cancel job).

async run()[source]

Run the trigger and wait for the job to complete.

If the task is cancelled while waiting, attempt to cancel the EMR Serverless job if cancel_on_kill is enabled and it’s safe to do so.

async on_kill()[source]

Cancel the EMR Serverless job when the trigger is cancelled by a user action.

This hook is available in Airflow 3.3+ via BaseTrigger.on_kill(). For older Airflow versions, the CancelledError handler in run() provides the same cancellation behavior.

class airflow.providers.amazon.aws.triggers.emr.EmrServerlessDeleteApplicationTrigger(application_id, waiter_delay=30, waiter_max_attempts=60, aws_conn_id='aws_default', region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Poll an Emr Serverless application and wait for it to be deleted.

Parameters:
Waiter_delay:

polling period in seconds to check for the status

aws_hook_class[source]
class airflow.providers.amazon.aws.triggers.emr.EmrServerlessCancelJobsTrigger(application_id, aws_conn_id, waiter_delay, waiter_max_attempts, region_name=None, verify=None, botocore_config=None)[source]

Bases: airflow.providers.amazon.aws.triggers.base.AwsBaseWaiterTrigger

Trigger for canceling a list of jobs in an EMR Serverless application.

Parameters:
aws_hook_class[source]
property hook_instance: airflow.providers.amazon.aws.hooks.base_aws.AwsGenericHook[source]

This property is added for backward compatibility.

Was this entry helpful?