如何配置Airflow DAG同时按小时调度并在Dataset更新后触发?
解决方案:Airflow DAG同时支持小时调度与Dataset触发
针对你需要DAG B仅在对应小时数据(由DAG A产出)就绪后才运行的场景,这里提供两种实用实现方式,核心思路是通过带时间维度的Dataset关联上下游,再结合调度或任务条件判断控制执行逻辑。
1. 先给DAG A配置按小时划分的Dataset
首先要让DAG A在同步完某小时的数据后,标记对应唯一的Dataset,这样DAG B能精准关联到目标小时的数据。
from airflow import Dataset from airflow.decorators import dag, task from datetime import datetime # DAG A:低频同步数据,产出按小时区分的Dataset @dag( schedule="@every 4h", # 按实际低频需求设置,比如每4小时跑一次 start_date=datetime(2024, 1, 1), catchup=False ) def dag_a_sync_hourly_data(): @task(outlets=[Dataset(f"s3://your-data-bucket/hourly-data/{{ execution_date.strftime('%Y%m%d%H') }}")]) def rsync_hourly_data(): # 这里写rsync同步逻辑,比如同步对应小时的数据源到存储 print(f"完成 {{{{ execution_date.strftime('%Y-%m-%d %H:00') }}}} 时段的数据同步") rsync_hourly_data() dag_a_sync_hourly_data()
这里用outlets声明任务产出的Dataset,URI里通过execution_date格式化出小时维度的路径,确保每个小时的数据对应独立的Dataset标识。
2. 配置DAG B的双重触发逻辑
根据你的需求,有两种实现方式可选:
方式一:小时调度+任务级条件判断(保留小时调度实例,仅在数据就绪时执行)
如果需要保留DAG B的小时调度实例,但仅在对应Dataset更新后才执行实际处理任务,可以通过分支任务判断Dataset状态:
from airflow import Dataset from airflow.decorators import dag, task from airflow.operators.empty import EmptyOperator from datetime import datetime from airflow.utils.db import provide_session from airflow.models.dataset import DatasetEvent # 复用Dataset URI模板 DATASET_URI_TPL = "s3://your-data-bucket/hourly-data/{hour_str}" @dag( schedule="@hourly", # 按小时生成调度实例 start_date=datetime(2024, 1, 1), catchup=False ) def dag_b_process_hourly_data(): # 检查对应小时的Dataset是否在execution_date之后有更新 @provide_session def is_dataset_ready(execution_date, session=None): hour_str = execution_date.strftime("%Y%m%d%H") target_dataset = DATASET_URI_TPL.format(hour_str=hour_str) # 查询Dataset更新事件,确保事件时间晚于调度时间(避免旧数据触发) event = session.query(DatasetEvent).filter( DatasetEvent.dataset_uri == target_dataset, DatasetEvent.timestamp > execution_date ).first() return event is not None start = EmptyOperator(task_id="start") @task.branch(task_id="check_data_ready") def decide_execution(**context): if is_dataset_ready(context["execution_date"]): return "process_data" else: return "skip_process" @task() def process_data(**context): hour_str = context["execution_date"].strftime("%Y%m%d%H") target_dataset = DATASET_URI_TPL.format(hour_str=hour_str) # 这里写数据处理逻辑,比如读取Dataset对应的小时数据进行计算 print(f"开始处理 {target_dataset} 的小时数据") skip_process = EmptyOperator(task_id="skip_process") end = EmptyOperator(task_id="end", trigger_rule="none_failed_min_one_success") # 构建任务依赖 start >> decide_execution >> [process_data, skip_process] >> end dag_b_process_hourly_data()
- 核心是通过
is_dataset_ready函数查询Airflow的DatasetEvent表,确认目标小时的Dataset是否已更新。 - 用分支任务决定执行处理逻辑还是直接跳过,最后通过
trigger_rule保证DAG实例正常结束。
方式二:仅由Dataset触发(无无效调度实例)
如果不需要保留没有数据的小时调度实例,直接让DAG B仅在Dataset更新时触发对应小时的实例,这种方式更简洁:
from airflow import Dataset from airflow.decorators import dag, task from datetime import datetime DATASET_URI_TPL = "s3://your-data-bucket/hourly-data/{{ execution_date.strftime('%Y%m%d%H') }}" @dag( schedule=Dataset(DATASET_URI_TPL), # 仅由Dataset更新触发 start_date=datetime(2024, 1, 1), catchup=True # 开启catchup,确保历史未处理的Dataset也能触发对应实例 ) def dag_b_process_hourly_data(): @task() def process_data(**context): hour_str = context["execution_date"].strftime("%Y%m%d%H") target_dataset = DATASET_URI_TPL.replace("{{ execution_date.strftime('%Y%m%d%H') }}", hour_str) print(f"处理 {target_dataset} 的小时数据") process_data() dag_b_process_hourly_data()
- 这种方式下,DAG B的实例完全由DAG A的Dataset更新触发,且每个实例的
execution_date自动与Dataset对应的小时对齐。 - 开启
catchup=True可以处理DAG A提前产出的历史小时数据,确保没有遗漏。
内容的提问来源于stack exchange,提问作者qcha
相关产品推荐
相关产品推荐

