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

如何通过Elasticsearch HTTP输入级联调用Airflow API获取DAG状态?

级联调用Airflow API及获取失败DAG信息的方案

关于Elasticsearch HTTP输入实现级联调用的可行性

原生Elasticsearch的HTTP相关组件(比如Ingest Pipeline的HTTP处理器)无法直接完成这种先获取DAG列表、再逐个调用dagRuns接口的级联逻辑——因为Elasticsearch的ingest流程是单文档处理模式,不支持接收列表后自动遍历发起批量HTTP请求并聚合结果。

如果要基于Elastic Stack实现这个需求,需要配合外部脚本或其他采集组件:

  • 用Python/Shell脚本先调用Airflow /dags接口过滤出带指定标签的DAG ID列表,再遍历每个ID调用dags/{dag_id}/dagRuns接口获取状态,最后将数据批量发送到Elasticsearch
  • 借助Filebeat/Elastic Agent的自定义采集逻辑,通过脚本扩展实现级联请求,再将采集到的数据导入ES

获取失败DAG信息的替代方案

除了调用API,还有几种更高效的方式获取失败DAG数据:

1. Airflow CLI命令

直接用Airflow官方CLI工具筛选目标数据:

  • 获取带指定标签的DAG列表:
    airflow dags list --tags "你的标签名"
    
  • 针对单个DAG筛选失败的运行实例:
    airflow dags list-runs --dag-id "目标DAG ID" --state failed
    

可以把这些命令封装成定时脚本,将输出结果写入文件或直接POST到Elasticsearch。

2. 直接查询Airflow元数据库

Airflow的所有DAG运行状态都存在元数据数据库(PostgreSQL/MySQL等)中,直接查询效率更高:
示例SQL(以PostgreSQL为例):

SELECT 
  d.dag_id, 
  dr.run_id, 
  dr.start_date, 
  dr.end_date, 
  dr.state
FROM dag d
JOIN dag_tag dt ON d.dag_id = dt.dag_id
JOIN dag_run dr ON d.dag_id = dr.dag_id
WHERE dt.tag = '你的标签名' AND dr.state = 'failed';

可以用Elasticsearch的JDBC输入插件,直接连接Airflow数据库定时同步数据。

3. Airflow告警回调

通过Airflow的回调机制,在DAG失败时主动推送数据:

  • 配置DAG的on_failure_callback函数,当DAG运行失败时,直接将失败详情发送到Elasticsearch的HTTP接口
  • 或者使用Airflow的内置告警系统(如Email、Slack),再通过Elastic的采集组件接收这些告警信息并导入ES

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 01:54:54