You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何配置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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 21:10:37