Python 中的 Apache Airflow 入門
Mike Metzger
Data Engineer




all_success-所有前置工作皆成功,無略過all_failed-所有前置工作皆失敗all_done-所有前置工作皆完成,不論結果one_failed-至少一個工作失敗one_success-至少一個工作成功none_failed-所有前置工作成功或被略過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!')
retries-工作參數,定義 Airflow 在標記為失敗前要重試幾次retry_delay-工作參數,接受 timedelta,定義每次重試間的等待時間@task(
trigger_rule=TriggerRule.ONE_SUCCESS,
retries=3,
retry_delay=timedelta(minutes=5)
)
from airflow.providers.standard.operators.trigger_dagrun
import TriggerDagRunOperatortrigger_dag_id-必須與 dag_id 相同wait_for_completion-阻擋後續工作直到完成poke_interval-檢查頻率conf-傳給子 Dag 的資料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 入門