如何获取Airflow当前DAG运行实例的所有任务及状态列表
Airflow 当前DAG运行实例任务状态统计实现
核心实现逻辑是通过PythonOperator执行时注入的上下文对象,直接关联当前DAG Run拉取所有关联的任务实例数据,不需要硬编码DAG标识或运行时间参数,不会出现数据串跑的问题。
具体实现步骤
- 编写统计用的Python callable函数,从上下文参数中获取当前运行的DAG、DAG Run、当前任务ID信息,查询关联的所有任务实例遍历统计
- 给末尾的统计任务配置正确的触发规则,保证无论上游任务成功失败,统计任务都会正常执行
- 按需处理统计结果,可输出日志、写入XCom供下游使用、或推送到外部存储
可直接复用的代码示例
from airflow.utils.state import State from airflow.utils.trigger_rule import TriggerRule from airflow.operators.python import PythonOperator def count_task_status(**context): # 从上下文获取当前运行的核心标识 current_task_id = context["task"].task_id # 拉取当前DAG Run下所有任务实例 task_instances = context["dag_run"].get_task_instances() status_stat = { State.SUCCESS: 0, State.FAILED: 0 } task_status_detail = [] for ti in task_instances: # 跳过当前正在运行的统计任务本身,避免干扰统计结果 if ti.task_id == current_task_id: continue task_status_detail.append( {"task_id": ti.task_id, "state": ti.state} ) # 仅统计目标的成功、失败状态 if ti.state in status_stat: status_stat[ti.state] += 1 # 结果处理:打印日志、返回值自动写入XCom print(f"全量任务状态明细:{task_status_detail}") print(f"统计结果:成功任务{status_stat[State.SUCCESS]}个,失败任务{status_stat[State.FAILED]}个") return {"detail": task_status_detail, "stat": status_stat} # DAG定义内的任务配置示例 with DAG( dag_id="your_dag_id", # 其余DAG基础参数按现有配置填写 ) as dag: # 其他业务任务按原有逻辑定义... # 末尾的统计任务 task_status_count = PythonOperator( task_id="count_task_status", python_callable=count_task_status, # 核心配置:所有上游任务执行完成(无论成败)就触发当前任务 trigger_rule=TriggerRule.ALL_DONE ) # 配置所有业务任务为统计任务的上游,按实际DAG依赖编写即可 # [task1, task2, task3...] >> task_status_count
注意事项
- 必须配置
trigger_rule=TriggerRule.ALL_DONE,如果用默认的ALL_SUCCESS规则,只要有任意一个上游任务失败,统计任务会直接被标记为跳过,无法执行统计 - 不要直接全局查询TaskInstance表做过滤,通过
context['dag_run'].get_task_instances()获取的实例是和当前运行强绑定的,不会查到其他DAG、其他历史运行批次的任务数据,准确性最高 - 如果需要统计跳过、上游失败、重试中等其他状态,直接在
status_stat字典里添加对应的State常量即可 - Airflow 1.x版本使用时需要给PythonOperator加上
provide_context=True参数,才能正常拿到上下文对象;2.x版本默认开启上下文注入,不需要额外配置 - 动态生成的任务、任务组内的子任务都可以被正常统计,该方法是从运行时的任务实例表拉取数据,不是静态解析DAG文件
内容的提问来源于stack exchange,提问作者Vrusha
相关产品推荐
相关产品推荐

