sla_miss_callback to send email on missing task SLA in Apache Airflow

Viewed 838

I have a DAG A that is being triggered by a parent DAG B. So DAG A doesn't have any schedule interval defined in it.

1.I would like to set up a sla_miss_callback on one of the task in DAG A.

2.I would like to get an e-mail notification whenever the task misses it's SLA.

I have tried methods available in google and stackoverflow. The e-mail is not getting triggered as expected.

Sharing the sample code I have used for testing.

from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import timedelta, datetime
import logging

def print_sla_miss(**kwargs):
    logging.info("SLA missed")


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2021, 1, 1),
    'email': 'sample@xxx.com',
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 0
}

with DAG('sla_test', schedule_interval=None, max_active_runs=1, catchup=False,sla_miss_callback=print_sla_miss, default_args=default_args) as dag:

    sleep = BashOperator(
        task_id='timeout',
        sla=timedelta(seconds=5),
        bash_command='sleep 15',
        retries=0,
        dag=dag,
    )

Thanks in advance.

2 Answers

SLAs will only be evaluated on scheduled DAG Runs. Since you have schedule_interval=None the SLA you set is not being evaluated for this DAG.

If there is a certain amount of time you expect the triggered DAG to finish, you could set that SLA in the sensor task in the parent DAG that checks when the child DAG is finished.

Another possible workaround is to set up a Slack notification for when the child DAG finishes entirely, or when a certain task starts/finishes so you can evaluate if it has been running for too long.

To achieve my requirement, I have created a seperate DAG that watches the task run status every 5 mins and notifies through e-mail based on the run status as below.To do this I am sending the execution date of my main DAG to an airflow variable.

#importing operators and modules
from airflow import DAG
from airflow.operators.python_operator import BranchPythonOperator
from airflow.operators.python_operator import PythonOperator
from airflow.operators.email_operator import EmailOperator
from airflow.api.common.experimental.get_task_instance import get_task_instance
from airflow.models import Variable
from datetime import datetime,timedelta,timezone
import dateutil

#setting default arguments
default_args = {
    'owner': 'test',
    'depends_on_past': False,
    'start_date': datetime(2021, 1, 1),
    'email': ['abc@example.com'],
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 0
}

#getting current status of task in main DAG
exec_date = dateutil.parser.parse(Variable.get('main_dag_execution_date'))
ti = get_task_instance('main_dag', 'task_to_check', exec_date)
state = ti.current_state()
start_date = ti.start_date
end_date = ti.end_date
print("start_date",start_date," end_date",end_date, " execution_date",exec_date)

#deciding the action based on status of the task
def check_task_status(**kwargs):

    if state == 'running' and datetime.now(timezone.utc) > start_date + timedelta(minutes = 10):
        breach_mail = 'breach_mail'
        return breach_mail
    elif state == 'failed':
        failure_mail = 'failure_mail'
        return failure_mail
    else:
        other_state = 'other_state'
        return other_state

#print statement when status is not in breached or failed state
def print_current_state(**context):
    if start_date is None:
        print("task is in wait state")
    else:
        print("task is in " + state + " state")


with DAG('sla_check', schedule_interval='0-59/5 9-23 * * *', max_active_runs=1, catchup=False,default_args=default_args) as dag:

    check_task_status = BranchPythonOperator(task_id='check_task_status', python_callable=check_task_status,
                                    provide_context=True,
                                    dag=dag)

    breach_mail = EmailOperator(task_id='breach_mail', to='Abc@example.com',
                                      subject='SLA for task breached',
                                      html_content="<p>Hi,<br><br>task running belyond SLA<br>", dag=dag)

    failure_mail = EmailOperator(task_id='failure_mail', to='Abc@example.com',
                                      subject='task failed',
                                      html_content="<p>Hi,<br><br>task failed. Please check.<br>", dag=dag)

    other_state = PythonOperator(task_id='other_state', python_callable=print_current_state,
                                    provide_context=True,
                                    dag=dag)

    check_task_status >> breach_mail
    check_task_status >> failure_mail
    check_task_status >> other_state
Related