Apache Airflow เบื้องต้นด้วย Python
Mike Metzger
Data Engineer

on_failure_callback - Dag run / task ล้มเหลวon_success_callback - Dag run / task สำเร็จon_retry_callbackon_skipped_callbackon_execute_callback
context ให้โดยอัตโนมัติcontext มีข้อมูลเกี่ยวกับ Dag run / taskcontext["dag"].dag_id - ชื่อของ Dagcontext["task_instance"].task_id - ชื่อของ taskcontext["logical_date"] - วันที่ของ Dag runcontextdef 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 channel
airflow.providers.smtp.notifications.smtpfrom_email และ tosubject, 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'
)
)

Apache Airflow เบื้องต้นด้วย Python