开始使用免费开始使用

触发子 Dag

您注意到部分工作流使用了相似的组件,于是想到可以把通用任务拆分到各自的 Dag 中。这样就能在需要时独立运行这些组件,而无需维护多份拷贝。您决定在当前工作流中的某个任务里执行一个子 Dag。

dagtaskdatetime 已为您导入。

本练习是课程的一部分

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()
编辑并运行代码