觸發子 Dag
你發現有些工作流程使用了相似的元件,於是想到可以把共用的任務拆成獨立的 Dag。這樣就能在需要時執行這些元件,而不必維護多份拷貝。你決定在目前的工作流程中,透過一個任務來執行子 Dag。
dag、task 和 datetime 元件已為你匯入。
本練習屬於課程
Python 中的 Apache Airflow 入門
練習說明
- 匯入可在你的工作流程中啟動 Dag 的 operator。
- 將該 operator 設定為觸發名為
child_pipeline的 Dag。 - 確保父 Dag 會在觸發的 Dag 完成後再繼續。
- 設定該 operator 檢查子 Dag 是否完成的頻率。
動手互動練習
試著完成這個範例程式碼,體驗一下這個練習。
# Import TriggerDagRunOperator
from airflow.providers.standard.operators.trigger_dagrun import ____
@dag(start_date=datetime(2026, 1, 1))
def parent_orchestrator_dag():
# Trigger child_pipeline and wait for it to complete
trigger_child = TriggerDagRunOperator(
task_id="trigger_child_pipeline",
trigger_dag_id="____",
____=True,
____=30,
conf={"source": "s3://my-bucket/raw/"})
validate() >> trigger_child >> post_trigger_summary()
parent_orchestrator_dag()