ÎncepețiÎncepe gratuit

Implementarea unei funcții de callback

Ai primit recent sarcina de a adăuga callback-uri de eșec la DAG-urile create de echipa ta. Pentru început, vrei să adaugi un callback simplu care scrie un mesaj în jurnalul de audit atunci când sales_etl_dag eșuează.

Obiectele dag și task sunt deja importate, iar taskurile get_sales_data și process_sales_data au fost deja create.

Acest exercițiu face parte din cursul

Introducere în Apache Airflow în Python

Vezi cursul

Instrucțiuni pentru exercițiu

  • Creează o funcție de callback numită alert_on_failure.
  • Definește funcția astfel încât să accepte orice obiecte transmise de Airflow.
  • Specifică un callback de eșec folosind funcția alert_on_failure.

Exercițiu interactiv practic

Încearcă acest exercițiu completând acest cod de exemplu.

# 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()
Editează și rulează codul