Implémenter une fonction de rappel
On vous a récemment confié d'ajouter des fonctions de rappel en cas d'échec aux Dags créés par votre équipe. Pour commencer, vous souhaitez ajouter une simple fonction de rappel qui écrit un message dans le journal d'audit lorsque le sales_etl_dag échoue.
Les objets dag et task sont déjà importés et les tâches get_sales_data et process_sales_data ont été créées.
Cette activité fait partie du cours
Introduction à Apache Airflow en Python
Instructions de l’exercice
- Créez une fonction de rappel nommée
alert_on_failure. - Définissez la fonction pour accepter tout objet qu'Airflow lui transmet.
- Spécifiez une fonction de rappel en cas d'échec en utilisant la fonction
alert_on_failure.
Exercice interactif pratique
Essayez cet exercice en complétant ce code d’exemple.
# Create the callback function
def ____(____):
dag_id = context["dag"].dag_id
task_id = context["task_instance"].task_id
print(f"Task {task_id} in Dag {dag_id} has failed.")
# Specify the Dag with a failure callback
@dag(dag_id='sales_etl_dag',
____=alert_on_failure
)
def sales_etl_dag():
get_sales_data() >> process_sales_data()
sales_etl_dag()