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依赖 |
| Dataset | Airflow 2.4+ | 低 | 依赖最新数据产出的跨DAG依赖 |
内容的提问来源于stack exchange,提问作者user2262504
相关产品推荐
相关产品推荐

