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

如何在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_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:22:21