如何基于Airflow中DAG的运行状态创建Python DataFrame
Airflow DAG状态DataFrame实现方案
两种方案按需选择即可,均基于Airflow原生能力实现,不需要额外开发复杂组件。
方案1:直连元数据库查询(推荐,性能最高)
适合脚本部署在Airflow集群内网、可以直接访问Airflow元数据库的场景。
- 前置依赖安装:
pip install pandas sqlalchemy - 实现逻辑:Airflow所有DAG运行记录都存在元数据库的
dag_run表中,通过窗口函数取每个DAG最新一次运行的状态,直接映射为你需要的三类中文状态即可,原生状态和目标状态对应关系:success→ 执行成功failed→ 执行失败running→ 正在运行
- 可直接运行的参考代码:
import pandas as pd from sqlalchemy import create_engine # 替换为实际的Airflow元数据库连接串,支持PostgreSQL/MySQL等 engine = create_engine("postgresql://账号:密码@数据库地址:端口/airflow库名") # 若只需要指定20个DAG,在WHERE子句中加 dag_id IN ('dag1','dag2'...) 过滤即可 query_sql = """ SELECT dag_id, state FROM ( SELECT dag_id, state, ROW_NUMBER() OVER (PARTITION BY dag_id ORDER BY execution_date DESC) AS run_rank FROM dag_run WHERE state IN ('success', 'failed', 'running') ) ranked_runs WHERE run_rank = 1 """ # 读取数据并做字段、值映射 df = pd.read_sql(query_sql, engine) state_mapping = {"success": "执行成功", "failed": "执行失败", "running": "正在运行"} df = df.replace({"state": state_mapping}) df.columns = ["DAG名称", "运行状态"]
- 补充:如果遇到排队(queued)、重试中(up_for_retry)这类中间状态,直接归到「正在运行」分类即可,在state_mapping里加对应键值对就行。如果要把这个逻辑做成Airflow内部的监控DAG,直接用Airflow自带的
PostgresHook/MySqlHook替代sqlalchemy连接即可,不需要硬编码数据库账号密码,更安全。
方案2:调用Airflow REST API(适合远程访问场景)
如果没法直连元数据库,用Airflow 2.0+自带的稳定REST API拉取数据即可,不需要额外开放数据库权限。
- 前置依赖安装:
pip install pandas requests - 实现逻辑:先获取DAG列表,再逐个拉取每个DAG最新的运行记录,做状态映射后转成DataFrame。
- 参考代码:
import requests import pandas as pd # 替换为实际的Airflow Web服务地址和登录账号 BASE_URL = "http://airflow-web地址:端口/api/v1" AUTH_INFO = ("登录用户名", "登录密码") # 如果是固定20个DAG,直接把下面的dag_list替换成你的DAG名称列表即可,不用拉全量 all_dags = requests.get(f"{BASE_URL}/dags", auth=AUTH_INFO, params={"limit": 100}).json()["dags"] dag_list = [dag["dag_id"] for dag in all_dags] state_mapping = {"success": "执行成功", "failed": "执行失败", "running": "正在运行"} result_rows = [] for dag_id in dag_list: # 按执行时间倒序取最新1条运行记录 latest_run = requests.get( f"{BASE_URL}/dags/{dag_id}/dagRuns", auth=AUTH_INFO, params={"limit": 1, "order_by": "-execution_date"} ).json()["dag_runs"] if not latest_run: continue # 无运行记录的DAG可按需自定义处理逻辑 raw_state = latest_run[0]["state"] if raw_state in state_mapping: result_rows.append({"DAG名称": dag_id, "运行状态": state_mapping[raw_state]}) df = pd.DataFrame(result_rows, columns=["DAG名称", "运行状态"])
- 补充:如果要降低接口调用频率,固定20个DAG的场景下直接硬编码DAG列表即可,不用每次拉全量DAG信息。
内容的提问来源于stack exchange,提问作者h.b
相关产品推荐
相关产品推荐

