Airflow中DagRun.find为何返回含None属性的对象?
Airflow数据感知调度中DagRun.find返回None或空属性的原因解析
背景场景
我正在使用Airflow的数据感知调度(data-aware scheduling)运行任务,其中一个依赖更新数据集的任务需要获取触发当前流程的DAG运行的最早和最晚时间戳,具体步骤如下:
- 获取
triggering_dataset_events - 提取触发任务的
run_id - 通过
DagRun.find获取对应的DagRun对象 - 从DagRun对象中提取开始和结束时间戳
现有实现代码
import sys from toolz import seq from airflow.models import DagRun @task() def get_start_and_end_timestamp(ti=None): template_context = ti.get_template_context() triggering_ds = template_context["triggering_dataset_events"] times = ( seq(list(triggering_ds.values())) .flatten() # 需展平,因为原始格式为[[ds_event1,...]] .map(lambda ds_event: ds_event.source_run_id) .map(lambda source_run_id: DagRun.find(run_id=source_run_id)) .flatten() # 需展平,因为DagRun.find返回DagRun实例列表 .filter( lambda dag_run: dag_run is not None and dag_run.start_date is not None and dag_run.end_date is not None ) .map( lambda dag_run: ( dag_run.start_date.timestamp(), dag_run.end_date.timestamp(), ) ) .reduce( lambda x, y: (min(x[0], y[0]), max(x[1], y[1])), (sys.float_info.max, sys.float_info.min), ) ) # 校验结果有效性 assert times[0] != sys.float_info.max, "未找到有效的起始时间戳" assert times[1] != sys.float_info.min, "未找到有效的结束时间戳" # 将时间戳转换为微秒级整数 return (int(times[0] * 1e6), int(times[1] * 1e6))
遇到的问题
在实现过程中,发现DagRun.find返回的部分DagRun实例存在None属性(如start_date或end_date为空),必须通过过滤逻辑排除无效实例才能正常运行。
原因解析
1. DAG运行状态未完成
Airflow中,DagRun.end_date只有在整个DAG的所有任务都执行完成(成功或失败)后才会被赋值;如果触发的DAG运行仍在执行中,对应的end_date会是None。此外,极端情况下DAG刚启动时start_date也可能未被及时初始化(比如元数据库写入延迟)。
2. 元数据记录异常
如果对应的DAG运行记录被手动删除、或者元数据库出现数据损坏,DagRun.find可能返回属性缺失的实例;或者run_id对应的DAG运行已被取消,其start_date/end_date字段在数据库中被置为NULL。
3. 数据感知事件的延迟或失效
数据感知调度的triggering_dataset_events可能包含已失效的source_run_id(比如触发事件生成后,对应的DAG运行被删除),此时DagRun.find返回的实例属性会不完整。
DagRun.find返回内容说明
DagRun.find是Airflow提供的类方法,用于从元数据库中查询匹配条件的DAG运行记录:
- 返回类型为
list[DagRun],每个元素对应一条dag_run表中的记录 - 如果没有匹配条件的记录,返回空列表
- 实例的属性值完全对应元数据库中的字段:如果数据库中字段为NULL,实例的对应属性就是
None
内容的提问来源于stack exchange,提问作者TrendSpark
相关产品推荐
相关产品推荐

