Zacznij terazZacznij za darmo

Implementacja funkcji zwrotnej

Dostałeś ostatnio zadanie dodania funkcji zwrotnych przy błędach do DAG-ów tworzonych przez twój zespół. Na początek chcesz dodać prostą funkcję zwrotną, która zapisuje komunikat do dziennika audytu, gdy sales_etl_dag zakończy się niepowodzeniem.

Obiekt dag i task są już zaimportowane, a zadania get_sales_data i process_sales_data zostały utworzone.

To ćwiczenie jest częścią kursu

Wprowadzenie do Apache Airflow w Pythonie

Zobacz kurs

Instrukcje do ćwiczenia

  • Utwórz funkcję zwrotną o nazwie alert_on_failure.
  • Zdefiniuj tę funkcję tak, aby przyjmowała dowolne obiekty przekazywane przez Airflow.
  • Wskaż funkcję alert_on_failure jako funkcję zwrotną wywoływaną przy błędzie.

Interaktywne ćwiczenie praktyczne

Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.

# 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()
Edytuj i uruchom kod