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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 10:03:20