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 入门