如何通过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
相关产品推荐
相关产品推荐

