如何在Airflow中统计当前DAG Run的上游失败任务数量?
解决同一执行日期重跑DAG中失败任务统计的问题
我来帮你搞定这个问题!其实核心关键点就在于TaskInstance和DagRun之间的直接关联字段,你之前没找到的关联方式,Airflow早就给你准备好了。
核心原理:TaskInstance与DagRun的关联
每个TaskInstance记录都带有dag_run_id字段,这个字段直接对应DagRun表的主键ID。也就是说,不管你同一执行日期重跑多少次,每次重跑都会生成一个唯一的DagRun实例,而该实例下的所有TaskInstance都会绑定这个唯一的dag_run_id——这就是精准区分不同重跑记录的关键!
具体实现步骤
1. 在任务中获取当前DagRun的ID
在你的那个设置了ALL_DONE规则的任务里,可以通过Airflow的任务上下文(context)直接拿到当前的DagRun对象,进而获取它的唯一ID:
def your_all_done_task(**context): # 从上下文获取当前DAG运行实例 current_dag_run = context['dag_run'] current_dag_run_id = current_dag_run.id # 后续统计逻辑...
2. 统计当前DagRun中的失败任务
你可以选择用Airflow的ORM模型查询,或者直接写SQL,两种方式都能精准过滤出当前运行的失败任务:
方式一:使用Airflow ORM(推荐,更贴合Airflow生态)
from airflow.models import TaskInstance from airflow.utils.state import TaskInstanceState def count_failed_tasks(**context): current_dag_run_id = context['dag_run'].id # 过滤当前dag_run下的失败任务 failed_task_count = TaskInstance.query.filter( TaskInstance.dag_run_id == current_dag_run_id, TaskInstance.state == TaskInstanceState.FAILED ).count() print(f"当前DAG重跑实例中,失败任务数量为:{failed_task_count}") # 这里可以把结果存入变量或者做后续处理
方式二:直接执行SQL查询(适合自定义需求场景)
如果你的Airflow元数据库是PostgreSQL,SQL语句大概是这样的(替换<current_dag_run_id>为你拿到的ID):
SELECT COUNT(*) FROM task_instance WHERE dag_run_id = <current_dag_run_id> AND state = 'failed';
3. 关键注意事项
- 绝对不要只靠
execution_date过滤:同一执行日期可能有多个DagRun(重跑实例),只有结合dag_run_id才能精准定位到你当前正在运行的那一次。 - 手动重跑/自动重跑都会生成新的DagRun:每次重跑Airflow都会创建新的DagRun记录,所以
dag_run_id是完全唯一的,不会和历史运行混淆。 - 确保任务上下文可用:你的
ALL_DONE任务需要开启provide_context=True(在PythonOperator中设置),这样才能从context中拿到dag_run对象。
内容的提问来源于stack exchange,提问作者mad_
相关产品推荐
相关产品推荐

