Airflow 1.10.12中如何从DAG A查询DAG B运行状态与配置JSON
在Airflow 1.10.12中从DAG A查询DAG B的运行状态与触发配置
你可以通过Airflow的核心ORM模型在DAG A的任务中直接查询DAG B的运行状态和触发配置,具体实现如下:
1. 查询DAG B是否正在运行
在DAG A中添加一个PythonOperator任务,通过DagRun模型筛选出状态为RUNNING的DAG B实例:
from airflow import DAG from airflow.operators.python_operator import PythonOperator from airflow.models import DagRun from datetime import datetime def check_dag_b_status(**context): # 替换为你的DAG B的dag_id dag_b_id = "dag_b_id" # 查询所有处于RUNNING状态的DAG B实例 running_dag_runs = DagRun.find(dag_id=dag_b_id, state="RUNNING") if running_dag_runs: print(f"DAG B 当前有 {len(running_dag_runs)} 个实例正在运行") # 可以进一步处理每个运行实例 for run in running_dag_runs: print(f"运行实例ID: {run.run_id}, 启动时间: {run.start_date}") else: print("DAG B 当前没有运行中的实例") # 定义DAG A default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1), } with DAG( 'dag_a', default_args=default_args, schedule_interval='*/720 * * * *', # 每12小时运行一次 catchup=False ) as dag: check_b_status_task = PythonOperator( task_id='check_dag_b_status', python_callable=check_dag_b_status, provide_context=True )
2. 获取DAG B的触发配置JSON
每个DagRun对象的conf属性就是触发时传入的配置JSON,你可以在查询到运行中的DAG B实例后直接读取:
def get_dag_b_config(**context): dag_b_id = "dag_b_id" running_dag_runs = DagRun.find(dag_id=dag_b_id, state="RUNNING") if running_dag_runs: for run in running_dag_runs: # 获取触发配置 trigger_config = run.conf print(f"DAG B 实例 {run.run_id} 的触发配置: {trigger_config}") # 如果需要提取具体字段,比如config里的key # specific_value = trigger_config.get("your_key") else: print("DAG B 当前没有运行中的实例") # 在DAG A中添加该任务 get_b_config_task = PythonOperator( task_id='get_dag_b_config', python_callable=get_dag_b_config, provide_context=True )
注意事项
- 如果DAG B存在多个同时运行的实例,上述代码会返回所有运行实例的状态和配置,你可以根据
start_date或run_id筛选最新的实例 - 确保执行DAG A的Airflow用户拥有访问DAG B元数据的权限(默认同一实例下权限是互通的)
- Airflow 1.10.x中
DagRun.find()方法支持按dag_id和state过滤,无需额外编写SQL查询
内容的提问来源于stack exchange,提问作者TreeWater
相关产品推荐
相关产品推荐

