【发布时间】:2021-06-21 15:44:34
【问题描述】:
我有多个依赖于 initial_dag 的 DAG 首先运行,之后我希望依赖的 DAG 一个接一个地运行。这是我所拥有的:
dag = DAG(
dag_id=DAG_NAME,
default_args=default_args,
schedule_interval=None,
start_date=airflow.utils.dates.days_ago(1)
)
initial_dag = BashOperator(
task_id='initial_dag',
bash_command="python /home/airflow/gcs/dags/task.py",
dag=dag
)
dependent_dag1 = TriggerDagRunOperator(
task_id="dependent_dag1",
trigger_dag_id="dependent_dag1",
wait_for_completion=True,
dag=dag
)
dependent_dag2 = TriggerDagRunOperator(
task_id="dependent_dag2",
trigger_dag_id="dependent_dag2",
wait_for_completion=True,
dag=dag
)
dependent_dag3 = TriggerDagRunOperator(
task_id="dependent_dag3",
trigger_dag_id="dependent_dag3",
wait_for_completion=True,
dag=dag
)
initial_dag >> dependent_dag1 >> dependent_dag2 >> dependent_dag3
我认为wait_for_completion=True 会在触发下一个 DAG 之前完成每个 DAG 的运行。例如。 initial_dag 运行并完成,然后触发 dependent_dag1 并等待其完成以触发后续任务。
触发 DAG 的顺序是正确的,但它似乎并没有等待前一个 DAG 先完成,例如dependent_dag2 在 dependent_dag1 完成之前被触发。
我错过了什么吗?
【问题讨论】: