监控告警MCP Server:让AI查看系统状态
摘要:监控告警MCP Server开发,让AI查看系统状态和告警信息,集成Prometheus、Grafana等监控系统的MCP工具开发方案。
监控告警MCP Server 让AI查看系统状态
本文是MCP协议全栈实战专栏第54篇。标签: MCP, 监控告警, Prometheus, 系统状态, MCP Server
开头聊两句
上个月我们线上出了次故障,凌晨三点被叫醒。打开电脑看Grafana,一堆仪表板翻来翻去,CPU、内存、网络、磁盘、各个服务的错误率,眼睛都看花了。当时我就想,要是能直接问AI"现在系统有什么异常",AI自动去查Prometheus和Grafana,把异常指标汇总给我,那该多好。
第二天我就开始写监控告警MCP Server。第一版很快写完了,能查Prometheus的指标和告警。但测试的时候发现一个要命的问题,AI不知道PromQL怎么写。我问AI"查一下CPU使用率",AI生成的PromQL是cpu_usage,但Prometheus里这个指标叫node_cpu_seconds_total,而且需要用1 - avg(rate(...))来算使用率。AI直接查cpu_usage返回空结果,然后AI说"系统没有CPU监控数据"。
后来我做了两个改进。一是预置了一批常用查询模板,AI只需要说"查CPU",我就用预定义的PromQL去查。二是加了一个指标发现工具,AI可以先查有哪些指标可用,再构造查询。这两个改进让AI查监控的准确率从30%提升到了90%以上。
这篇文章就把监控告警MCP Server的完整实现写出来,包括Prometheus集成、告警查询、指标聚合和事件通知。
核心知识 监控系统对接MCP的设计
AI查监控的三个难点
第一个难点是查询语言。Prometheus用PromQL,Grafana也有自己的查询语法。AI不一定能正确生成这些查询语句。解决方案是预置常用查询模板,把复杂的PromQL封装成简单的参数化查询。
第二个难点是数据量大。Prometheus里可能有成千上万个时间序列,一次查询可能返回几万个数据点。直接全部返回给AI会超出token限制。解决方案是做数据降采样,把1分钟的精度聚合成5分钟或10分钟,只返回趋势数据。
第三个难点是告警噪音。一个生产环境可能有几百条告警规则,同时触发的告警可能有几十条。AI如果逐条分析会非常慢。解决方案是做告警聚合,按服务、按严重级别分组,先给AI一个概览,AI需要细节时再展开。
架构设计
监控告警MCP Server对接两个数据源。Prometheus提供原始指标数据和告警状态,Grafana提供仪表板视图和历史图表。
Server对外提供六个工具。指标查询工具查任意Prometheus指标,支持瞬时查询和范围查询。常用指标工具提供预置的CPU、内存、磁盘、网络等常用查询模板。告警查询工具获取当前活跃告警列表。告警规则工具列出所有告警规则及其状态。仪表板查询工具获取Grafana仪表板的数据。事件订阅工具注册告警事件通知。
数据在返回给AI之前会做两层处理。第一层是降采样,把高频数据聚合成低频趋势数据。第二层是摘要生成,对大量告警做分组汇总,只返回摘要而不是全量列表。
完整代码 监控告警MCP Server
"""
监控告警MCP Server - 让AI查看系统状态
功能: Prometheus指标查询、告警状态、Grafana仪表板、事件通知
依赖: pip install mcp httpx
"""
import asyncio
import json
import time
import logging
from typing import Any, Optional
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from collections import defaultdict
import httpx
from mcp.server import Server
from mcp.server.stdio import stdio_server
from mcp.types import Tool, TextContent
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("monitoring_mcp_server")
# ============================================================
# 第一部分: Prometheus客户端
# 封装Prometheus HTTP API的调用
# 支持瞬时查询、范围查询、告警查询、规则查询
# ============================================================
class PrometheusClient:
"""
Prometheus HTTP API客户端
封装常用的Prometheus查询接口
所有请求通过httpx异步发送
"""
def __init__(
self,
base_url: str = "http://localhost:9090",
auth_token: str = None,
timeout: int = 30
):
# Prometheus API基础URL
self.base_url = base_url.rstrip("/")
# 认证Token (如果Prometheus启用了认证)
self.auth_token = auth_token
# 请求超时时间
self.timeout = timeout
# HTTP客户端, 共享连接池
self._client = httpx.AsyncClient(
timeout=httpx.Timeout(timeout),
limits=httpx.Limits(max_connections=10)
)
async def _request(
self,
path: str,
params: dict = None
) -> dict:
"""
发送HTTP请求到Prometheus API
统一处理认证、错误和响应解析
"""
url = f"{self.base_url}{path}"
headers = {}
# 如果有认证Token, 加到请求头
if self.auth_token:
headers["Authorization"] = f"Bearer {self.auth_token}"
try:
response = await self._client.get(
url, params=params, headers=headers
)
response.raise_for_status()
return response.json()
except httpx.TimeoutException:
return {"error": "Prometheus请求超时"}
except httpx.HTTPStatusError as e:
return {
"error": f"Prometheus返回错误: {e.response.status_code}",
"detail": e.response.text[:500]
}
except Exception as e:
return {"error": f"请求失败: {str(e)}"}
async def instant_query(
self,
query: str,
time: str = None
) -> dict:
"""
瞬时查询 - 查询某个时间点的指标值
query: PromQL查询语句
time: 查询时间点(RFC3339格式), 不传则查当前
"""
params = {"query": query}
if time:
params["time"] = time
result = await self._request("/api/v1/query", params)
if "error" in result:
return result
# 解析Prometheus返回的结果
return self._parse_query_result(result)
async def range_query(
self,
query: str,
start: str,
end: str,
step: str = "60s"
) -> dict:
"""
范围查询 - 查询一段时间内的指标趋势
query: PromQL查询语句
start: 开始时间 (Unix时间戳或RFC3339)
end: 结束时间
step: 采样间隔, 如 60s, 5m, 1h
"""
params = {
"query": query,
"start": start,
"end": end,
"step": step,
}
result = await self._request("/api/v1/query_range", params)
if "error" in result:
return result
return self._parse_range_result(result)
async def get_alerts(self) -> dict:
"""
获取当前所有告警状态
返回所有正在触发的告警
"""
result = await self._request("/api/v1/alerts")
if "error" in result:
return result
# 解析告警列表
alerts = result.get("data", {}).get("alerts", [])
parsed_alerts = []
for alert in alerts:
parsed_alerts.append({
"name": alert.get("labels", {}).get("alertname", "unknown"),
"severity": alert.get("labels", {}).get("severity", "unknown"),
"state": alert.get("state", "unknown"), # firing, pending
"instance": alert.get("labels", {}).get("instance", ""),
"job": alert.get("labels", {}).get("job", ""),
"summary": alert.get("annotations", {}).get("summary", ""),
"description": alert.get("annotations", {}).get("description", ""),
"starts_at": alert.get("startsAt", ""),
"ends_at": alert.get("endsAt", ""),
"active_at": alert.get("activeAt", ""),
})
return {
"total": len(parsed_alerts),
"firing": sum(1 for a in parsed_alerts if a["state"] == "firing"),
"pending": sum(1 for a in parsed_alerts if a["state"] == "pending"),
"alerts": parsed_alerts,
}
async def get_rules(self) -> dict:
"""
获取所有告警规则
包括规则名称、表达式、状态
"""
result = await self._request("/api/v1/rules")
if "error" in result:
return result
groups = result.get("data", {}).get("groups", [])
rules = []
for group in groups:
for rule in group.get("rules", []):
if rule.get("type") == "alerting":
rules.append({
"name": rule.get("name", ""),
"query": rule.get("query", ""),
"state": rule.get("state", ""),
"group": group.get("name", ""),
"labels": rule.get("labels", {}),
"annotations": rule.get("annotations", {}),
"health": rule.get("health", ""),
"last_eval": rule.get("lastEvaluation", ""),
})
return {
"total": len(rules),
"rules": rules,
}
async def get_targets(self) -> dict:
"""
获取所有监控目标的状态
查看哪些服务正在被监控, 哪些挂了
"""
result = await self._request("/api/v1/targets")
if "error" in result:
return result
active = result.get("data", {}).get("activeTargets", [])
targets = []
for target in active:
targets.append({
"job": target.get("labels", {}).get("job", ""),
"instance": target.get("labels", {}).get("instance", ""),
"health": target.get("health", ""), # up, down
"last_error": target.get("lastError", ""),
"last_scrape": target.get("lastScrape", ""),
"scrape_duration": target.get("scrapeDuration", ""),
})
up_count = sum(1 for t in targets if t["health"] == "up")
down_count = sum(1 for t in targets if t["health"] == "down")
return {
"total": len(targets),
"up": up_count,
"down": down_count,
"targets": targets,
}
def _parse_query_result(self, result: dict) -> dict:
"""解析瞬时查询结果"""
data = result.get("data", {})
result_type = data.get("resultType", "")
results = data.get("result", [])
if result_type == "vector":
# 瞬时向量, 每个结果包含标签和值
values = []
for item in results:
values.append({
"labels": item.get("metric", {}),
"value": float(item.get("value", [0, "0"])[1]),
"timestamp": item.get("value", [0, "0"])[0],
})
return {"type": "vector", "results": values}
elif result_type == "matrix":
# 矩阵, 多个时间序列
return {"type": "matrix", "results": results}
elif result_type == "scalar":
# 标量值
value = data.get("result", [0, "0"])
return {"type": "scalar", "value": float(value[1])}
return {"type": result_type, "results": results}
def _parse_range_result(self, result: dict) -> dict:
"""
解析范围查询结果
对数据做降采样, 减少返回给AI的数据量
"""
data = result.get("data", {})
results = data.get("result", [])
parsed = []
for series in results:
labels = series.get("metric", {})
values = series.get("values", [])
# 降采样: 如果数据点超过60个, 做聚合
if len(values) > 60:
values = self._downsample(values, target_count=60)
# 把时间戳转成可读时间
formatted_values = []
for ts, val in values:
formatted_values.append({
"time": datetime.fromtimestamp(float(ts)).strftime(
"%H:%M:%S"
),
"value": float(val),
})
parsed.append({
"labels": labels,
"point_count": len(formatted_values),
"values": formatted_values,
})
return {"series_count": len(parsed), "results": parsed}
def _downsample(
self,
values: list,
target_count: int = 60
) -> list:
"""
降采样: 把大量数据点减少到目标数量
用简单的等间隔采样, 每N个点取一个
"""
if len(values) <= target_count:
return values
# 计算采样间隔
step = len(values) // target_count
# 等间隔采样
sampled = values[::step][:target_count]
return sampled
async def close(self):
"""关闭HTTP客户端"""
await self._client.aclose()
# ============================================================
# 第二部分: 预置查询模板
# 封装常用的PromQL查询, AI不需要自己写PromQL
# ============================================================
class QueryTemplates:
"""
预置查询模板
把复杂的PromQL封装成简单的参数化查询
AI只需要说"查CPU使用率", 不用知道PromQL怎么写
"""
# 模板定义, key是模板名, value是查询配置
TEMPLATES = {
# CPU使用率 - 用irate计算每秒变化率
"cpu_usage": {
"description": "查询CPU使用率, 返回0-1之间的值",
"query": '1 - avg(rate(node_cpu_seconds_total{{mode="idle"}}[5m])) by (instance)',
"unit": "ratio",
"type": "instant",
},
# 内存使用率
"memory_usage": {
"description": "查询内存使用率",
"query": '1 - (node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes)',
"unit": "ratio",
"type": "instant",
},
# 磁盘使用率
"disk_usage": {
"description": "查询磁盘空间使用率",
"query": '1 - (node_filesystem_avail_bytes{{fstype!~"tmpfs|overlay"}} / node_filesystem_size_bytes{{fstype!~"tmpfs|overlay"}})',
"unit": "ratio",
"type": "instant",
},
# 网络入流量
"network_rx": {
"description": "查询网络接收流量(每秒字节数)",
"query": 'rate(node_network_receive_bytes_total{{device!~"lo|veth.*"}}[5m])',
"unit": "bytes/sec",
"type": "instant",
},
# HTTP请求率
"http_rate": {
"description": "查询HTTP请求速率(每秒请求数)",
"query": 'sum(rate(http_requests_total[5m])) by (handler)',
"unit": "req/sec",
"type": "instant",
},
# HTTP错误率
"http_error_rate": {
"description": "查询HTTP 5xx错误率",
"query": 'sum(rate(http_requests_total{{status=~"5.."}}[5m])) / sum(rate(http_requests_total[5m]))',
"unit": "ratio",
"type": "instant",
},
# 服务状态
"service_up": {
"description": "查询服务是否在线(up=1, down=0)",
"query": 'up',
"unit": "boolean",
"type": "instant",
},
}
@classmethod
def get_template(cls, name: str) -> Optional[dict]:
"""获取指定名称的查询模板"""
return cls.TEMPLATES.get(name)
@classmethod
def list_templates(cls) -> list:
"""列出所有可用的查询模板"""
return [
{
"name": name,
"description": tpl["description"],
"unit": tpl["unit"],
"type": tpl["type"],
}
for name, tpl in cls.TEMPLATES.items()
]
# ============================================================
# 第三部分: Grafana客户端
# 查询Grafana仪表板数据
# ============================================================
class GrafanaClient:
"""
Grafana API客户端
查询仪表板列表和仪表板数据
"""
def __init__(
self,
base_url: str = "http://localhost:3000",
api_key: str = None,
timeout: int = 30
):
self.base_url = base_url.rstrip("/")
self.api_key = api_key
self.timeout = timeout
self._client = httpx.AsyncClient(
timeout=httpx.Timeout(timeout),
limits=httpx.Limits(max_connections=5)
)
async def _request(self, path: str) -> dict:
"""发送请求到Grafana API"""
url = f"{self.base_url}{path}"
headers = {}
if self.api_key:
headers["Authorization"] = f"Bearer {self.api_key}"
try:
response = await self._client.get(url, headers=headers)
response.raise_for_status()
return response.json()
except Exception as e:
return {"error": str(e)}
async def list_dashboards(self) -> dict:
"""获取所有仪表板列表"""
result = await self._request("/api/search?type=dash-db")
if "error" in result:
return result
dashboards = []
for item in result:
dashboards.append({
"id": item.get("id"),
"uid": item.get("uid"),
"title": item.get("title"),
"url": item.get("url"),
"tags": item.get("tags", []),
"folder": item.get("folderTitle", ""),
})
return {"total": len(dashboards), "dashboards": dashboards}
async def get_dashboard(self, uid: str) -> dict:
"""获取指定仪表板的详细信息"""
result = await self._request(f"/api/dashboards/uid/{uid}")
if "error" in result:
return result
dashboard = result.get("dashboard", {})
panels = dashboard.get("panels", [])
# 提取面板信息, 不返回完整数据(太大)
panel_list = []
for panel in panels:
panel_list.append({
"id": panel.get("id"),
"title": panel.get("title"),
"type": panel.get("type"),
"datasource": panel.get("datasource"),
"queries": [
t.get("expr", "") for t in panel.get("targets", [])
if t.get("expr")
],
})
return {
"title": dashboard.get("title"),
"tags": dashboard.get("tags", []),
"panel_count": len(panel_list),
"panels": panel_list,
}
async def close(self):
await self._client.aclose()
# ============================================================
# 第四部分: 事件通知管理器
# 管理告警事件订阅和通知
# ============================================================
class EventNotifier:
"""
事件通知管理器
允许AI订阅特定类型的告警事件
当告警状态变化时触发通知
"""
def __init__(self):
# 订阅列表, 每个订阅记录关注条件
self._subscriptions: list[dict] = []
# 事件历史, 存储最近的通知
self._event_history: list[dict] = []
# 最大历史记录数
self._max_history = 100
def subscribe(
self,
name: str,
severity: str = None,
service: str = None,
callback_url: str = None
) -> dict:
"""
订阅告警事件
severity: 只关注特定级别的告警 (critical/warning/info)
service: 只关注特定服务的告警
callback_url: 事件触发时回调的URL
"""
sub = {
"id": f"sub_{len(self._subscriptions) + 1}",
"name": name,
"severity": severity,
"service": service,
"callback_url": callback_url,
"created_at": datetime.now().isoformat(),
"triggered_count": 0,
}
self._subscriptions.append(sub)
logger.info(f"事件订阅创建: {sub['id']} ({name})")
return sub
def unsubscribe(self, sub_id: str) -> bool:
"""取消订阅"""
before = len(self._subscriptions)
self._subscriptions = [
s for s in self._subscriptions if s["id"] != sub_id
]
return len(self._subscriptions) < before
def check_and_notify(self, alerts: list) -> list:
"""
检查告警列表, 对匹配订阅的告警生成通知
返回触发的通知列表
"""
notifications = []
for alert in alerts:
for sub in self._subscriptions:
# 检查是否匹配订阅条件
matched = True
if sub.get("severity"):
if alert.get("severity") != sub["severity"]:
matched = False
if sub.get("service"):
if sub["service"] not in alert.get("job", ""):
matched = False
if matched:
notification = {
"subscription_id": sub["id"],
"subscription_name": sub["name"],
"alert": alert,
"notified_at": datetime.now().isoformat(),
}
notifications.append(notification)
sub["triggered_count"] += 1
self._add_to_history(notification)
return notifications
def _add_to_history(self, event: dict):
"""添加事件到历史记录"""
self._event_history.append(event)
# 保持历史记录在最大数量内
if len(self._event_history) > self._max_history:
self._event_history = self._event_history[-self._max_history:]
def get_history(self, limit: int = 20) -> list:
"""获取最近的事件通知历史"""
return self._event_history[-limit:]
def list_subscriptions(self) -> list:
"""列出所有订阅"""
return self._subscriptions
# ============================================================
# 第五部分: MCP Server主程序
# 整合Prometheus、Grafana和事件通知
# ============================================================
class MonitoringMCPServer:
"""
监控告警MCP Server
工具: 常用指标查询、自定义PromQL查询、告警状态、
告警规则、监控目标、仪表板查询、事件订阅
"""
def __init__(
self,
prometheus_url: str = "http://localhost:9090",
grafana_url: str = "http://localhost:3000",
grafana_api_key: str = None,
prometheus_token: str = None
):
# 初始化各客户端
self.prometheus = PrometheusClient(
base_url=prometheus_url,
auth_token=prometheus_token
)
self.grafana = GrafanaClient(
base_url=grafana_url,
api_key=grafana_api_key
)
self.notifier = EventNotifier()
# MCP Server
self.server = Server("monitoring-server")
self._setup_handlers()
def _setup_handlers(self):
"""注册MCP工具"""
@self.server.list_tools()
async def handle_list_tools() -> list[Tool]:
return [
Tool(
name="mon_metrics",
description=(
"查询常用监控指标。提供预置的指标查询模板, "
"不需要写PromQL。可选模板: cpu_usage, "
"memory_usage, disk_usage, network_rx, "
"http_rate, http_error_rate, service_up。"
),
inputSchema={
"type": "object",
"properties": {
"metric": {
"type": "string",
"description": "指标名称, 如 cpu_usage"
}
},
"required": ["metric"]
}
),
Tool(
name="mon_query",
description=(
"执行自定义PromQL查询。"
"如果你知道PromQL语法可以用这个工具查询任意指标。"
"返回瞬时查询结果。"
),
inputSchema={
"type": "object",
"properties": {
"query": {
"type": "string",
"description": "PromQL查询语句"
}
},
"required": ["query"]
}
),
Tool(
name="mon_query_range",
description=(
"执行范围查询, 获取一段时间内的指标趋势。"
"返回降采样后的时间序列数据。"
),
inputSchema={
"type": "object",
"properties": {
"query": {
"type": "string",
"description": "PromQL查询语句"
},
"duration": {
"type": "string",
"description": "时间范围, 如 1h, 6h, 24h",
"default": "1h"
},
"step": {
"type": "string",
"description": "采样间隔, 如 60s, 5m",
"default": "60s"
}
},
"required": ["query"]
}
),
Tool(
name="mon_alerts",
description=(
"查询当前所有活跃告警。"
"返回告警名称、级别、状态和描述。"
"按严重级别排序。"
),
inputSchema={
"type": "object",
"properties": {}
}
),
Tool(
name="mon_rules",
description=(
"列出所有告警规则及其当前状态。"
),
inputSchema={
"type": "object",
"properties": {}
}
),
Tool(
name="mon_targets",
description=(
"查看监控目标状态, 哪些服务在线, 哪些离线。"
),
inputSchema={
"type": "object",
"properties": {}
}
),
Tool(
name="mon_dashboards",
description=(
"查询Grafana仪表板列表或指定仪表板的详情。"
"不传uid则返回所有仪表板列表。"
),
inputSchema={
"type": "object",
"properties": {
"uid": {
"type": "string",
"description": "仪表板UID, 不传则列全部"
}
}
}
),
Tool(
name="mon_subscribe",
description=(
"订阅告警事件通知。当匹配条件的告警触发时会收到通知。"
"可以按严重级别和服务名过滤。"
),
inputSchema={
"type": "object",
"properties": {
"name": {
"type": "string",
"description": "订阅名称"
},
"severity": {
"type": "string",
"description": "告警级别: critical/warning"
},
"service": {
"type": "string",
"description": "服务名过滤"
}
},
"required": ["name"]
}
),
]
@self.server.call_tool()
async def handle_call_tool(
name: str,
arguments: dict
) -> list[TextContent]:
if name == "mon_metrics":
result = await self._handle_metrics(arguments)
elif name == "mon_query":
result = await self.prometheus.instant_query(
arguments["query"]
)
elif name == "mon_query_range":
result = await self._handle_range_query(arguments)
elif name == "mon_alerts":
result = await self._handle_alerts()
elif name == "mon_rules":
result = await self.prometheus.get_rules()
elif name == "mon_targets":
result = await self.prometheus.get_targets()
elif name == "mon_dashboards":
result = await self._handle_dashboards(arguments)
elif name == "mon_subscribe":
result = self._handle_subscribe(arguments)
else:
result = {"error": f"未知工具: {name}"}
return [TextContent(
type="text",
text=json.dumps(result, ensure_ascii=False, indent=2, default=str)
)]
async def _handle_metrics(self, arguments: dict) -> dict:
"""处理常用指标查询"""
metric_name = arguments.get("metric", "")
# 查找预置模板
template = QueryTemplates.get_template(metric_name)
if not template:
return {
"error": f"未知的指标模板: {metric_name}",
"available": QueryTemplates.list_templates()
}
# 执行查询
result = await self.prometheus.instant_query(template["query"])
if "error" in result:
return result
# 附加模板信息
result["metric"] = metric_name
result["description"] = template["description"]
result["unit"] = template["unit"]
return result
async def _handle_range_query(self, arguments: dict) -> dict:
"""处理范围查询"""
query = arguments["query"]
duration = arguments.get("duration", "1h")
step = arguments.get("step", "60s")
# 把duration转成开始和结束时间
now = datetime.now()
end_time = now.timestamp()
# 解析duration字符串 (如 1h, 6h, 24h)
duration_seconds = self._parse_duration(duration)
start_time = end_time - duration_seconds
result = await self.prometheus.range_query(
query=query,
start=str(start_time),
end=str(end_time),
step=step
)
if "error" not in result:
result["duration"] = duration
result["step"] = step
return result
def _parse_duration(self, duration: str) -> float:
"""把时间字符串转成秒数"""
# 支持 1h, 6h, 24h, 30m, 1d 等格式
if duration.endswith("h"):
return float(duration[:-1]) * 3600
elif duration.endswith("m"):
return float(duration[:-1]) * 60
elif duration.endswith("d"):
return float(duration[:-1]) * 86400
elif duration.endswith("s"):
return float(duration[:-1])
else:
return 3600 # 默认1小时
async def _handle_alerts(self) -> dict:
"""处理告警查询, 做聚合和排序"""
result = await self.prometheus.get_alerts()
if "error" in result:
return result
alerts = result.get("alerts", [])
# 按严重级别排序: critical > warning > info
severity_order = {"critical": 0, "warning": 1, "info": 2}
alerts.sort(
key=lambda a: severity_order.get(a.get("severity", ""), 3)
)
# 按服务分组汇总
service_summary = defaultdict(lambda: {"critical": 0, "warning": 0, "info": 0})
for alert in alerts:
service = alert.get("job", "unknown")
severity = alert.get("severity", "info")
service_summary[service][severity] += 1
# 检查事件订阅
notifications = self.notifier.check_and_notify(alerts)
return {
"total": len(alerts),
"firing": result.get("firing", 0),
"pending": result.get("pending", 0),
"service_summary": dict(service_summary),
"notifications_triggered": len(notifications),
"alerts": alerts,
}
async def _handle_dashboards(self, arguments: dict) -> dict:
"""处理仪表板查询"""
uid = arguments.get("uid")
if uid:
return await self.grafana.get_dashboard(uid)
else:
return await self.grafana.list_dashboards()
def _handle_subscribe(self, arguments: dict) -> dict:
"""处理事件订阅"""
sub = self.notifier.subscribe(
name=arguments["name"],
severity=arguments.get("severity"),
service=arguments.get("service"),
)
return {"subscription": sub}
async def run(self):
"""启动MCP Server"""
logger.info("监控告警MCP Server启动")
async with stdio_server() as (read_stream, write_stream):
await self.server.run(
read_stream,
write_stream,
self.server.create_initialization_options()
)
async def shutdown(self):
"""关闭Server"""
await self.prometheus.close()
await self.grafana.close()
logger.info("监控告警MCP Server已关闭")
# ============================================================
# 入口
# ============================================================
async def main():
"""
主函数
配置Prometheus和Grafana地址, 启动MCP Server
"""
server = MonitoringMCPServer(
prometheus_url="http://localhost:9090",
grafana_url="http://localhost:3000",
grafana_api_key="your_grafana_api_key", # 从环境变量读取
prometheus_token=None, # 如果Prometheus没有认证则不传
)
try:
await server.run()
except KeyboardInterrupt:
pass
finally:
await server.shutdown()
if __name__ == "__main__":
asyncio.run(main())
对比分析 监控对接方案对比
| 方案 | 查询能力 | 开发成本 | AI使用难度 | 实时性 | 适用场景 |
|---|---|---|---|---|---|
| 直接调Prometheus API | 完整PromQL | 低 | 高(需懂PromQL) | 实时 | 运维工程师 |
| 预置模板+自定义查询(本文) | 模板+原生PromQL | 中 | 低(模板优先) | 实时 | AI助手 |
| Grafana截图给AI看 | 不可查询 | 低 | 极低 | 延迟 | 简单场景 |
| 商业AIOps平台 | 智能分析 | 低(配置) | 极低 | 实时 | 大企业 |
| 日志直接搜ELK | 日志级 | 中 | 中 | 近实时 | 故障排查 |
我的方案在预置模板和原生PromQL之间做了平衡。常用指标用模板,AI不需要写PromQL。特殊需求用自定义查询,AI可以写PromQL。这种"模板优先,自定义兜底"的策略适合大部分场景。
踩坑经验 范围查询返回海量数据把AI撑爆
这个坑发生在我让AI查"过去24小时的CPU趋势"时。
AI调用了mon_query_range工具,传了duration=24h和step=60s。Prometheus返回了24小时的数据,每分钟一个点,共1440个数据点。每个数据点包含时间戳和值,JSON序列化后约30KB。
问题是有4台服务器,返回了4条时间序列,总共5760个数据点。JSON大小超过120KB。Claude Desktop收到后虽然没报错,但AI处理这么大的数据时明显变慢了,回复花了将近30秒。
更夸张的是,如果查的是多维度指标(比如按pod分组的CPU使用率),一个查询可能返回几十条时间序列,数据量直接爆炸。
我的解决方案是在_parse_range_result方法里做降采样。不管原始数据有多少点,统一降到60个点以内。24小时的数据变成每24分钟一个点,足够看趋势了。如果AI需要更精细的数据,可以缩短查询范围,比如只查过去1小时。
降采样后的对比:
| 查询范围 | 原始数据点数 | 降采样后 | 原始JSON大小 | 降采样后大小 |
|---|---|---|---|---|
| 1小时(60s步长) | 60 | 60(不变) | 2KB | 2KB |
| 6小时(60s步长) | 360 | 60 | 12KB | 2KB |
| 24小时(60s步长) | 1440 | 60 | 48KB | 2KB |
| 24小时多序列(4条) | 5760 | 240 | 120KB | 8KB |
| 7天(5m步长) | 2016 | 60 | 67KB | 2KB |
降采样用的是最简单的等间隔采样,每N个点取一个。这种方法的缺点是可能跳过峰值。如果AI需要找峰值,可以查瞬时值或者用max_over_time函数。更好的方案是用LTTB(Largest-Triangle-Three-Buckets)算法做降采样,能保留视觉特征,但实现复杂度高,对于AI分析趋势来说等间隔够用了。
还有一个关于告警的坑。Prometheus的/api/v1/alerts接口返回的是当前活跃的告警,不包括已经恢复的告警。有一次AI查告警发现没有异常,但实际上5分钟前有一波告警已经自动恢复了。如果需要看历史告警,要查Alertmanager的API而不是Prometheus的。后来我加了一个mon_subscribe工具,AI可以订阅告警事件,这样即使告警恢复了,订阅历史里也有记录。
小结
这篇写了监控告警MCP Server的完整实现。六个核心工具覆盖了AI查看系统状态的主要需求,常用指标查询用预置模板,自定义查询支持原生PromQL,告警查询做分组聚合,仪表板查询对接Grafana,事件订阅支持告警通知。
预置查询模板是这套方案的关键创新。AI不需要学PromQL就能查常用指标,大幅降低了使用门槛。如果AI需要查模板没覆盖的指标,可以用自定义查询工具写PromQL。
降采样是处理范围查询数据爆炸的核心手段。不管原始数据多大,统一降到60个点以内,保证返回给AI的数据量可控。代价是丢失了一些细节,但对于趋势分析来说够用了。
这个系列写到这篇,我们从架构设计到各个具体的MCP Server实现都覆盖了。从企业级的工具网关,到数据库、文件系统、API网关、Git仓库、监控告警,一套完整的MCP工具体系就搭起来了。希望这个专栏能帮到正在做MCP开发的你。
相关推荐
更多推荐

所有评论(0)