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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 04:17:38