MCP监控和自动诊断系统
实时监控、异常检测和自动修复 Apache Kafka·Apache Spark·HDFS大数据管道
   ](https://docker.com)   
不需要人为干预。 该系统检测异常,诊断根本原因,并自动修复。 通过内置的AI聊天模式,用简单的英语问它任何问题。
______________________________________________________________________
概述
这个项目是 生产级监控和自动诊断平台 专为大数据管道而建。它模拟了一个真实世界的环境 Apache Kafka, Apache Spark,以及 HDFS --所有这些都在不断地将实时指标导出到 普罗米修斯,可视化于 格拉法纳,并由智能保护 MCP(模型上下文协议)服务器 即:
- 接收 通过webhook从Alertmanager发出警报
- 诊断 使用YAML Runbook自动查找根本原因
- 补救措施 通过Docker API执行修复剧本的问题
- 日志 采取的每一项行动都有完整的审计历史记录
- 暴露AI代理 用于自然语言集群管理
这 修正代理CLI 提供了一个交互式终端界面,用于实时监控、手动干预、指标查询、人工智能驱动的诊断和对话式人工智能聊天,使其成为实时演示或操作屏幕的理想选择。
______________________________________________________________________
系统架构
┌──────────────────────────────────────────────────────────────────────┐
│ Simulated Big Data Pipeline │
│ │
│ ┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐ │
│ │ Kafka Exporter │ │ Spark Exporter │ │ HDFS Exporter │ │
│ │ (Python) │ │ (Python) │ │ (Python) │ │
│ │ port :8001 │ │ port :8002 │ │ port :8003 │ │
│ └────────┬────────┘ └────────┬────────┘ └────────┬────────┘ │
└────────────┼────────────────────┼────────────────────┼───────────────┘
│ /metrics │ /metrics │ /metrics
▼ ▼ ▼
┌──────────────────────────────────────────────────────────────────────┐
│ Prometheus :9090 │
│ • Scrapes all exporters every 15 seconds │
│ • Evaluates 12 alert rules continuously │
│ • Stores time-series metric data │
└──────────────────────────┬───────────────────────────────────────────┘
│ Alert fired (threshold breached)
▼
┌──────────────────────────────────────────────────────────────────────┐
│ Alertmanager :9093 │
│ • Groups and deduplicates alerts │
│ • Routes by severity (critical / warning) │
│ • Sends webhook POST to MCP Server │
└──────────────────────────┬───────────────────────────────────────────┘
│ POST /webhook/alert
▼
┌──────────────────────────────────────────────────────────────────────┐
│ MCP Server :8888 │
│ • FastAPI application (Model Context Protocol) │
│ • Receives and stores all firing/resolved alerts │
│ • Auto-triggers remediation engine per alert action label │
│ • 11 YAML runbooks for structured diagnosis │
│ • Exposes 11 REST tools for agents and automation │
│ • Exposes its own /metrics endpoint (scraped by Prometheus) │
└──────────────┬────────────────────────────┬─────────────────────────┘
│ │
▼ ▼
┌──────────────────────┐ ┌─────────────────────────────────┐
│ Grafana :3000 │ │ Remediation Agent CLI │
│ │ │ │
│ • 11-panel dashboard│ │ • Service health table │
│ • Live time-series │ │ • Active alerts view │
│ • Color thresholds │ │ • Remediation history │
│ • Auto-provisioned │ │ • Prometheus metric query │
│ │ │ • Manual trigger │
│ │ │ • Live monitor (10s refresh) │
│ │ │ • AI diagnosis (option 7) │
│ │ │ • AI chat mode (option 8) │
└──────────────────────┘ └─────────────────────────────────┘______________________________________________________________________
主要特点
完全Docker化的堆栈
一 docker compose up 命令启动整个基础架构:Prometheus、Grafana、Alertmanager、Node Exporter、cAdvisor和所有3个自定义度量导出器。
实时Grafana仪表板(11个面板)
- Kafka消费者滞后(带尖峰检测的时间序列)
- Kafka Broker状态(向上/向下统计面板)
- Kafka每秒消息数(速率图)
- Spark活动任务(带阈值的统计)
- 火花内存使用率%(仪表,0–100%)
- Spark失败作业(关键统计数据)
- HDFS磁盘使用率%(仪表,0-100%)
- HDFS数据节点状态(向上/向下统计)
- 活动警报表(实时PromQL)
- 主机CPU使用率(时间序列)
- 主机内存使用情况(时间序列)
12普罗米修斯警报规则
涵盖Kafka、Spark、HDFS和主机级异常——每个异常都有一个 severity 标签和特定 action 修复引擎使用的标签。
配备11个REST工具的MCP服务器
实现模型上下文协议模式的自定义FastAPI应用程序。它充当大脑——接收警报、运行修复剧本,并暴露任何人工智能代理或自动化都可以调用的监控工具。
11个YAML运行手册
每个警报都有一个结构化的runbook,包括:症状描述、诊断步骤和具有安全分类的补救措施。AI代理阅读这些内容以解释问题并采取行动。
自动修复引擎
当警报触发时,MCP服务器立即通过Docker API运行特定的修复行动手册:
- Kafka延迟尖峰 → 重新启动消费者群体/扩大消费者规模
- Spark作业失败 → 从上一个检查点重试
- 火花存储器高 → 增加执行器内存配置
- HDFS磁盘已满 → 清理旧的临时文件和档案
- 数据节点关闭 → 重新启动DataNode并触发复制
- 经纪人下跌 → 重启Kafka代理并重新分配分区
人工智能诊断(选项7)
读取所有活动警报和匹配的运行手册,然后将其发送到LLM(通过OpenRouter发送Qwen)进行结构化分析:最关键的问题、根本原因、建议的行动、业务影响和级联风险。
AI聊天模式(选项8)
由实时MCP工具数据支持的完整对话式自然语言界面。用简单的英语提问、获得建议和触发补救措施。
33 Pytest测试
跨端点的完整测试覆盖率、警报规则验证和runbook完整性。
完全可观察性循环
MCP服务器本身导出Prometheus指标(mcp_alerts_received_total, mcp_remediations_triggered_total, mcp_remediation_duration_seconds)--这样你就可以监控监控系统了。
______________________________________________________________________
快速开始
先决条件
- 已安装并正在运行
- python 3.11+
- macOS或Linux(在macOS苹果Silicon M2上测试)
1.克隆存储库
git clone https://github.com/AbdAllAh950/mcp-monitoring-project.git
cd mcp-monitoring-project2.安装Python依赖项
make setup3.启动完整的Docker监控栈
cd monitoring
docker compose up -d --build第一次运行需要3-5分钟来提取图像并建立导出器。
4.启动MCP服务器
# Open a new terminal tab — keep this running
cd mcp-server
python3 mcp_server.py等待: Uvicorn running on http://0.0.0.0:8888
5.启动修正代理CLI
# Open another new terminal tab
cd remediation-agent
python3 agent.py6.打开仪表板
make demo| 服务 | URL | 登录 |
|---|---|---|
| Grafana仪表板 | http://localhost:3000 | admin / admin123 |
| 普罗米修斯UI | http://localhost:9090 | — |
| MCP服务器API文档 | http://localhost:8888/docs | — |
| 警报管理器UI | http://localhost:9093 | — |
| Kafka度量 | http://localhost:8001/metrics | — |
| Spark指标 | http://localhost:8002/metrics | — |
| HDFS指标 | http://localhost:8003/metrics | — |
______________________________________________________________________
警报规则参考
| 警报 | 服务 | 严重性 | 触发条件 | 自动修复 |
|---|---|---|---|---|
KafkaConsumerLagHigh | Kafka | 警告 | lag > 5,000 1分钟 | 重新启动消费者组 |
KafkaConsumerLagCritical | 卡夫卡 | 批判 | lag > 15,000 持续2分钟 | 扩展消费者实例 |
KafkaBrokerDown | 卡夫卡 | 批判 | broker_up == 0 针对30s | 重新启动Kafka代理程序(Docker API) |
KafkaUnderReplicatedPartitions | Kafka | 警告 | under_replicated > 0 1分钟 | 检查复制因子 |
SparkJobFailed | Spark | 关键 | failed_jobs > 0 立即 | 从检查点重试作业 |
SparkExecutorMemoryHigh | 火花 | 警告 | memory > 85% 持续2分钟 | 增加执行器内存 |
SparkActiveTasksLow | 火花 | 警告 | active_tasks 80% 2分钟 | 清理旧文件 |
HDFSDataNodeDown | HDFS | 关键 | datanode_up == 0 1分钟 | 重新启动DataNode(Docker API) |
HDFSReplicationLow | HDFS | 警告 | under_replicated > 100 持续2分钟 | 触发复制恢复 |
HighCPUUsage | 系统 | 警告 | cpu > 80% 2分钟 | 调查高CPU进程 |
HighMemoryUsage | 系统 | 警告 | memory > 85% 持续2分钟 | 检查内存泄漏 |
______________________________________________________________________
MCP服务器API参考
所有工具均可在 http://localhost:8888/docs (Swagger用户界面):
| 工具 | 方法 | 端点 | 描述 |
|---|---|---|---|
| 健康检查 | GET | /health | 服务器活性检查 |
| 警报Webhook | POST | /webhook/alert | 接收Alertmanager webhooks |
| 活动警报 | GET | /tools/get_active_alerts | 当前所有发射警报 |
| 服务健康 | GET | /tools/get_service_health | 每项服务的健康摘要 |
| 查询Prometheus | POST | /tools/query_prometheus | 执行任何PromQL查询 |
| 获取指标 | GET | /tools/get_metrics | 当前指标快照 |
| 补救历史 | GET | /tools/get_remediation_history | 所有操作的完整审计日志 |
| 触发补救 | POST | /tools/trigger_remediation | 手动触发修复 |
| 列出Runbook | GET | /tools/list_runbooks | 列出所有11本Runbook |
| 获取Runbook | GET | /tools/get_runbook | 获取特定警报的runbook |
| MCP自身指标 | GET | /metrics | MCP服务器的Prometheus指标 |
示例:触发演示的手动警报
curl -X POST http://localhost:8888/webhook/alert \
-H "Content-Type: application/json" \
-d '{
"receiver": "mcp-webhook",
"status": "firing",
"alerts": [{
"status": "firing",
"labels": {
"alertname": "HDFSDataNodeDown",
"severity": "critical",
"service": "hdfs",
"action": "restart_datanode"
},
"annotations": {
"summary": "HDFS DataNode is DOWN",
"description": "DataNode unreachable for 1 minute"
}
}],
"groupLabels": {},
"commonLabels": {},
"commonAnnotations": {},
"externalURL": ""
}'______________________________________________________________________
AI代理功能
Remediation Agent CLI包括两种基于OpenRouter构建的AI驱动模式(兼容Qwen/GPT-4o)。
选项7-AI诊断
读取所有活动警报及其匹配的运行手册,然后生成结构化的LLM分析:
- 确定的最关键问题
- 可能的根本原因已得到解释
- 建议采取确切的补救措施
- 未解决的业务影响
- 级联风险警告
选项8——AI聊天模式
全对话式自然语言界面。AI可以访问每条消息上的实时MCP工具数据。
You: What alerts are firing right now?
AI: 2 active alerts:
- KafkaBrokerDown [critical]: Broker unreachable for 30 seconds
- KafkaConsumerLagHigh [warning]: Consumer group lagging 14,123 on topic transactions
Data source: GET /tools/get_active_alerts
You: What should I do about the kafka broker?
AI: 1. Verify broker status via kafka_broker_up metric
2. Review remediation history — restart was already attempted
3. Trigger restart_broker via /tools/trigger_remediation
4. Follow KafkaBrokerDown runbook for deep diagnostics
Data source: GET /tools/get_active_alerts, get_service_health, get_remediation_history
You: Fix it
AI: Triggered restart_broker for KafkaBrokerDown
Status: executed
Action taken: POST /tools/trigger_remediationAI功能设置
创建 remediation-agent/.env:
cp remediation-agent/.env.example remediation-agent/.env
# Add your OpenRouter API key — free at https://openrouter.ai内容:
OPENAI_API_KEY=sk-or-v1-...
OPENAI_BASE_URL=https://openrouter.ai/api/v1
OPENAI_MODEL=qwen/qwen3-8b______________________________________________________________________
自动修复的工作原理
Prometheus detects: kafka_consumer_lag > 5000 for 1 minute
↓
Alertmanager fires: KafkaConsumerLagHigh (severity=warning, action=restart_consumer)
↓
MCP Server receives POST /webhook/alert
↓
Runbook lookup: KafkaConsumerLagHigh → symptom, diagnosis steps, safe actions
↓
Remediation Engine reads alert.labels.action = "restart_consumer"
↓
Runs playbook:
Step 1: Detect affected consumer group via Prometheus query
Step 2: Pause consumer group temporarily
Step 3: Reset consumer offset to latest checkpoint
Step 4: Restart consumer group
Step 5: Verify lag is decreasing
↓
Result logged: { success: true, duration: 2.0s, steps: [...] }
↓
Alert auto-resolves when lag drops below 5000______________________________________________________________________
运行测试
# Start MCP server first, then:
python3 -m pytest tests/ -v3个模块的33个测试全部通过:
| 模块 | 测试 | 它涵盖了什么 |
|---|---|---|
test_alert_rules.py | 12 | 每个警报都有严重性、服务、操作、表达式和摘要 |
test_mcp_endpoints.py | 11 | 健康、webhook、所有9个工具、Prometheus格式 |
test_runbook_coverage.py | 10 | 每个runbook都有症状、诊断步骤和安全措施 |
______________________________________________________________________
现场演示指南
# Tab 1: Docker stack already running in background
# Tab 2: MCP Server
cd mcp-server && python3 mcp_server.py
# Tab 3: Agent
cd remediation-agent && python3 agent.py
# Press 6 for live auto-refresh monitor
# Press 7 for AI diagnosis
# Press 8 for AI chat mode一次打开所有演示选项卡:
make demo
# Opens: Grafana, Prometheus /alerts, MCP API /docs, Alertmanager触发实时演示警报:
curl -X POST http://localhost:8888/webhook/alert \
-H "Content-Type: application/json" \
-d '{"receiver":"mcp-webhook","status":"firing",
"alerts":[{"status":"firing",
"labels":{"alertname":"HDFSDataNodeDown","severity":"critical",
"service":"hdfs","action":"restart_datanode"},
"annotations":{"summary":"HDFS DataNode is DOWN",
"description":"DataNode unreachable for 1 minute"}}],
"groupLabels":{},"commonLabels":{},"commonAnnotations":{},"externalURL":""}'观看实时监视器:HDFS从健康状态切换到危急状态,修复运行,状态显示完成——所有这些都在10秒内完成。
______________________________________________________________________
项目结构
mcp-monitoring-project/
│
├── Makefile # Control panel: setup/start/stop/demo
├── README.md
├── .gitignore
│
├── monitoring/ # Full Docker monitoring stack
│ ├── docker-compose.yml # 8 services in one file
│ ├── prometheus/
│ │ ├── prometheus.yml # Scrape configs for all targets
│ │ ├── alert_rules.yml # 12 alert rules (Kafka/Spark/HDFS/System)
│ │ └── alertmanager.yml # Webhook routing to MCP Server
│ └── grafana/
│ ├── provisioning/
│ │ ├── datasources/prometheus.yml # Auto-connects Prometheus datasource
│ │ └── dashboards/dashboards.yml # Auto-loads dashboard on startup
│ └── dashboards/
│ └── mcp-monitoring.json # 11-panel live dashboard definition
│
├── exporters/ # Metric simulators
│ ├── kafka_exporter.py # Kafka: lag, broker, messages, partitions
│ ├── Dockerfile.kafka
│ ├── spark_exporter.py # Spark: tasks, memory, jobs, executors
│ ├── Dockerfile.spark
│ ├── hdfs_exporter.py # HDFS: disk, datanodes, blocks, files
│ └── Dockerfile.hdfs
│
├── mcp-server/ # MCP Server (the brain)
│ ├── mcp_server.py # FastAPI app: webhook + tools + remediation
│ ├── runbooks.yaml # 11 structured runbooks for every alert
│ └── requirements.txt
│
├── remediation-agent/ # Interactive CLI agent
│ ├── agent.py # Rich terminal UI with AI chat (options 1-8)
│ ├── .env.example # Template for OpenRouter API key
│ └── requirements.txt
│
└── tests/ # 33 pytest tests
├── conftest.py
├── test_alert_rules.py # 12 tests: alert rule validation
├── test_mcp_endpoints.py # 11 tests: REST API endpoints
└── test_runbook_coverage.py # 10 tests: runbook completeness______________________________________________________________________
Makefile命令
make setup # Install all Python dependencies for MCP server and agent
make start # Start Docker stack + MCP server
make stop # Stop all services cleanly
make restart # Stop then start everything
make logs # Follow all Docker container logs live
make mcp-server # Start only the MCP server
make agent # Start only the Remediation Agent CLI
make status # Check health of all services
make demo # Open all 4 demo URLs in your browser
make clean # Stop everything and delete all Docker volumes
make incident-kafka # Simulate Kafka broker failure
make incident-spark # Simulate Spark failure
make incident-hdfs # Simulate HDFS failure
make incident-stop # Recover all services______________________________________________________________________
技术栈
| 类别 | 技术 | 版本 |
|---|---|---|
| 度量与警报 | 普罗米修斯 | 2.51.0 |
| 可视化 | 格拉法纳 | 10.4.2 |
| 警报路由 | Alertmanager | 0.27.0 |
| MCP服务器框架 | FastAPI+Uvicorn | 0.115/0.32 |
| AI代理 | OpenAI SDK+Qwen | 通过OpenRouter |
| 修正代理UI | Python丰富 | 13.9+ |
| HTTP客户端 | httpx | 0.27 |
| 数据验证 | Pydantic | 2.10+ |
| 公制出口商 | 普罗米修斯客户 | 0.21 |
| 主机指标 | 节点导出器 | 1.7.0 |
| 容器度量 | cAdvisor | 0.49.1 |
| 容器化 | Docker+Compose | 最新 |
| 语言 | Python | 3.11+ |
| 测试 | pytest | 8.1+ |
______________________________________________________________________
作者
阿卜杜拉 — @阿卜杜拉赫950
______________________________________________________________________
课程
该项目是作为 大数据与机器学习 节目在 圣光机大学 (第三学期)。
项目3——通过MCP服务器和修复代理使用Prometheus+Grafana进行监控和自动诊断
______________________________________________________________________
许可证
该项目旨在教育目的,作为 ITMO大学大数据与机器学习课程.
此存储库旨在作为参考和学习资源。
