
1. 项目缘起从一次深夜告警失联说起去年底我们团队负责的一套核心业务系统在凌晨两点发生了数据库连接池耗尽的问题。当时虽然Elastic StackELK的日志收集和Kibana的可视化监控都已就位但告警机制却完全失灵。我们依赖的是Kibana自带的邮件告警而邮件服务器恰好在那个时段进行维护。结果就是直到第二天早上用户投诉电话打进来我们才发现系统已经“静默”宕机了数小时。这次事故让我们痛定思痛告警通道必须可靠、即时且触达力强。邮件太被动短信成本高而团队日常沟通全在钉钉。于是一个明确的需求浮出水面能否让Kibana的日志告警直接推送到钉钉群调研后发现Kibana自带的告警功能如Watcher在基础版中能力有限而将告警事件通过Webhook推送到外部系统如钉钉机器人是白金版Platinum及以上版本才提供的企业级功能。对于许多中小团队或个人开发者而言订阅昂贵的白金版许可并不现实。因此探索一种在现有开源或基础版Kibana上实现类似“白金版”Webhook告警推送的方案就成了一项极具实用价值的技术实践。本文记录的就是我如何通过一系列技术组合拳在Kibana 8.5环境下搭建起一套稳定、可用的日志告警钉钉推送体系。请注意本文讨论的技术方案旨在学习和研究告警集成流程所有操作应在您拥有合法使用权的软件和环境上进行。2. 技术栈选型与架构设计思路要实现“Kibana触发告警 - 钉钉接收消息”这个链路核心在于解决两个问题告警规则的评估与触发以及告警信息的格式化与投递。Kibana 8.x版本后其告警框架Alerting Framework已经相当强大但免费版确实锁住了Action执行动作中的部分连接器Connectors比如Webhook。我们的思路是“曲线救国”。2.1 核心组件拆解告警引擎Alerting Engine 仍然使用Kibana内置的。它负责根据我们定义的规则Rule持续查询Elasticsearch中的数据一旦条件满足就生成告警实例Alert Instance。这部分功能在基础版中是可用的。动作执行器Action Executor 这是被限制的部分。理想情况下告警实例应直接触发一个配置好的Webhook动作。既然此路不通我们需要一个“中间人”。“中间人”代理服务Proxy Service 这是本方案的核心。我们需要一个轻量级、可靠的服务它能够以某种方式接收到Kibana告警引擎产生的告警事件然后将其转换为钉钉机器人要求的格式并通过HTTP请求发送出去。2.2 两种主流实现路径对比基于上述拆解主要有两种技术路径路径A利用Kibana的服务器日志Server Log与日志收集器。配置Kibana的告警规则当其触发时将告警信息以特定格式写入Kibana自身的服务器日志文件。然后使用Filebeat等日志采集器监控该日志文件将新增的告警日志发送到Elasticsearch中的一个特定索引。最后再创建一个Kibana告警规则监控这个“告警日志索引”并触发一个……等等这陷入了循环。此路径的瓶颈在于我们最终还是需要一个“出口”将事件推出去。它更适合用于告警事件的二次归档与审计而非直接对外通知。路径B使用Elasticsearch的写入后触发与外部服务。这是更优雅和直接的方式。Kibana告警触发后其状态和上下文信息会保存在一个特殊的Elasticsearch索引中通常是.kibana-alerting-*。我们可以编写一个外部服务例如用Python Flask、Node.js或Go编写定期或通过Elasticsearch的变更流Change StreamAPI监听这个索引的新增文档。一旦发现新的告警记录服务立即处理并调用钉钉Webhook。显然路径B在架构上更解耦可靠性更高也是企业级系统常见的“事件驱动”模式。我们选择这条路径。整个架构流程图如下所示Kibana Alerting Rule | v (触发告警写入ES) | v Elasticsearch Index: .kibana-alerting-* | v [ 外部代理服务 ] (定期轮询或监听ES变化) | v (格式化消息HTTP POST) | v 钉钉群机器人Webhook | v 钉钉群内告警消息3. 实战部署构建外部代理服务我们选择使用Python的Flask框架来构建这个代理服务因为它轻量、开发速度快且有成熟的Elasticsearch客户端和HTTP请求库。3.1 环境准备与依赖安装首先准备一台可以与Elasticsearch集群和钉钉网络互通的服务可以与Kibana同机也可独立。确保安装Python3和pip。# 创建项目目录并进入 mkdir kibana-dingtalk-proxy cd kibana-dingtalk-proxy # 创建虚拟环境推荐 python3 -m venv venv source venv/bin/activate # Linux/macOS # venv\Scripts\activate # Windows # 安装核心依赖 pip install elasticsearch flask requestselasticsearch: 用于连接和查询Elasticsearch。flask: 用于提供简单的API接口可选用于健康检查和组织应用。requests: 用于向钉钉机器人发送HTTP请求。3.2 核心服务代码实现创建一个名为app.py的文件以下是核心代码及详细注释import json import time import hashlib import base64 import hmac from datetime import datetime, timedelta from elasticsearch import Elasticsearch, exceptions from flask import Flask import requests import threading import logging # 配置日志 logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s) logger logging.getLogger(__name__) app Flask(__name__) # ---------- 配置区 ---------- # Elasticsearch 连接配置 ES_HOSTS [http://your-elasticsearch-host:9200] # 替换为你的ES地址 ES_USER elastic # 如果有认证 ES_PASSWORD your-password # 替换为你的密码 ALERT_INDEX_PATTERN .kibana-alerting-* # Kibana告警索引模式 # 钉钉机器人配置 DINGDING_WEBHOOK https://oapi.dingtalk.com/robot/send?access_tokenYOUR_TOKEN # 替换 DINGDING_SECRET YOUR_SECRET # 如果设置了加签替换 # 服务运行配置 POLL_INTERVAL 10 # 轮询ES的间隔秒数 LAST_PROCESSED_TIME_FILE last_processed.time # 记录上次处理时间戳的文件 # ---------- 初始化ES客户端 ---------- es Elasticsearch( hostsES_HOSTS, http_auth(ES_USER, ES_PASSWORD) if ES_USER and ES_PASSWORD else None, timeout30 ) def get_dingtalk_signature(): 生成钉钉机器人加签签名如果启用了加签 if not DINGDING_SECRET: return timestamp str(round(time.time() * 1000)) secret_enc DINGDING_SECRET.encode(utf-8) string_to_sign f{timestamp}\n{DINGDING_SECRET} string_to_sign_enc string_to_sign.encode(utf-8) hmac_code hmac.new(secret_enc, string_to_sign_enc, digestmodhashlib.sha256).digest() sign base64.b64encode(hmac_code).decode(utf-8) return ftimestamp{timestamp}sign{sign} def send_dingtalk_message(alert_body): 发送消息到钉钉机器人 webhook_url DINGDING_WEBHOOK get_dingtalk_signature() headers {Content-Type: application/json} # 构建钉钉Markdown消息格式 alert_instance alert_body.get(_source, {}).get(kibana.alert.instance, {}) rule_name alert_body.get(_source, {}).get(kibana.alert.rule.name, 未知规则) context alert_body.get(_source, {}).get(kibana.alert.context, {}) # 从context中尝试提取更具体的消息例如查询条件或结果 message context.get(message) or f告警规则【{rule_name}】被触发。 # 构建Markdown内容 dingtalk_msg { msgtype: markdown, markdown: { title: f Kibana告警 - {rule_name}, text: f### Kibana告警通知 **规则名称**{rule_name} **触发时间**{datetime.now().strftime(%Y-%m-%d %H:%M:%S)} **告警详情** {message} **关联实例**{json.dumps(alert_instance, ensure_asciiFalse)} **原始数据ID**{alert_body.get(_id)} 请相关同事及时处理。 }, at: { isAtAll: False # 不所有人可以根据规则严重性动态设置 # atMobiles: [138xxxxxxx] # 可以指定手机号 } } try: resp requests.post(webhook_url, headersheaders, datajson.dumps(dingtalk_msg), timeout10) resp.raise_for_status() result resp.json() if result.get(errcode) 0: logger.info(f钉钉消息发送成功: {rule_name}) return True else: logger.error(f钉钉接口返回错误: {result}) return False except Exception as e: logger.error(f发送钉钉消息失败: {e}) return False def load_last_processed_time(): 从文件加载上次处理的时间戳 try: with open(LAST_PROCESSED_TIME_FILE, r) as f: return datetime.fromisoformat(f.read().strip()) except FileNotFoundError: # 如果是第一次运行默认处理过去5分钟的数据 return datetime.utcnow() - timedelta(minutes5) def save_last_processed_time(t): 保存本次处理的时间戳 with open(LAST_PROCESSED_TIME_FILE, w) as f: f.write(t.isoformat()) def poll_alerts(): 轮询Elasticsearch获取新增告警 while True: try: last_time load_last_processed_time() current_time datetime.utcnow() # 构建ES查询DSL查询指定时间范围后新产生的告警 # kibana.alert.start 是告警开始时间字段 query { query: { range: { kibana.alert.start: { gt: last_time.isoformat(), lte: current_time.isoformat() } } }, sort: [{kibana.alert.start: {order: asc}}], size: 100 } response es.search(indexALERT_INDEX_PATTERN, bodyquery) hits response.get(hits, {}).get(hits, []) if hits: logger.info(f发现 {len(hits)} 条新告警) for hit in hits: # 发送到钉钉 success send_dingtalk_message(hit) # 可以根据success决定是否重试或记录失败 # 此处简单处理发送即视为处理 time.sleep(0.5) # 避免对钉钉接口请求过快 logger.info(f本轮告警处理完成) else: logger.debug(未发现新告警) # 更新处理时间为当前轮询开始时间避免遗漏 save_last_processed_time(current_time) except exceptions.ConnectionError as e: logger.error(f无法连接到Elasticsearch: {e}) except Exception as e: logger.error(f轮询过程中发生未知错误: {e}) # 等待下一个轮询周期 time.sleep(POLL_INTERVAL) app.route(/health) def health(): 健康检查端点 return {status: ok, service: kibana-dingtalk-proxy}, 200 if __name__ __main__: # 启动轮询线程 poll_thread threading.Thread(targetpoll_alerts, daemonTrue) poll_thread.start() logger.info(Kibana告警钉钉代理服务启动轮询线程已运行。) # 启动Flask Web服务可选主要用于健康检查 app.run(host0.0.0.0, port5000, debugFalse)3.3 关键配置与安全说明Elasticsearch连接安全 代码中使用了HTTP基础认证。在生产环境中强烈建议使用HTTPS并配置证书或者将服务部署在与ES集群同一安全网络内。如果ES开启了安全特性如API密钥请使用Elasticsearch客户端对应的认证方式。钉钉机器人创建在钉钉群 - 群设置 - 智能群助手 - 添加机器人 - 自定义机器人。设置机器人名称和关键词可选如果设置消息中需包含关键词。安全设置至关重要选择“加签”。将生成的secret填入代码的DINGDING_SECRET。Webhook URL中的access_token填入DINGDING_WEBHOOK。加签能有效防止Webhook被恶意调用。告警索引权限 运行此服务的账户需要对.kibana-alerting-*索引有读取权限。建议创建一个专属的ES用户并赋予最小必要权限。时间戳与容错 服务通过本地文件记录上次处理时间。如果服务重启会从文件中读取时间处理从那个时间点之后的新告警。这保证了告警不会丢失但可能导致服务重启后短时间内重复处理少量告警取决于轮询间隔。对于严格不重不漏的场景可以考虑将处理状态如告警ID持久化到数据库。3.4 服务部署与运行# 后台运行服务使用 nohup 或 systemd nohup python app.py proxy.log 21 # 查看日志 tail -f proxy.log建议使用systemd或supervisor等进程管理工具来托管此服务确保其高可用。4. Kibana告警规则配置详解外部服务就绪后我们需要在Kibana中配置能真正触发告警的规则。这里以最常见的“日志错误率超标”为例。4.1 创建告警规则登录Kibana进入Management - Stack Management。在左侧导航栏选择Rules and Connectors - Rules。点击Create rule。4.2 定义规则条件Rule ConditionsRule name:生产服务错误日志告警Check every:1m(每分钟检查一次)Index selection: 选择你的应用日志索引例如app-logs-*Time field:timestamp接下来是关键的定义查询和阈值部分。在Define rule conditions区域我们使用Kibana Query Language (KQL) 或 Lucene查询。Aggregation: 选择Count因为我们关心错误日志的数量。Group by: 选择Top values在Field中选择能区分服务或主机字段如service.name。Size设置为10。这样告警可以按服务分组避免一个服务出错淹没其他服务的告警。Filter: 输入查询条件例如log.level: ERROR。Threshold: 设置is above10。表示**每个分组每个服务**在时间窗口内错误日志超过10条即触发。注意这里的阈值逻辑是核心。Group by之后阈值判断是针对每个独立的分组进行的。如果不分组阈值就是针对所有匹配日志的总数。4.3 配置告警行动Actions- 模拟“Webhook”由于我们没有真正的Webhook连接器这里的Actions配置主要是为了让告警规则能够运行并生成告警实例这些实例会被写入.kibana-alerting-*索引从而被我们的代理服务捕获。在Actions部分点击Add action。Action type: 选择你可用的一种例如Email如果你配置了邮件服务器或者Index将告警记录到另一个索引。选择哪个并不重要因为我们不依赖它来通知只是为了触发规则。这里为了演示可以选择Index。Index: 可以指定一个如kibana-alert-history的索引用于归档。Message: 简单填写即可如Rule {{rule.name}} triggered for group {{context.group}}。关键点这个Action的执行与否不影响告警实例的生成。只要规则条件满足Kibana就会在内部生成告警实例并保存。我们配置一个Action只是为了满足UI的必填项并让规则可以“运行”。4.4 保存并启用填写完所有信息后保存规则并启用它。现在当你的日志索引中出现符合条件如某个服务错误日志在1分钟内超过10条时Kibana就会在后台生成告警记录。5. 故障排查与优化心得在实际部署和运行过程中我遇到了几个典型问题以下是排查思路和解决方案。5.1 代理服务无法连接到Elasticsearch现象服务日志持续报ConnectionError。排查网络连通性 在代理服务主机上使用curl -u user:password http://es-host:9200测试。认证信息 检查ES_USER和ES_PASSWORD是否正确特别是密码中的特殊字符是否需要转义。防火墙与安全组 确认9200端口对代理服务主机开放。ES集群健康状态 在Kibana Dev Tools中执行GET /_cluster/health查看状态。解决 确保网络可达并使用正确的认证凭据。对于生产环境建议使用服务账号和API密钥而非用户名密码。5.2 轮询不到新告警现象Kibana界面显示告警已触发但钉钉收不到消息代理服务日志显示“未发现新告警”。排查检查时间字段 Kibana 8.x告警实例的时间字段可能是kibana.alert.start或timestamp。需要确认.kibana-alerting-*索引中实际使用的字段名。可以在Kibana Dev Tools中查询该索引的映射GET .kibana-alerting-*/_mapping或查看一条样例数据。调整时间范围 代码中默认处理过去5分钟的数据。如果告警是更早之前触发的可能被过滤。可以临时修改load_last_processed_time函数返回一个更早的时间进行测试。确认索引权限 用于连接ES的账户是否有权读取.kibana-alerting-*索引尝试用该账户直接执行查询GET .kibana-alerting-*/_search。查看告警实例数据 在Kibana中进入Management - Stack Management - Rules and Connectors - Alerts查看触发的告警列表。确认告警确实已生成。解决 根据查询结果调整代码中的时间字段和查询逻辑。确保服务账户有足够权限。5.3 钉钉机器人收不到消息或格式错误现象 代理服务日志显示“发送钉钉消息失败”或返回错误码。排查Webhook URL与Secret 确认DINGDING_WEBHOOK和DINGDING_SECRET填写正确特别是Secret用于加签时代码中的签名算法必须与钉钉要求一致。关键词/IP安全设置 如果钉钉机器人设置了关键词则发送的消息内容中必须包含该关键词。如果设置了IP白名单代理服务出口IP必须加入白名单。消息格式 钉钉机器人对JSON格式要求严格。使用json.dumps(..., ensure_asciiFalse)确保中文正常显示。用print或日志输出最终要发送的JSON字符串在线JSON格式化工具检查其有效性。网络出口 代理服务所在机器是否能正常访问外网和oapi.dingtalk.com域名。解决 仔细核对机器人配置并在代码中增加更详细的请求和响应日志便于定位问题。5.4 性能与可靠性优化建议使用Elasticsearch的PITPoint in Time与Search After 对于大量告警的场景轮询使用from/size分页可能不高效。建议改用PIT API和search_after参数进行深分页能获得更稳定的遍历性能。引入消息队列解耦 当前方案是代理服务直接处理并发送。在高频告警场景下可能对钉钉接口造成压力或被限流。可以引入一个内部轻量级消息队列如Redis List代理服务将告警事件推入队列由另一个专门的发送器Sender从队列消费并发送实现削峰填谷。告警去重与抑制 Kibana告警规则可能短时间内连续触发。可以在代理服务层增加简单的去重逻辑例如对于同一规则同一分组rule.idgroup key的告警在5分钟内只发送一条。更复杂的抑制规则可以记录状态到Redis。完善监控 监控代理服务本身为其添加心跳、处理延迟、失败次数等指标并可以通过它自身或另一个独立的通道报告异常避免“告警系统本身挂了无人知晓”的尴尬局面。这套方案虽然绕过了Kibana白金版的许可限制但引入了一个需要自行维护的外部组件。它的优势在于灵活、成本低并且将告警逻辑在Kibana和通知逻辑在代理服务解耦未来可以很容易地扩展支持企业微信、飞书等其他通知渠道。对于资源有限但又有强告警需求的团队不失为一个务实的选择。