เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

การใช้งาน 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()
แก้ไขและรันโค้ด