Airflow数据感知调度:单数据集更新触发DAG调度的实现方式
实现任意数据集更新触发DAG:无需依赖传统TriggerDagRunOperator
不用非得靠传统的TriggerDagRunOperator来实现这个需求!Airflow的数据集调度功能本身就能支持“任意一个数据集更新即触发DAG”的逻辑,不同版本还有更贴合场景的简洁实现方式:
1. Airflow 2.5+版本:用UnionDataset直接实现OR逻辑
Airflow 2.5及以上版本引入了UnionDataset,专门用来处理“任意一个数据集更新就触发”的场景,完美替代默认的AND逻辑。代码示例如下:
from airflow.datasets import Dataset, UnionDataset from datetime import datetime from airflow.operators.bash import BashOperator # 定义你的数据集 example_dataset_1 = Dataset("s3://your-bucket/path1") example_dataset_2 = Dataset("s3://your-bucket/path2") example_dataset_3 = Dataset("s3://your-bucket/path3") with DAG( dag_id='any_dataset_trigger_example', schedule=UnionDataset([example_dataset_1, example_dataset_2, example_dataset_3]), start_date=datetime(2023, 1, 1), catchup=False, tags=['dataset-trigger'] ): # 示例任务逻辑 sample_task = BashOperator( task_id='print_trigger_msg', bash_command='echo "DAG triggered by an updated dataset!"' )
当列表里的任意一个数据集自上次运行后更新过,这个DAG就会被自动触发,完全匹配你的需求。
2. Airflow 2.4版本:轻量结合TriggerDagRunOperator(非传统用法)
如果你还在使用2.4.0版本(也就是你提到的官方文档版本),UnionDataset尚未上线,这时候可以用一种更贴合数据集调度模型的方式,而非传统的手动触发逻辑:
- 创建3个极简的“监听DAG”,每个DAG的
schedule只绑定一个目标数据集; - 每个监听DAG里仅放置一个
TriggerDagRunOperator,用来触发你真正要运行的业务DAG。
比如,其中一个监听DAG的代码:
from airflow.datasets import Dataset from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime example_dataset_1 = Dataset("s3://your-bucket/path1") with DAG( dag_id='listener_for_dataset_1', schedule=[example_dataset_1], start_date=datetime(2023, 1, 1), catchup=False, ): trigger_business_dag = TriggerDagRunOperator( task_id='trigger_target_dag', trigger_dag_id='your_core_business_dag_id', wait_for_completion=False )
另外两个监听DAG依葫芦画瓢,分别绑定example_dataset_2和example_dataset_3即可。这样,任意一个数据集更新时,对应的监听DAG会自动运行,触发你的业务DAG。
这种方式和传统的TriggerDagRunOperator用法不同——它完全基于数据集的更新事件触发,而非手动指定触发条件,更符合Airflow的数据集调度理念。
内容的提问来源于stack exchange,提问作者searain
相关产品推荐
相关产品推荐

