触发子 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()