觸發器與錯誤容忍

Python 中的 Apache Airflow 入門

Mike Metzger

Data Engineer

打造健壯的 pipeline

  • 觸發器(Triggers)-控制工作何時執行
    • 定義在特定條件下執行的工作
    • 預設觸發條件:全部成功

說明觸發器如何控制工作何時執行

  • 錯誤容忍(Fault tolerance)-允許在單一工作內復原
    • 重新嘗試個別工作

說明錯誤容忍如何重試失敗的工作

Python 中的 Apache Airflow 入門

觸發規則

  • 在工作開始前檢查前置工作的狀態
  • 讓你決定工作是否應繼續
  • 預設需所有前置工作成功後才繼續
  • 觸發規則可改變此行為

說明觸發規則如何檢查上游工作的狀態

Python 中的 Apache Airflow 入門

觸發規則的用途

  • 通知工作
  • 清理工作
  • 條件式執行

示意觸發規則的應用情境,如通知與清理工作

Python 中的 Apache Airflow 入門

關鍵觸發規則

  • all_success-所有前置工作皆成功,無略過
  • all_failed-所有前置工作皆失敗
  • all_done-所有前置工作皆完成,不論結果
  • one_failed-至少一個工作失敗
  • one_success-至少一個工作成功
  • none_failed-所有前置工作成功或被略過
Python 中的 Apache Airflow 入門

觸發規則實作

  • from airflow.utils.trigger_rule import TriggerRule
  • 將觸發規則套用在 @task 裝飾器
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Python 中的 Apache Airflow 入門

錯誤容忍屬性

  • retries-工作參數,定義 Airflow 在標記為失敗前要重試幾次
  • retry_delay-工作參數,接受 timedelta,定義每次重試間的等待時間
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Python 中的 Apache Airflow 入門

TriggerDagRunOperator

  • 允許一個 Dag 觸發另一個 Dag
  • 可傳遞程式碼與 Dags
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id-必須與 dag_id 相同
  • wait_for_completion-阻擋後續工作直到完成
  • poke_interval-檢查頻率
  • conf-傳給子 Dag 的資料
Python 中的 Apache Airflow 入門

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    # 先執行 task1,再觸發子 Dag
trigger_child >> cleanup  # 子 Dag 完成後執行清理工作
Python 中的 Apache Airflow 入門

一起來練習吧!

Python 中的 Apache Airflow 入門

Preparing Video For Download...