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

如何通过Airflow SDK使用DAG run id查询对应DAG id及相关信息

解答:通过Airflow Python SDK查询指定DAG Run的相关信息

完全可以实现,Airflow提供了两种场景的查询方式:

场景1:在Airflow部署环境内部(scheduler/worker/webserver节点)查询

直接调用Airflow内置的ORM模型即可,无需额外鉴权,示例代码如下:

from airflow.models import DagRun, DagModel

# 替换为你要查询的DAG Run ID
target_run_id = "scheduled__2021-11-30T09:30:00+00:00"
# 按run_id查询
dag_run_list = DagRun.find(run_id=target_run_id)

if not dag_run_list:
    print("未找到匹配的DAG Run")
else:
    # *DAG Run ID全局唯一,查询结果最多只有1条*
    target_dag_run = dag_run_list[0]
    # 1. 获取对应DAG ID
    dag_id = target_dag_run.dag_id
    print(f"匹配的DAG ID为:{dag_id}")
    # 2. 提取DAG Run的其他信息
    print(f"DAG Run状态:{target_dag_run.state}")
    print(f"DAG Run执行日期:{target_dag_run.execution_date}")
    print(f"DAG Run启动时间:{target_dag_run.start_date}")
    print(f"DAG Run结束时间:{target_dag_run.end_date}")
    print(f"DAG Run携带的参数:{target_dag_run.conf}")
    # 3. 查询对应DAG的详细信息
    target_dag = DagModel.get_current(dag_id)
    print(f"DAG所有者:{target_dag.owners}")
    print(f"DAG调度周期:{target_dag.schedule_interval}")
    print(f"DAG是否暂停:{target_dag.is_paused}")
    print(f"DAG标签:{target_dag.tags}")

场景2:在外部环境通过Airflow Python SDK查询

适用于Airflow 2.3+版本,需要提前安装官方Python客户端、配置Airflow API的访问权限,示例代码如下:

from airflow_client.client import ApiClient, Configuration
from airflow_client.client.api import DAGRunApi, DAGApi

# 配置你的Airflow实例连接信息
config = Configuration(
    host="http://<你的Airflow webserver地址>:<端口>/api/v1",
    username="<你的Airflow登录账号>",
    password="<你的Airflow登录密码>"
)

with ApiClient(config) as client:
    # 全局查询匹配run_id的DAG Run
    run_api = DAGRunApi(client)
    query_result = run_api.get_dag_runs(run_id="scheduled__2021-11-30T09:30:00+00:00")
    if not query_result.dag_runs:
        print("未找到匹配的DAG Run")
    else:
        target_dag_run = query_result.dag_runs[0]
        # 获取DAG ID
        dag_id = target_dag_run.dag_id
        print(f"匹配的DAG ID为:{dag_id}")
        # 查询DAG详细信息
        dag_api = DAGApi(client)
        target_dag = dag_api.get_dag(dag_id=dag_id)
        print(f"DAG调度周期:{target_dag.schedule_interval}")

注意事项

  • 低于2.3版本的Airflow对外API不支持全局按run_id过滤DAG Run,建议直接用场景1的方式查询
  • 如果查询不到结果,先确认你使用的账号有对应DAG的访问权限、run_id输入无误

内容的提问来源于stack exchange,提问作者Jwala Sowmika

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 05:24:07