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

Airflow跨DAG依赖咨询:每日任务依赖小时任务的实现方案

Airflow不同调度频率跨DAG依赖实现方案

针对你描述的场景(每日5点运行的DAG1的taskB需同时依赖自身taskA和每小时15分运行的DAG2的taskY),以下是两种可靠的实现方式:

方案一:使用ExternalTaskSensor(兼容所有Airflow版本)

这是跨DAG依赖的经典实现方式,通过传感器监听目标DAG的指定任务是否完成,需精确匹配两个DAG的execution_date。

代码示例

DAG1(每日5点调度)

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.sensors.external_task import ExternalTaskSensor
from datetime import datetime, timedelta

def calculate_dag2_exec_date(dag1_exec_date):
    # DAG1的execution_date为调度日的前一天(如5点运行时,execution_date是前一天的日期)
    # 需等待DAG2在当日4点15分运行的taskY,对应DAG2的execution_date为当日4点整
    return dag1_exec_date + timedelta(days=1, hours=4)

with DAG(
    dag_id="dag1",
    schedule_interval="0 5 * * *",
    start_date=datetime(2024, 6, 1),
    catchup=False
) as dag:
    taskA = DummyOperator(task_id="taskA")

    wait_for_dag2_taskY = ExternalTaskSensor(
        task_id="wait_for_dag2_taskY",
        external_dag_id="dag2",
        external_task_id="taskY",
        execution_date_fn=calculate_dag2_exec_date,
        mode="reschedule",  # 释放worker资源,避免长期占用
        timeout=3600,  # 超时时间1小时,覆盖DAG2的运行窗口
        poke_interval=60  # 每分钟检查一次任务状态
    )

    taskB = DummyOperator(task_id="taskB")

    taskA >> wait_for_dag2_taskY >> taskB

关键说明

  • execution_date_fn 用于计算DAG2对应的execution_date,需根据你的调度规则调整逻辑;
  • 选用reschedule模式而非默认的poke模式,减少worker资源占用;
  • 如果需要补历史任务,需确保catchup=True时,execution_date的匹配逻辑依然正确。

方案二:使用Dataset(Airflow 2.4+推荐)

Airflow 2.4引入的Dataset功能,通过数据资产的产出/依赖关系实现跨DAG调度,无需手动处理execution_date匹配,更简洁优雅。

代码示例

DAG2(每小时15分调度)

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.datasets import Dataset
from datetime import datetime

# 定义数据集标识
dag2_taskY_dataset = Dataset("dataset://dag2_taskY_output")

with DAG(
    dag_id="dag2",
    schedule_interval="15 * * * *",
    start_date=datetime(2024, 6, 1),
    catchup=False
) as dag:
    taskX = DummyOperator(task_id="taskX")
    taskY = DummyOperator(
        task_id="taskY",
        outlets=[dag2_taskY_dataset]  # 声明该任务产出数据集
    )
    taskX >> taskY

DAG1(每日5点调度)

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.datasets import Dataset
from datetime import datetime

dag2_taskY_dataset = Dataset("dataset://dag2_taskY_output")

with DAG(
    dag_id="dag1",
    schedule_interval="0 5 * * *",
    start_date=datetime(2024, 6, 1),
    catchup=False
) as dag:
    taskA = DummyOperator(task_id="taskA")
    taskB = DummyOperator(
        task_id="taskB",
        inlets=[dag2_taskY_dataset]  # 声明该任务依赖数据集
    )
    # 同时依赖taskA完成
    taskA >> taskB

关键说明

  • Dataset会自动跟踪最新的产出版本,DAG1的taskB会在taskA完成且DAG2的taskY最新产出可用时执行;
  • 无需编写时间匹配逻辑,Airflow自动处理不同调度频率的依赖关系;
  • 适合依赖"最新数据产出"的场景,若需严格匹配特定时间窗口的任务,仍建议使用ExternalTaskSensor。

方案对比

方案兼容版本复杂度适用场景
ExternalTaskSensor所有Airflow版本中等需严格匹配时间窗口的跨DAG依赖
DatasetAirflow 2.4+低依赖最新数据产出的跨DAG依赖

内容的提问来源于stack exchange,提问作者user2262504

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 11:57:46