Airflow原生实现跨DAG依赖调度的正确方案
Airflow 原生跨DAG依赖调度实现方案
调度规则基础配置
- 3个DT1类DAG统一设置调度间隔为
0 * * * *,即每小时整点触发,单实例运行完成后标记整体状态为success即可。 - DT2类DAG按需设置调度间隔为
0 22 * * *(22点触发)或0 23 * * *(23点触发),示例以22点运行为准。 - 所有DAG统一配置相同的
start_date,关闭catchup参数,避免历史补数实例干扰依赖匹配逻辑。
核心依赖校验实现
直接使用Airflow原生自带的ExternalTaskSensor算子做跨DAG状态检测,不需要安装任何第三方插件,也不需要自定义算子:
- 在DT2的DAG文件中初始化3个
ExternalTaskSensor任务,每个任务对应一个DT1 DAG的状态检测。 - 每个Sensor配置
external_dag_id为对应DT1的DAG ID,external_task_id设为None,代表检测整个DT1 DAG实例的运行状态,不需要绑定DT1内部的具体任务。 - 配置
execution_delta=timedelta(hours=1),自动匹配DT2当前调度时点往前推1小时对应的DT1运行实例,不需要手动拼接计算执行时间。 - 配置
mode="reschedule",Sensor在等待期间会释放worker资源,不会长期占用工作槽位;poke_interval设为60即每分钟检测一次依赖状态,timeout设为3600即最长等待1小时,超时直接标记失败触发告警。 - 配置
allowed_states=["success"],仅当对应DT1实例状态为success时判定依赖满足。
任务流依赖配置
将DT2内部所有数据转换、聚合类业务任务,全部设置为3个ExternalTaskSensor任务的下游,只有3个依赖检测任务全部通过,才会启动后续业务处理。
核心实现代码参考:
from datetime import datetime, timedelta from airflow import DAG from airflow.sensors.external_task import ExternalTaskSensor from airflow.operators.python import PythonOperator # 替换为实际的3个DT1的DAG ID DT1_DAG_LIST = ["dt1_source_sync_1", "dt1_source_sync_2", "dt1_source_sync_3"] default_args = { "owner": "data_team", "depends_on_past": False, "retries": 1 } with DAG( dag_id="dt2_daily_agg", default_args=default_args, schedule_interval="0 22 * * *", start_date=datetime(2024, 1, 1), catchup=False, tags=["data_process"] ) as dag: # 初始化3个DT1依赖检测任务 dt1_sensor_list = [] for dt1_id in DT1_DAG_LIST: sensor = ExternalTaskSensor( task_id=f"wait_{dt1_id}_success", external_dag_id=dt1_id, execution_delta=timedelta(hours=1), allowed_states=["success"], mode="reschedule", poke_interval=60, timeout=3600 ) dt1_sensor_list.append(sensor) # 实际业务处理任务示例 def run_data_transform(): # 此处替换为实际的数据转换、聚合逻辑 print("开始从数据湖读取数据执行计算") transform_task = PythonOperator( task_id="run_data_aggregation", python_callable=run_data_transform ) # 配置依赖关系 dt1_sensor_list >> transform_task
避坑说明
- 不要用
TriggerDagRunOperator从DT1侧触发DT2,不符合DT2按自身定时规则调度的要求,且无法保证3个DT1全部成功后才触发。 - 不要把Sensor的
mode设为poke,会长期占用worker资源,高负载场景下会导致worker阻塞。 - 若配置了DAG序列化、权限管控,需要确保DT2的运行账号有读取3个DT1 DAG实例状态的权限,否则Sensor会报权限错误。
内容的提问来源于stack exchange,提问作者Takito Isumoro
相关产品推荐
相关产品推荐

