Začněte nyníZačněte zdarma

Implementace callback funkce

Nedávno ti bylo přiděleno doplnit failure callbacky do DAGů vytvořených tvým týmem. Pro začátek chceš přidat jednoduchý failure callback, který zapíše zprávu do audit logu vždy, když sales_etl_dag selže.

Objekty dag a task jsou již naimportovány a úlohy get_sales_data a process_sales_data jsou připraveny.

Toto cvičení je součástí kurzu

Úvod do Apache Airflow v Pythonu

Zobrazit kurz

Pokyny k cvičení

  • Vytvoř callback funkci s názvem alert_on_failure.
  • Definuj funkci tak, aby přijímala libovolné objekty, které jí Airflow předá.
  • Zadej failure callback pomocí funkce alert_on_failure.

Interaktivní cvičení na vyzkoušení si v praxi

Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.

# 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()
Upravit a spustit kód