Triggers และ Fault Tolerance

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

Mike Metzger

Data Engineer

สร้าง pipeline ที่มีความทนทาน

  • Triggers - ควบคุมว่าเมื่อใด task จะสามารถทำงานได้
    • กำหนด task ที่รันภายใต้เงื่อนไขเฉพาะ
    • Trigger เริ่มต้น - all success

ภาพประกอบแสดง triggers ที่ควบคุมการทำงานของ task

  • Fault tolerance - ช่วยให้ task กู้คืนได้เมื่อเกิดข้อผิดพลาด
    • ลองรัน task ที่ล้มเหลวใหม่ได้

ภาพประกอบแสดง fault tolerance ที่รีทรี task ที่ล้มเหลว

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

Trigger rules

  • ตรวจสอบสถานะของ task ก่อนหน้าก่อนที่ task จะเริ่มทำงาน
  • กำหนดได้ว่า task ควรดำเนินต่อหรือไม่
  • ค่าเริ่มต้น: task ก่อนหน้าทุก task ต้องสำเร็จก่อนจึงจะดำเนินต่อ
  • Trigger rules ช่วยให้เปลี่ยนพฤติกรรมนี้ได้

ภาพประกอบแสดง trigger rules ที่ตรวจสอบสถานะของ task ต้นน้ำ

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

การใช้งาน trigger rules

  • Task แจ้งเตือน
  • Task ล้างข้อมูล
  • การรันแบบมีเงื่อนไข

ภาพประกอบแสดงกรณีใช้งาน trigger rules เช่น task แจ้งเตือนและล้างข้อมูล

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

Trigger rules ที่ใช้บ่อย

  • all_success - ทุก task ก่อนหน้าสำเร็จ ไม่มีการข้าม
  • all_failed - ทุก task ก่อนหน้าล้มเหลว
  • all_done - ทุก task ก่อนหน้าเสร็จสิ้น ไม่ว่าผลจะเป็นอย่างไร
  • one_failed - มีอย่างน้อยหนึ่ง task ที่ล้มเหลว
  • one_success - มีอย่างน้อยหนึ่ง task ที่สำเร็จ
  • none_failed - ทุก task ก่อนหน้าสำเร็จหรือถูกข้าม
Apache Airflow เบื้องต้นด้วย Python

การใช้งาน trigger rule

  • from airflow.utils.trigger_rule import TriggerRule
  • ใช้ trigger rules กับ decorator @task
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Apache Airflow เบื้องต้นด้วย Python

แอตทริบิวต์ fault tolerance

  • retries - พารามิเตอร์ของ task ที่กำหนดจำนวนครั้งที่ Airflow จะลองรัน task ที่ล้มเหลวใหม่ก่อนจะมาร์กว่าล้มเหลว
  • retry_delay - พารามิเตอร์ของ task ที่รับค่า timedelta เพื่อกำหนดเวลารอระหว่างการลองใหม่แต่ละครั้ง
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Apache Airflow เบื้องต้นด้วย Python

TriggerDagRunOperator

  • ให้ DAG หนึ่งเรียกใช้อีก DAG ได้
  • มีโค้ดและ DAG ให้พร้อม
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id - ต้องตรงกับ dag_id
  • wait_for_completion - บล็อก task ถัดไปจนกว่าจะเสร็จ
  • poke_interval - ความถี่ในการตรวจสอบ
  • conf - ข้อมูลที่ส่งไปยัง DAG ลูก
Apache Airflow เบื้องต้นด้วย Python

ตัวอย่าง TriggerDagRun

trigger_child = TriggerDagRunOperator(
        task_id="trigger_child_pipeline",
        trigger_dag_id="child_pipeline_dag",   
        wait_for_completion=True,              
        poke_interval=30,                      
        conf={                                 
            "source": "s3://my-bucket/raw/",
        },
    )

task1 >> trigger_child    # Run first task, then kick off child Dag
trigger_child >> cleanup  # Run a cleanup task after child Dag completes
Apache Airflow เบื้องต้นด้วย Python

มาฝึกกันเถอะ!

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

Preparing Video For Download...