利用上游参数判断下游数据感知DAG是否应运行
基于数据集更新与上游参数控制的Airflow DAG调度实现
需求说明
仅在以下两个条件同时满足时,才触发当前DAG的调度:
- 目标数据集表完成更新
- 上游触发DAG的运行参数
kick_off_downstream_dag被设置为true
实现代码
@dag( start_date=pendulum.datetime(year=2023, month=1, day=1, tz="America/New_York"), schedule=dataset if "{{ (triggering_dataset_events.values() | first | first).source_dag_run.conf['kick_off_downstream_dag']}}" else None, description=None, catchup=False, max_active_runs=1, default_args=DAG_ARGS, ) def your_dag_name(): # 在此编写DAG的任务逻辑 pass
关键逻辑解释
schedule参数通过条件判断实现双约束:利用Jinja模板语法从上游触发事件中提取kick_off_downstream_dag参数值,判断是否满足触发条件- 当参数为真时,以
dataset作为调度触发器(监听数据集更新);不满足条件时,设置schedule=None,即不自动触发调度 - 额外配置:关闭回溯执行(
catchup=False)、限制最大活跃运行数为1(max_active_runs=1),保障DAG运行的稳定性
内容的提问来源于stack exchange,提问作者David Whiting
相关产品推荐
相关产品推荐

