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

