Airflow 2.8.2:如何从DagBag获取数据感知调度DAG的Dataset URI?
解决方案
一、修复DagBag加载后dag.schedule为空的问题
Airflow 2.8.x中,Data-Aware调度的DAG会将schedule_interval设为"Dataset"作为标识,真正的Dataset依赖存储在dag.schedule属性中。直接用DagBag()加载时可能未触发Dataset解析逻辑,导致schedule为空,可通过两种方式处理:
- 显式触发调度解析:加载DagBag后遍历DAG,调用
resolve_schedule()强制解析Dataset依赖:from airflow.models import DagBag dagbag = DagBag(dag_folder="/path/to/your/dags") for dag_id, dag in dagbag.dags.items(): if dag.schedule_interval == "Dataset": dag.resolve_schedule() # 此时dag.schedule会填充对应Dataset集合 - 使用测试专用DagBag工具:改用
airflow.utils.testing.get_test_dagbag(),它会自动处理初始化逻辑:from airflow.utils.testing import get_test_dagbag dagbag = get_test_dagbag(dag_folder="/path/to/your/dags") # Data-Aware调度的DAG的schedule属性已自动加载完成
二、解决DagModel.schedule_datasets的DetachedInstanceError问题
该错误是SQLAlchemy会话脱离导致的——DagModel对象不在活跃会话中,无法延迟加载关联的schedule_datasets,解决方案如下:
- 预先加载关联数据:查询
DagModel时用joinedload主动加载关联的Dataset,避免延迟加载出错:from airflow.models import DagModel from airflow.settings import Session from sqlalchemy.orm import joinedload session = Session() try: dag_models = session.query(DagModel).options( joinedload(DagModel.schedule_datasets) ).filter(DagModel.schedule_interval == "Dataset").all() for dm in dag_models: # 可直接访问dm.schedule_datasets print(f"DAG {dm.dag_id}依赖的Dataset: {[ds.uri for ds in dm.schedule_datasets]}") finally: session.close() - 重新绑定对象到活跃会话:如果已经执行过
dagbag.sync_to_db(),可将脱离的DagModel对象重新合并到当前会话:from airflow.models import DagModel from airflow.settings import Session dagbag.sync_to_db() session = Session() # 假设dm是之前获取的脱离会话的DagModel对象 dm = session.merge(dm) # 现在可正常访问dm.schedule_datasets
三、完整测试用例示例
结合上述两点,编写验证Dataset引用合法性的测试代码:
import pytest from airflow.models import DagBag, Dataset from airflow.settings import Session from airflow.utils.testing import get_test_dagbag from sqlalchemy.orm import joinedload def test_data_aware_datasets_exist_in_producers(): dagbag = get_test_dagbag(dag_folder="/path/to/your/dags") session = Session() try: # 收集所有Task outlets中产出的Dataset produced_datasets = set() for dag in dagbag.dags.values(): for task in dag.tasks: produced_datasets.update( ds for ds in task.outlets if isinstance(ds, Dataset) ) # 验证Data-Aware DAG的调度Dataset均为已产出的Dataset # 方式1:通过DagBag中的DAG对象验证 for dag in dagbag.dags.values(): if dag.schedule_interval == "Dataset": dag.resolve_schedule() for scheduled_ds in dag.schedule: assert scheduled_ds in produced_datasets, \ f"DAG {dag.dag_id}引用的Dataset {scheduled_ds.uri}未被任何Task产出" # 方式2:通过DagModel验证 dag_models = session.query(DagModel).options( joinedload(DagModel.schedule_datasets) ).filter(DagModel.schedule_interval == "Dataset").all() for dm in dag_models: for scheduled_ds in dm.schedule_datasets: assert scheduled_ds in produced_datasets, \ f"DAG {dm.dag_id}引用的Dataset {scheduled_ds.uri}未被任何Task产出" finally: session.close()
内容的提问来源于stack exchange,提问作者Erik Ogan
相关产品推荐
相关产品推荐

