Airflow如何从外部DAG中获取指定任务的当前运行状态
跨DAG查询指定任务实时状态的实现方案
你之前只查询到了DAG运行实例(DagRun)的状态,任务级别的状态存储在任务实例(TaskInstance)对象中,直接查询该表即可获取实时状态,不会像Sensor一样阻塞等待。
核心实现代码
首先导入依赖的模型类,封装查询函数:
from typing import Optional from datetime import datetime from airflow.models import DagRun, TaskInstance def get_external_task_state( external_dag_id: str, external_task_id: str, exec_dt: datetime ) -> Optional[str]: # 查询匹配的DAG运行实例 dag_runs = DagRun.find( dag_id=external_dag_id, execution_date=exec_dt, no_backfills=True ) if not dag_runs: return None target_dag_run = dag_runs[0] # 查询对应任务的实例 task_instances = TaskInstance.find( dag_id=external_dag_id, task_id=external_task_id, execution_date=target_dag_run.execution_date ) return task_instances[0].state if task_instances else None
场景适配方案
你可以通过ShortCircuitOperator调用上面的查询函数,实现DAG2任务运行时自动跳过DAG1同名任务的逻辑:
from airflow.operators.python import ShortCircuitOperator def skip_if_dag2_task_running(**context): exec_dt = context["execution_date"] task_state = get_external_task_state( external_dag_id="你的DAG2的ID", external_task_id="你要查询的同名任务ID", exec_dt=exec_dt ) # 返回False时会跳过后续所有下游任务 return task_state != "running" # 定义检查任务,放在需要跳过的同名任务上游 check_task = ShortCircuitOperator( task_id="check_dag2_task_state", python_callable=skip_if_dag2_task_running, provide_context=True, dag=dag # 这里替换为你的DAG1实例对象 ) # 配置依赖关系 check_task >> 你DAG1中的同名任务实例
注意事项
- Airflow 2.x版本可直接用
TaskInstance.get_task_instance()方法替代TaskInstance.find(),调用逻辑一致 - 若开启了DAG权限控制,需要确保DAG1的执行角色拥有DAG2的只读访问权限
- 该方法直接查询Airflow元数据库,返回的是查询瞬间的实时状态,不会产生阻塞等待
内容的提问来源于stack exchange,提问作者sks
相关产品推荐
相关产品推荐

