การใช้งาน Trigger Rule
หลังจากสร้าง workflow แล้ว คุณพบว่า Dag จะมีประโยชน์มากขึ้นหากมีการแจ้งเตือนเมื่อมี task อย่างน้อยหนึ่งรายการล้มเหลว จึงตัดสินใจเพิ่ม task ที่ใช้การตรวจสอบแบบ one failed เพื่อแจ้งเตือนเมื่อมี task ใดใน Dag ล้มเหลว
ทุก task ที่เหลือถูกกำหนดไว้แล้ว และได้ import ออบเจ็กต์ task และ dag มาให้เรียบร้อยแล้ว
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
Apache Airflow เบื้องต้นด้วย Python
คำแนะนำการฝึกหัด
- Import ไลบรารีที่เหมาะสมสำหรับใช้งาน trigger rules
- เพิ่ม attribute ของ trigger rule ที่ถูกต้องให้กับ task
notify_on_failure - กำหนดค่า attribute ให้ task นี้ทำงานเมื่อมี upstream task อย่างน้อยหนึ่งรายการล้มเหลว
- กำหนดให้
notify_on_failureเป็น downstream dependency ของ transform task ทั้งสอง
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# Import TriggerRule
from airflow.utils.____ import ____
@dag(schedule="@daily", start_date=datetime(2026, 5, 1))
def etl_pipeline():
# Trigger notify_on_failure when any upstream task fails
@task(____=TriggerRule.____)
def notify_on_failure(**context) -> None:
dag_id = context["dag"].dag_id
run_id = context["run_id"]
print(f"ALERT: A task failed in DAG '{dag_id}', run '{run_id}'. Sending notification...")
# Set notify_on_failure downstream of both transform tasks
[transform_users(), transform_orders()] ____ notify_on_failure()
etl_pipeline()