НачатьНачать бесплатно

Реализация функции обратного вызова

Вам недавно поручили добавить функции обратного вызова при сбоях в DAG'и, созданные вашей командой. Для начала вы хотите добавить простую функцию обратного вызова, которая записывает сообщение в журнал аудита при сбое sales_etl_dag.

Объекты dag и task уже импортированы, а задачи get_sales_data и process_sales_data созданы.

Это упражнение является частью курса

Введение в Apache Airflow на Python

Посмотреть курс

Инструкции к упражнению

  • Создайте функцию обратного вызова с именем alert_on_failure.
  • Определите функцию так, чтобы она принимала любые объекты, которые передаёт ей Airflow.
  • Укажите функцию alert_on_failure в качестве обратного вызова при сбое.

Интерактивное практическое упражнение

Попробуйте выполнить это упражнение, дополнив этот пример кода.

# 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()
Редактировать и запускать код