Python 中的 Apache Airflow 入門
Mike Metzger
Data Engineer

on_failure_callback-Dag 執行/任務失敗on_success_callback-Dag 執行/任務成功on_retry_callbackon_skipped_callbackon_execute_callback
context 字典context 包含該次 Dag 執行/任務的資訊context["dag"].dag_id-Dag 名稱context["task_instance"].task_id-任務名稱context["logical_date"]-Dag 執行的日期contextdef alert_on_failure(context):dag_id = context["dag"].dag_id task_id = context["task_instance"].task_id print(f"Task {task_id} in DAG {dag_id} has failed.")
def alert_on_failure(context): dag_id = context["dag"].dag_id task_id = context["task_instance"].task_id print(f"Task {task_id} in Dag {dag_id} has failed.")@dag(on_failure_callback=alert_on_failure) def sales_etl_dag(): @task() def data_import_task(): raise ValueError("Simulated failure")# Task data_import_task in Dag sales_etl_dag has failed.
SmtpNotifier-寄送電子郵件警示SlackNotifier-發送訊息到 Slack 頻道
airflow.providers.smtp.notifications.smtp 函式庫from_email 與 to 屬性subject、html_content 等@dag(dag_id=`sales_etl_dag`,
on_failure_callback=SmtpNotifier(
to='[email protected]',
from_email='[email protected]',
subject='Dag sales_etl_dag has failed'
)
)

Python 中的 Apache Airflow 入門