如何按指定所有者导出最新Airflow DAG及任务状态的表格数据
按所有者筛选导出Airflow DAG及任务最新状态实现方案
核心实现逻辑
无需手动指定DAG ID,直接通过Airflow REST API的原生过滤能力实现按所有者筛选,全流程分4步完成:
- 第一步:调用获取DAG列表的接口,通过
owners查询参数直接筛选出指定所有者的所有DAG,返回结果中包含所有匹配的DAG ID基础信息 - 第二步:调用DAG运行记录接口,传入上一步获取的DAG ID列表,通过排序和限制参数直接拉取每个DAG的最新一次运行记录,获取DAG级别的最新状态
- 第三步:调用任务实例接口,传入对应DAG ID和最新运行记录的ID,批量拉取所有DAG下属任务的最新执行状态
- 第四步:将获取到的DAG状态、任务状态按需要的字段整理,可直接导出为CSV/Excel表格格式,或直接对接监控看板作为数据源
关键实现代码示例
筛选指定所有者的DAG
import requests # 配置认证信息,根据你的Airflow实际认证方式调整 headers = {"Authorization": "Bearer 你的Airflow API Token"} # 替换为你的Airflow服务地址、目标所有者名称 airflow_host = "http://你的Airflow服务地址" target_owner = "指定的所有者名称" # 拉取指定所有者的所有DAG dag_resp = requests.get( f"{airflow_host}/api/v1/dags?owners={target_owner}", headers=headers ) dag_list = [dag["dag_id"] for dag in dag_resp.json()["dags"]]
批量获取DAG最新运行状态
# 传入第一步获取的所有DAG ID,按执行时间倒序拉取最新运行记录 dag_ids_param = ",".join(dag_list) run_resp = requests.get( f"{airflow_host}/api/v1/dagRuns?dag_ids={dag_ids_param}&order_by=-execution_date&limit={len(dag_list)}", headers=headers ) latest_dag_runs = run_resp.json()["dag_runs"]
拉取对应任务的最新状态
task_status_data = [] for run in latest_dag_runs: dag_id = run["dag_id"] run_id = run["dag_run_id"] # 拉取单个DAG运行的所有任务状态 task_resp = requests.get( f"{airflow_host}/api/v1/dags/{dag_id}/dagRuns/{run_id}/taskInstances", headers=headers ) for task in task_resp.json()["task_instances"]: task_status_data.append({ "DAG ID": dag_id, "DAG最新状态": run["state"], "DAG执行时间": run["execution_date"], "任务ID": task["task_id"], "任务最新状态": task["state"], "任务开始时间": task["start_date"], "任务结束时间": task["end_date"] })
导出为表格文件
import pandas as pd pd.DataFrame(task_status_data).to_excel("Airflow工作流状态统计表.xlsx", index=False)
注意事项
- Airflow REST API默认需要认证,可根据你的部署配置选择基础认证、Bearer Token认证等方式,请求时需携带有效认证信息
- 若匹配的DAG数量较多,可调整接口的分页参数拉取全量数据,避免遗漏
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

