【发布时间】:2020-06-25 20:32:37
【问题描述】:
要求是让 DAG 一个接一个地运行,并且每个 DAG 都成功
我有一个 Master DAG,我在其中调用所有 DAG 以依次执行
此外,在每个 dag_A、dag_B、dag_C 中,我必须给定 schedule_interval = None 并在 GUI 中手动打开
我正在使用 ExternalTaskSensor,因为即使在第一个 dag_A 中的所有任务完成之前,它也会启动第二个 dag_B,为避免此类问题,我正在使用 ExternalTaskSensor。如果有更好的实现,请告诉我
不知道我在这里缺少什么
代码:master_dag.py
import datetime
import os
from datetime import timedelta
from airflow.models import DAG, Variable
from airflow.operators.dagrun_operator import TriggerDagRunOperator
from airflow.operators.sensors import ExternalTaskSensor
default_args = {
'owner': 'airflow',
'start_date': datetime.datetime(2020, 1, 7),
'provide_context': True,
'execution_timeout': None,
'retries': 0,
'retry_delay': timedelta(minutes=3),
'retry_exponential_backoff': True,
'email_on_retry': False,
}
dag = DAG(
dag_id='master_dag',
schedule_interval='7 3 * * *',
default_args=default_args,
max_active_runs=1,
catchup=False,
)
trigger_dag_A = TriggerDagRunOperator(
task_id='trigger_dag_A',
trigger_dag_id='dag_A',
dag=dag,
)
wait_for_dag_A = ExternalTaskSensor(
task_id='wait_for_dag_A',
external_dag_id='dag_A',
external_task_id='proc_success',
poke_interval=60,
allowed_states=['success'],
dag=dag,
)
trigger_dag_B = TriggerDagRunOperator(
task_id='trigger_dag_B',
trigger_dag_id='dag_B',
dag=dag,
)
wait_for_dag_B = ExternalTaskSensor(
task_id='wait_for_dag_B',
external_dag_id='dag_B',
external_task_id='proc_success',
poke_interval=60,
allowed_states=['success'],
dag=dag)
trigger_dag_C = TriggerDagRunOperator(
task_id='trigger_dag_C',
trigger_dag_id='dag_C',
dag=dag,
)
trigger_dag_A >> wait_dag_A >> trigger_dag_B >> wait_dag_B >> trigger_dag_C
每个 DAG 都有多个任务在运行,最后一个任务是 proc_success
【问题讨论】:
标签: airflow