
性能监控与智能告警系统
大约 19 分钟
性能监控与智能告警系统
前言:让数据"说话"的艺术
还记得我第一次做性能测试时,就像一个盲人摸象,只能看到最终的结果,却不知道测试过程中发生了什么。后来我意识到,好的监控系统就像给测试装上了一双"慧眼",能够实时洞察系统的每一个细微变化。
今天我们就来探讨如何构建一个智能的性能监控与告警系统,让性能测试过程变得透明可控,让问题在萌芽状态就被发现和解决。
监控系统设计理念
设计目标
"""
性能监控系统的设计目标
就像给汽车装上仪表盘,让驾驶变得安全可控
"""
monitoring_goals = {
"实时性": "毫秒级的数据采集和展示",
"全面性": "覆盖系统、应用、业务各个层面",
"准确性": "精确的指标计算和异常检测",
"可视化": "直观的图表和仪表盘展示",
"智能化": "自动异常检测和智能告警",
"可扩展": "支持自定义指标和告警规则"
}监控架构概览
"""
监控系统架构
多层次、全方位的监控体系
"""
monitoring_architecture = {
"数据采集层": {
"系统指标": "CPU、内存、网络、磁盘",
"应用指标": "QPS、响应时间、错误率",
"业务指标": "用户数、订单量、转化率",
"自定义指标": "业务特定的关键指标"
},
"数据处理层": {
"数据清洗": "过滤异常数据,补充缺失值",
"数据聚合": "按时间窗口聚合统计",
"数据计算": "计算衍生指标和趋势",
"数据存储": "时序数据库存储"
},
"分析决策层": {
"异常检测": "基于规则和机器学习的异常检测",
"趋势分析": "性能趋势和容量规划",
"根因分析": "问题定位和影响分析",
"智能告警": "多级告警和降噪处理"
},
"展示交互层": {
"实时仪表盘": "关键指标的实时展示",
"历史报告": "性能趋势和对比分析",
"告警中心": "告警管理和处理跟踪",
"API接口": "数据查询和集成接口"
}
}指标收集系统:数据的"采集器"
1. 多维度指标收集
# core/metrics_collector.py
"""
指标收集器 - 性能数据的"传感器"
全方位收集系统和应用的性能指标
"""
import time
import threading
import psutil
import requests
from typing import Dict, List, Any, Optional, Callable
from dataclasses import dataclass, asdict
from datetime import datetime
import json
import queue
@dataclass
class MetricPoint:
"""指标数据点"""
name: str
value: float
timestamp: datetime
tags: Dict[str, str]
unit: str = ""
def to_dict(self) -> Dict[str, Any]:
"""转换为字典"""
return {
'name': self.name,
'value': self.value,
'timestamp': self.timestamp.isoformat(),
'tags': self.tags,
'unit': self.unit
}
class MetricsCollector:
"""
指标收集器
就像一个勤劳的数据采集员,
不断收集各种性能指标
"""
def __init__(self):
self.collectors: Dict[str, Callable] = {}
self.metrics_queue = queue.Queue(maxsize=10000)
self.is_collecting = False
self.collect_thread = None
self.collect_interval = 5 # 收集间隔(秒)
# 注册默认收集器
self._register_default_collectors()
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger("MetricsCollector")
logger.setLevel(logging.INFO)
return logger
def _register_default_collectors(self):
"""注册默认的指标收集器"""
self.register_collector("system", self._collect_system_metrics)
self.register_collector("process", self._collect_process_metrics)
self.register_collector("network", self._collect_network_metrics)
def register_collector(self, name: str, collector_func: Callable):
"""注册指标收集器"""
self.collectors[name] = collector_func
self.logger.info(f"📊 注册指标收集器: {name}")
def start_collecting(self):
"""开始收集指标"""
if self.is_collecting:
return
self.is_collecting = True
self.collect_thread = threading.Thread(target=self._collect_loop, daemon=True)
self.collect_thread.start()
self.logger.info("📊 指标收集启动")
def stop_collecting(self):
"""停止收集指标"""
self.is_collecting = False
self.logger.info("📊 指标收集停止")
def _collect_loop(self):
"""收集循环"""
while self.is_collecting:
try:
# 执行所有收集器
for name, collector in self.collectors.items():
try:
metrics = collector()
if metrics:
for metric in metrics:
if not self.metrics_queue.full():
self.metrics_queue.put(metric)
else:
self.logger.warning("⚠️ 指标队列已满,丢弃数据")
except Exception as e:
self.logger.error(f"❌ 收集器 {name} 异常: {e}")
time.sleep(self.collect_interval)
except Exception as e:
self.logger.error(f"❌ 收集循环异常: {e}")
def _collect_system_metrics(self) -> List[MetricPoint]:
"""收集系统指标"""
metrics = []
timestamp = datetime.now()
try:
# CPU指标
cpu_percent = psutil.cpu_percent(interval=1)
metrics.append(MetricPoint(
name="system.cpu.usage",
value=cpu_percent,
timestamp=timestamp,
tags={"type": "system"},
unit="percent"
))
# 内存指标
memory = psutil.virtual_memory()
metrics.append(MetricPoint(
name="system.memory.usage",
value=memory.percent,
timestamp=timestamp,
tags={"type": "system"},
unit="percent"
))
metrics.append(MetricPoint(
name="system.memory.available",
value=memory.available / (1024**3), # GB
timestamp=timestamp,
tags={"type": "system"},
unit="GB"
))
# 磁盘指标
disk = psutil.disk_usage('/')
metrics.append(MetricPoint(
name="system.disk.usage",
value=(disk.used / disk.total) * 100,
timestamp=timestamp,
tags={"type": "system", "mount": "/"},
unit="percent"
))
except Exception as e:
self.logger.error(f"❌ 收集系统指标失败: {e}")
return metrics
def _collect_process_metrics(self) -> List[MetricPoint]:
"""收集进程指标"""
metrics = []
timestamp = datetime.now()
try:
process = psutil.Process()
# 进程CPU使用率
cpu_percent = process.cpu_percent()
metrics.append(MetricPoint(
name="process.cpu.usage",
value=cpu_percent,
timestamp=timestamp,
tags={"type": "process", "pid": str(process.pid)},
unit="percent"
))
# 进程内存使用
memory_info = process.memory_info()
metrics.append(MetricPoint(
name="process.memory.rss",
value=memory_info.rss / (1024**2), # MB
timestamp=timestamp,
tags={"type": "process", "pid": str(process.pid)},
unit="MB"
))
# 进程线程数
num_threads = process.num_threads()
metrics.append(MetricPoint(
name="process.threads.count",
value=num_threads,
timestamp=timestamp,
tags={"type": "process", "pid": str(process.pid)},
unit="count"
))
except Exception as e:
self.logger.error(f"❌ 收集进程指标失败: {e}")
return metrics
def _collect_network_metrics(self) -> List[MetricPoint]:
"""收集网络指标"""
metrics = []
timestamp = datetime.now()
try:
# 网络IO统计
net_io = psutil.net_io_counters()
metrics.append(MetricPoint(
name="network.bytes.sent",
value=net_io.bytes_sent / (1024**2), # MB
timestamp=timestamp,
tags={"type": "network"},
unit="MB"
))
metrics.append(MetricPoint(
name="network.bytes.recv",
value=net_io.bytes_recv / (1024**2), # MB
timestamp=timestamp,
tags={"type": "network"},
unit="MB"
))
# 网络连接数
connections = len(psutil.net_connections())
metrics.append(MetricPoint(
name="network.connections.count",
value=connections,
timestamp=timestamp,
tags={"type": "network"},
unit="count"
))
except Exception as e:
self.logger.error(f"❌ 收集网络指标失败: {e}")
return metrics
def get_metrics(self, count: int = 100) -> List[MetricPoint]:
"""获取指标数据"""
metrics = []
for _ in range(min(count, self.metrics_queue.qsize())):
try:
metric = self.metrics_queue.get_nowait()
metrics.append(metric)
except queue.Empty:
break
return metrics
def add_custom_metric(self, name: str, value: float, tags: Dict[str, str] = None, unit: str = ""):
"""添加自定义指标"""
metric = MetricPoint(
name=name,
value=value,
timestamp=datetime.now(),
tags=tags or {},
unit=unit
)
if not self.metrics_queue.full():
self.metrics_queue.put(metric)
else:
self.logger.warning("⚠️ 指标队列已满,丢弃自定义指标")
# 全局指标收集器实例
metrics_collector = MetricsCollector()2. Locust集成指标收集
# core/locust_metrics.py
"""
Locust指标收集器 - 专门收集压测相关指标
"""
from locust import events
import time
from typing import Dict, Any, List
from datetime import datetime
class LocustMetricsCollector:
"""
Locust指标收集器
专门收集压测过程中的性能指标
"""
def __init__(self, metrics_collector):
self.metrics_collector = metrics_collector
self.request_stats = {}
self.user_stats = {}
self.error_stats = {}
# 注册Locust事件监听器
self._register_event_listeners()
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger("LocustMetricsCollector")
logger.setLevel(logging.INFO)
return logger
def _register_event_listeners(self):
"""注册事件监听器"""
events.request_success.add_listener(self._on_request_success)
events.request_failure.add_listener(self._on_request_failure)
events.user_add.add_listener(self._on_user_add)
events.user_remove.add_listener(self._on_user_remove)
events.test_start.add_listener(self._on_test_start)
events.test_stop.add_listener(self._on_test_stop)
def _on_request_success(self, request_type, name, response_time, response_length, **kwargs):
"""请求成功事件处理"""
# 记录请求指标
self.metrics_collector.add_custom_metric(
name="locust.request.response_time",
value=response_time,
tags={
"method": request_type,
"endpoint": name,
"status": "success"
},
unit="ms"
)
self.metrics_collector.add_custom_metric(
name="locust.request.response_size",
value=response_length,
tags={
"method": request_type,
"endpoint": name
},
unit="bytes"
)
# 更新统计
key = f"{request_type}:{name}"
if key not in self.request_stats:
self.request_stats[key] = {
"count": 0,
"total_time": 0,
"min_time": float('inf'),
"max_time": 0
}
stats = self.request_stats[key]
stats["count"] += 1
stats["total_time"] += response_time
stats["min_time"] = min(stats["min_time"], response_time)
stats["max_time"] = max(stats["max_time"], response_time)
def _on_request_failure(self, request_type, name, response_time, response_length, exception, **kwargs):
"""请求失败事件处理"""
self.metrics_collector.add_custom_metric(
name="locust.request.error",
value=1,
tags={
"method": request_type,
"endpoint": name,
"error": str(exception)[:100] # 限制错误信息长度
},
unit="count"
)
# 记录错误统计
error_key = f"{request_type}:{name}:{type(exception).__name__}"
self.error_stats[error_key] = self.error_stats.get(error_key, 0) + 1
def _on_user_add(self, user_instance, **kwargs):
"""用户增加事件处理"""
self.metrics_collector.add_custom_metric(
name="locust.users.active",
value=1,
tags={"action": "add", "user_class": user_instance.__class__.__name__},
unit="count"
)
def _on_user_remove(self, user_instance, **kwargs):
"""用户移除事件处理"""
self.metrics_collector.add_custom_metric(
name="locust.users.active",
value=-1,
tags={"action": "remove", "user_class": user_instance.__class__.__name__},
unit="count"
)
def _on_test_start(self, environment, **kwargs):
"""测试开始事件处理"""
self.logger.info("🚀 Locust测试开始,指标收集启动")
self.metrics_collector.add_custom_metric(
name="locust.test.status",
value=1,
tags={"status": "started"},
unit="boolean"
)
def _on_test_stop(self, environment, **kwargs):
"""测试停止事件处理"""
self.logger.info("🛑 Locust测试结束,生成最终指标")
self.metrics_collector.add_custom_metric(
name="locust.test.status",
value=0,
tags={"status": "stopped"},
unit="boolean"
)
# 生成测试摘要指标
self._generate_summary_metrics(environment)
def _generate_summary_metrics(self, environment):
"""生成测试摘要指标"""
stats = environment.runner.stats
# 总体统计
self.metrics_collector.add_custom_metric(
name="locust.summary.total_requests",
value=stats.total.num_requests,
tags={"type": "summary"},
unit="count"
)
self.metrics_collector.add_custom_metric(
name="locust.summary.total_failures",
value=stats.total.num_failures,
tags={"type": "summary"},
unit="count"
)
self.metrics_collector.add_custom_metric(
name="locust.summary.avg_response_time",
value=stats.total.avg_response_time,
tags={"type": "summary"},
unit="ms"
)
self.metrics_collector.add_custom_metric(
name="locust.summary.max_response_time",
value=stats.total.max_response_time,
tags={"type": "summary"},
unit="ms"
)
# 成功率
success_rate = (1 - stats.total.num_failures / max(stats.total.num_requests, 1)) * 100
self.metrics_collector.add_custom_metric(
name="locust.summary.success_rate",
value=success_rate,
tags={"type": "summary"},
unit="percent"
)
# 初始化Locust指标收集器
locust_metrics = LocustMetricsCollector(metrics_collector)智能告警系统:问题的"预警雷达"
1. 多级告警引擎
# core/alert_engine.py
"""
智能告警引擎 - 问题的早期预警系统
基于规则和机器学习的智能告警
"""
import time
import threading
import json
from typing import Dict, List, Any, Optional, Callable
from dataclasses import dataclass, asdict
from datetime import datetime, timedelta
from enum import Enum
import statistics
import requests
class AlertLevel(Enum):
"""告警级别"""
INFO = "info"
WARNING = "warning"
CRITICAL = "critical"
EMERGENCY = "emergency"
class AlertStatus(Enum):
"""告警状态"""
ACTIVE = "active"
RESOLVED = "resolved"
SUPPRESSED = "suppressed"
@dataclass
class AlertRule:
"""告警规则"""
name: str
metric_name: str
condition: str # >, <, >=, <=, ==, !=
threshold: float
duration: int # 持续时间(秒)
level: AlertLevel
description: str
tags: Dict[str, str]
enabled: bool = True
def evaluate(self, value: float) -> bool:
"""评估告警条件"""
if self.condition == ">":
return value > self.threshold
elif self.condition == "<":
return value < self.threshold
elif self.condition == ">=":
return value >= self.threshold
elif self.condition == "<=":
return value <= self.threshold
elif self.condition == "==":
return value == self.threshold
elif self.condition == "!=":
return value != self.threshold
else:
return False
@dataclass
class Alert:
"""告警实例"""
id: str
rule_name: str
metric_name: str
current_value: float
threshold: float
level: AlertLevel
status: AlertStatus
message: str
start_time: datetime
end_time: Optional[datetime]
tags: Dict[str, str]
def to_dict(self) -> Dict[str, Any]:
"""转换为字典"""
return {
'id': self.id,
'rule_name': self.rule_name,
'metric_name': self.metric_name,
'current_value': self.current_value,
'threshold': self.threshold,
'level': self.level.value,
'status': self.status.value,
'message': self.message,
'start_time': self.start_time.isoformat(),
'end_time': self.end_time.isoformat() if self.end_time else None,
'tags': self.tags,
'duration': (self.end_time - self.start_time).total_seconds() if self.end_time else (datetime.now() - self.start_time).total_seconds()
}
class AlertEngine:
"""
智能告警引擎
就像一个智能的安全卫士,
时刻监控系统状态,及时发出预警
"""
def __init__(self, metrics_collector):
self.metrics_collector = metrics_collector
self.rules: Dict[str, AlertRule] = {}
self.active_alerts: Dict[str, Alert] = {}
self.alert_history: List[Alert] = []
self.notification_handlers: List[Callable] = []
self.is_monitoring = False
self.monitor_thread = None
self.check_interval = 10 # 检查间隔(秒)
# 告警抑制配置
self.suppression_rules = {}
self.cooldown_period = 300 # 冷却期(秒)
# 指标缓存
self.metric_cache: Dict[str, List[float]] = {}
self.cache_size = 100
self._register_default_rules()
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger("AlertEngine")
logger.setLevel(logging.INFO)
return logger
def _register_default_rules(self):
"""注册默认告警规则"""
default_rules = [
AlertRule(
name="high_cpu_usage",
metric_name="system.cpu.usage",
condition=">",
threshold=80.0,
duration=60,
level=AlertLevel.WARNING,
description="CPU使用率过高",
tags={"category": "system"}
),
AlertRule(
name="critical_cpu_usage",
metric_name="system.cpu.usage",
condition=">",
threshold=95.0,
duration=30,
level=AlertLevel.CRITICAL,
description="CPU使用率严重过高",
tags={"category": "system"}
),
AlertRule(
name="high_memory_usage",
metric_name="system.memory.usage",
condition=">",
threshold=85.0,
duration=120,
level=AlertLevel.WARNING,
description="内存使用率过高",
tags={"category": "system"}
),
AlertRule(
name="high_response_time",
metric_name="locust.request.response_time",
condition=">",
threshold=5000.0,
duration=60,
level=AlertLevel.WARNING,
description="响应时间过长",
tags={"category": "performance"}
),
AlertRule(
name="high_error_rate",
metric_name="locust.request.error",
condition=">",
threshold=5.0,
duration=30,
level=AlertLevel.CRITICAL,
description="错误率过高",
tags={"category": "reliability"}
)
]
for rule in default_rules:
self.add_rule(rule)
def add_rule(self, rule: AlertRule):
"""添加告警规则"""
self.rules[rule.name] = rule
self.logger.info(f"📋 添加告警规则: {rule.name}")
def remove_rule(self, rule_name: str):
"""移除告警规则"""
if rule_name in self.rules:
del self.rules[rule_name]
self.logger.info(f"🗑️ 移除告警规则: {rule_name}")
def add_notification_handler(self, handler: Callable[[Alert], None]):
"""添加通知处理器"""
self.notification_handlers.append(handler)
def start_monitoring(self):
"""开始监控"""
if self.is_monitoring:
return
self.is_monitoring = True
self.monitor_thread = threading.Thread(target=self._monitor_loop, daemon=True)
self.monitor_thread.start()
self.logger.info("🚨 告警监控启动")
def stop_monitoring(self):
"""停止监控"""
self.is_monitoring = False
self.logger.info("🚨 告警监控停止")
def _monitor_loop(self):
"""监控循环"""
while self.is_monitoring:
try:
# 获取最新指标
metrics = self.metrics_collector.get_metrics(1000)
# 更新指标缓存
self._update_metric_cache(metrics)
# 检查告警规则
self._check_alert_rules()
# 清理过期告警
self._cleanup_resolved_alerts()
time.sleep(self.check_interval)
except Exception as e:
self.logger.error(f"❌ 告警监控异常: {e}")
def _update_metric_cache(self, metrics: List[MetricPoint]):
"""更新指标缓存"""
for metric in metrics:
if metric.name not in self.metric_cache:
self.metric_cache[metric.name] = []
cache = self.metric_cache[metric.name]
cache.append(metric.value)
# 保持缓存大小
if len(cache) > self.cache_size:
cache.pop(0)
def _check_alert_rules(self):
"""检查告警规则"""
for rule_name, rule in self.rules.items():
if not rule.enabled:
continue
try:
self._evaluate_rule(rule)
except Exception as e:
self.logger.error(f"❌ 评估规则 {rule_name} 异常: {e}")
def _evaluate_rule(self, rule: AlertRule):
"""评估单个规则"""
if rule.metric_name not in self.metric_cache:
return
cache = self.metric_cache[rule.metric_name]
if not cache:
return
# 获取最新值
current_value = cache[-1]
# 检查是否触发告警
if rule.evaluate(current_value):
# 检查持续时间
if self._check_duration(rule, cache):
self._trigger_alert(rule, current_value)
else:
# 检查是否需要解除告警
self._resolve_alert(rule.name)
def _check_duration(self, rule: AlertRule, cache: List[float]) -> bool:
"""检查持续时间条件"""
if rule.duration <= 0:
return True
# 计算需要检查的数据点数量
points_needed = max(1, rule.duration // self.check_interval)
if len(cache) < points_needed:
return False
# 检查最近的数据点是否都满足条件
recent_values = cache[-points_needed:]
return all(rule.evaluate(value) for value in recent_values)
def _trigger_alert(self, rule: AlertRule, current_value: float):
"""触发告警"""
alert_id = f"{rule.name}_{int(time.time())}"
# 检查是否已有活跃告警
if rule.name in self.active_alerts:
return
# 检查冷却期
if self._is_in_cooldown(rule.name):
return
# 创建告警
alert = Alert(
id=alert_id,
rule_name=rule.name,
metric_name=rule.metric_name,
current_value=current_value,
threshold=rule.threshold,
level=rule.level,
status=AlertStatus.ACTIVE,
message=f"{rule.description}: 当前值 {current_value:.2f} {rule.condition} 阈值 {rule.threshold}",
start_time=datetime.now(),
end_time=None,
tags=rule.tags
)
self.active_alerts[rule.name] = alert
self.alert_history.append(alert)
# 发送通知
self._send_notifications(alert)
self.logger.warning(f"🚨 触发告警: {alert.message}")
def _resolve_alert(self, rule_name: str):
"""解除告警"""
if rule_name in self.active_alerts:
alert = self.active_alerts[rule_name]
alert.status = AlertStatus.RESOLVED
alert.end_time = datetime.now()
del self.active_alerts[rule_name]
# 发送解除通知
self._send_notifications(alert)
self.logger.info(f"✅ 解除告警: {alert.rule_name}")
def _is_in_cooldown(self, rule_name: str) -> bool:
"""检查是否在冷却期"""
# 查找最近的告警
recent_alerts = [
alert for alert in self.alert_history
if alert.rule_name == rule_name and alert.end_time
]
if not recent_alerts:
return False
latest_alert = max(recent_alerts, key=lambda a: a.end_time)
time_since_resolved = (datetime.now() - latest_alert.end_time).total_seconds()
return time_since_resolved < self.cooldown_period
def _send_notifications(self, alert: Alert):
"""发送通知"""
for handler in self.notification_handlers:
try:
handler(alert)
except Exception as e:
self.logger.error(f"❌ 通知处理器异常: {e}")
def _cleanup_resolved_alerts(self):
"""清理已解除的告警"""
# 保留最近1000条告警记录
if len(self.alert_history) > 1000:
self.alert_history = self.alert_history[-1000:]
def get_active_alerts(self) -> List[Alert]:
"""获取活跃告警"""
return list(self.active_alerts.values())
def get_alert_history(self, hours: int = 24) -> List[Alert]:
"""获取告警历史"""
cutoff_time = datetime.now() - timedelta(hours=hours)
return [
alert for alert in self.alert_history
if alert.start_time > cutoff_time
]
def get_alert_statistics(self) -> Dict[str, Any]:
"""获取告警统计"""
active_count = len(self.active_alerts)
# 按级别统计
level_stats = {}
for level in AlertLevel:
level_stats[level.value] = len([
alert for alert in self.active_alerts.values()
if alert.level == level
])
# 最近24小时统计
recent_alerts = self.get_alert_history(24)
recent_count = len(recent_alerts)
return {
"active_alerts": active_count,
"recent_alerts_24h": recent_count,
"alerts_by_level": level_stats,
"total_rules": len(self.rules),
"enabled_rules": len([r for r in self.rules.values() if r.enabled])
}
# 全局告警引擎实例
alert_engine = AlertEngine(metrics_collector)2. 通知处理器
# core/notification_handlers.py
"""
通知处理器 - 告警信息的"传令兵"
支持多种通知方式
"""
import requests
import json
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from typing import Dict, Any
class EmailNotificationHandler:
"""邮件通知处理器"""
def __init__(self, smtp_server: str, smtp_port: int, username: str, password: str):
self.smtp_server = smtp_server
self.smtp_port = smtp_port
self.username = username
self.password = password
def __call__(self, alert: Alert):
"""发送邮件通知"""
try:
# 创建邮件
msg = MIMEMultipart()
msg['From'] = self.username
msg['To'] = "admin@example.com" # 可配置
msg['Subject'] = f"[{alert.level.value.upper()}] {alert.rule_name}"
# 邮件内容
body = f"""
告警详情:
规则名称: {alert.rule_name}
告警级别: {alert.level.value}
指标名称: {alert.metric_name}
当前值: {alert.current_value}
阈值: {alert.threshold}
状态: {alert.status.value}
开始时间: {alert.start_time}
消息: {alert.message}
标签: {json.dumps(alert.tags, indent=2)}
"""
msg.attach(MIMEText(body, 'plain'))
# 发送邮件
server = smtplib.SMTP(self.smtp_server, self.smtp_port)
server.starttls()
server.login(self.username, self.password)
server.send_message(msg)
server.quit()
except Exception as e:
print(f"❌ 邮件发送失败: {e}")
class WebhookNotificationHandler:
"""Webhook通知处理器"""
def __init__(self, webhook_url: str, headers: Dict[str, str] = None):
self.webhook_url = webhook_url
self.headers = headers or {'Content-Type': 'application/json'}
def __call__(self, alert: Alert):
"""发送Webhook通知"""
try:
payload = {
"alert": alert.to_dict(),
"timestamp": alert.start_time.isoformat()
}
response = requests.post(
self.webhook_url,
json=payload,
headers=self.headers,
timeout=10
)
if response.status_code != 200:
print(f"⚠️ Webhook响应异常: {response.status_code}")
except Exception as e:
print(f"❌ Webhook发送失败: {e}")
class DingTalkNotificationHandler:
"""钉钉通知处理器"""
def __init__(self, webhook_url: str, secret: str = None):
self.webhook_url = webhook_url
self.secret = secret
def __call__(self, alert: Alert):
"""发送钉钉通知"""
try:
# 构造消息
level_emoji = {
'info': 'ℹ️',
'warning': '⚠️',
'critical': '🚨',
'emergency': '🔥'
}
emoji = level_emoji.get(alert.level.value, '📢')
message = f"""
{emoji} **性能告警通知**
**规则名称**: {alert.rule_name}
**告警级别**: {alert.level.value.upper()}
**指标名称**: {alert.metric_name}
**当前值**: {alert.current_value:.2f}
**阈值**: {alert.threshold}
**状态**: {alert.status.value}
**时间**: {alert.start_time.strftime('%Y-%m-%d %H:%M:%S')}
**详情**: {alert.message}
"""
payload = {
"msgtype": "markdown",
"markdown": {
"title": f"性能告警 - {alert.rule_name}",
"text": message
}
}
response = requests.post(
self.webhook_url,
json=payload,
headers={'Content-Type': 'application/json'},
timeout=10
)
if response.status_code != 200:
print(f"⚠️ 钉钉通知响应异常: {response.status_code}")
except Exception as e:
print(f"❌ 钉钉通知发送失败: {e}")
# 注册通知处理器示例
def setup_notification_handlers():
"""设置通知处理器"""
# 钉钉通知
dingtalk_handler = DingTalkNotificationHandler(
webhook_url="https://oapi.dingtalk.com/robot/send?access_token=YOUR_TOKEN"
)
alert_engine.add_notification_handler(dingtalk_handler)
# Webhook通知
webhook_handler = WebhookNotificationHandler(
webhook_url="https://your-webhook-endpoint.com/alerts"
)
alert_engine.add_notification_handler(webhook_handler)
print("📢 通知处理器设置完成")可视化监控面板:数据的"艺术展示"
1. 实时监控仪表盘
# core/dashboard.py
"""
监控仪表盘 - 数据的可视化展示
提供实时的性能监控界面
"""
from flask import Flask, render_template, jsonify, request
import json
from datetime import datetime, timedelta
from typing import Dict, List, Any
class MonitoringDashboard:
"""
监控仪表盘
就像汽车的仪表盘一样,
直观展示系统的各项指标
"""
def __init__(self, metrics_collector, alert_engine, cluster_manager=None):
self.app = Flask(__name__)
self.metrics_collector = metrics_collector
self.alert_engine = alert_engine
self.cluster_manager = cluster_manager
self._setup_routes()
def _setup_routes(self):
"""设置路由"""
@self.app.route('/')
def index():
"""主页"""
return render_template('dashboard.html')
@self.app.route('/api/metrics/current')
def get_current_metrics():
"""获取当前指标"""
try:
metrics = self.metrics_collector.get_metrics(100)
# 按指标名称分组
grouped_metrics = {}
for metric in metrics:
if metric.name not in grouped_metrics:
grouped_metrics[metric.name] = []
grouped_metrics[metric.name].append({
'value': metric.value,
'timestamp': metric.timestamp.isoformat(),
'tags': metric.tags,
'unit': metric.unit
})
return jsonify({
'status': 'success',
'data': grouped_metrics,
'timestamp': datetime.now().isoformat()
})
except Exception as e:
return jsonify({
'status': 'error',
'message': str(e)
}), 500
@self.app.route('/api/alerts/active')
def get_active_alerts():
"""获取活跃告警"""
try:
alerts = self.alert_engine.get_active_alerts()
return jsonify({
'status': 'success',
'data': [alert.to_dict() for alert in alerts],
'count': len(alerts)
})
except Exception as e:
return jsonify({
'status': 'error',
'message': str(e)
}), 500
@self.app.route('/api/alerts/history')
def get_alert_history():
"""获取告警历史"""
try:
hours = request.args.get('hours', 24, type=int)
alerts = self.alert_engine.get_alert_history(hours)
return jsonify({
'status': 'success',
'data': [alert.to_dict() for alert in alerts],
'count': len(alerts)
})
except Exception as e:
return jsonify({
'status': 'error',
'message': str(e)
}), 500
@self.app.route('/api/cluster/status')
def get_cluster_status():
"""获取集群状态"""
try:
if self.cluster_manager:
status = self.cluster_manager.get_cluster_status()
return jsonify({
'status': 'success',
'data': status
})
else:
return jsonify({
'status': 'error',
'message': '集群管理器未配置'
}), 404
except Exception as e:
return jsonify({
'status': 'error',
'message': str(e)
}), 500
@self.app.route('/api/statistics/summary')
def get_statistics_summary():
"""获取统计摘要"""
try:
# 告警统计
alert_stats = self.alert_engine.get_alert_statistics()
# 系统统计(示例)
system_stats = {
'uptime': '2 days 5 hours',
'total_requests': 1234567,
'avg_response_time': 245.6,
'success_rate': 99.8
}
return jsonify({
'status': 'success',
'data': {
'alerts': alert_stats,
'system': system_stats,
'timestamp': datetime.now().isoformat()
}
})
except Exception as e:
return jsonify({
'status': 'error',
'message': str(e)
}), 500
def run(self, host='0.0.0.0', port=5000, debug=False):
"""启动仪表盘"""
print(f"📊 监控仪表盘启动: http://{host}:{port}")
self.app.run(host=host, port=port, debug=debug)
# HTML模板示例
dashboard_html_template = """
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>性能监控仪表盘</title>
<script src="https://cdn.jsdelivr.net/npm/chart.js"></script>
<style>
body {
font-family: 'Segoe UI', Tahoma, Geneva, Verdana, sans-serif;
margin: 0;
padding: 20px;
background-color: #f5f5f5;
}
.dashboard {
display: grid;
grid-template-columns: repeat(auto-fit, minmax(300px, 1fr));
gap: 20px;
}
.card {
background: white;
border-radius: 8px;
padding: 20px;
box-shadow: 0 2px 4px rgba(0,0,0,0.1);
}
.card h3 {
margin-top: 0;
color: #333;
}
.metric-value {
font-size: 2em;
font-weight: bold;
color: #007bff;
}
.alert-item {
padding: 10px;
margin: 5px 0;
border-radius: 4px;
border-left: 4px solid;
}
.alert-warning {
background-color: #fff3cd;
border-color: #ffc107;
}
.alert-critical {
background-color: #f8d7da;
border-color: #dc3545;
}
.status-indicator {
display: inline-block;
width: 12px;
height: 12px;
border-radius: 50%;
margin-right: 8px;
}
.status-healthy { background-color: #28a745; }
.status-warning { background-color: #ffc107; }
.status-critical { background-color: #dc3545; }
</style>
</head>
<body>
<h1>🚀 性能监控仪表盘</h1>
<div class="dashboard">
<!-- 系统概览 -->
<div class="card">
<h3>📊 系统概览</h3>
<div id="system-overview">
<div>CPU使用率: <span class="metric-value" id="cpu-usage">--</span>%</div>
<div>内存使用率: <span class="metric-value" id="memory-usage">--</span>%</div>
<div>活跃用户: <span class="metric-value" id="active-users">--</span></div>
</div>
</div>
<!-- 性能指标 -->
<div class="card">
<h3>⚡ 性能指标</h3>
<canvas id="performance-chart" width="400" height="200"></canvas>
</div>
<!-- 活跃告警 -->
<div class="card">
<h3>🚨 活跃告警</h3>
<div id="active-alerts">
<div class="alert-item alert-warning">
<span class="status-indicator status-warning"></span>
暂无活跃告警
</div>
</div>
</div>
<!-- 集群状态 -->
<div class="card">
<h3>🏗️ 集群状态</h3>
<div id="cluster-status">
<div>总节点数: <span id="total-nodes">--</span></div>
<div>健康节点: <span id="healthy-nodes">--</span></div>
<div>总请求数: <span id="total-requests">--</span></div>
</div>
</div>
</div>
<script>
// 初始化图表
const ctx = document.getElementById('performance-chart').getContext('2d');
const performanceChart = new Chart(ctx, {
type: 'line',
data: {
labels: [],
datasets: [{
label: '响应时间 (ms)',
data: [],
borderColor: 'rgb(75, 192, 192)',
tension: 0.1
}]
},
options: {
responsive: true,
scales: {
y: {
beginAtZero: true
}
}
}
});
// 更新数据函数
function updateDashboard() {
// 获取当前指标
fetch('/api/metrics/current')
.then(response => response.json())
.then(data => {
if (data.status === 'success') {
updateSystemOverview(data.data);
updatePerformanceChart(data.data);
}
})
.catch(error => console.error('Error:', error));
// 获取活跃告警
fetch('/api/alerts/active')
.then(response => response.json())
.then(data => {
if (data.status === 'success') {
updateActiveAlerts(data.data);
}
})
.catch(error => console.error('Error:', error));
// 获取集群状态
fetch('/api/cluster/status')
.then(response => response.json())
.then(data => {
if (data.status === 'success') {
updateClusterStatus(data.data);
}
})
.catch(error => console.error('Error:', error));
}
function updateSystemOverview(metrics) {
if (metrics['system.cpu.usage']) {
const cpuUsage = metrics['system.cpu.usage'][0].value;
document.getElementById('cpu-usage').textContent = cpuUsage.toFixed(1);
}
if (metrics['system.memory.usage']) {
const memoryUsage = metrics['system.memory.usage'][0].value;
document.getElementById('memory-usage').textContent = memoryUsage.toFixed(1);
}
if (metrics['locust.users.active']) {
const activeUsers = metrics['locust.users.active'].length;
document.getElementById('active-users').textContent = activeUsers;
}
}
function updatePerformanceChart(metrics) {
if (metrics['locust.request.response_time']) {
const responseTimeData = metrics['locust.request.response_time'];
const labels = responseTimeData.map(item =>
new Date(item.timestamp).toLocaleTimeString()
);
const values = responseTimeData.map(item => item.value);
performanceChart.data.labels = labels.slice(-20); // 保留最近20个数据点
performanceChart.data.datasets[0].data = values.slice(-20);
performanceChart.update();
}
}
function updateActiveAlerts(alerts) {
const container = document.getElementById('active-alerts');
if (alerts.length === 0) {
container.innerHTML = '<div class="alert-item"><span class="status-indicator status-healthy"></span>暂无活跃告警</div>';
return;
}
container.innerHTML = alerts.map(alert => {
const levelClass = alert.level === 'critical' ? 'alert-critical' : 'alert-warning';
const statusClass = alert.level === 'critical' ? 'status-critical' : 'status-warning';
return `
<div class="alert-item ${levelClass}">
<span class="status-indicator ${statusClass}"></span>
<strong>${alert.rule_name}</strong><br>
${alert.message}
</div>
`;
}).join('');
}
function updateClusterStatus(clusterData) {
if (clusterData.cluster_info) {
const info = clusterData.cluster_info;
document.getElementById('total-nodes').textContent = info.total_workers || '--';
document.getElementById('healthy-nodes').textContent = info.healthy_workers || '--';
document.getElementById('total-requests').textContent = info.total_requests || '--';
}
}
// 定期更新数据
updateDashboard();
setInterval(updateDashboard, 5000); // 每5秒更新一次
</script>
</body>
</html>
"""2. Grafana集成配置
# monitoring/grafana/dashboards/locust-performance.json
{
"dashboard": {
"id": null,
"title": "Locust性能监控",
"tags": ["locust", "performance"],
"timezone": "browser",
"panels": [
{
"id": 1,
"title": "响应时间趋势",
"type": "graph",
"targets": [
{
"expr": "locust_request_response_time",
"legendFormat": "{{method}} {{endpoint}}"
}
],
"yAxes": [
{
"label": "响应时间 (ms)",
"min": 0
}
]
},
{
"id": 2,
"title": "请求成功率",
"type": "stat",
"targets": [
{
"expr": "rate(locust_request_success_total[5m]) / rate(locust_request_total[5m]) * 100"
}
]
},
{
"id": 3,
"title": "系统资源使用",
"type": "graph",
"targets": [
{
"expr": "system_cpu_usage",
"legendFormat": "CPU使用率"
},
{
"expr": "system_memory_usage",
"legendFormat": "内存使用率"
}
]
}
],
"time": {
"from": "now-1h",
"to": "now"
},
"refresh": "5s"
}
}实战部署与配置
1. 完整的监控系统启动脚本
# start_monitoring.py
"""
监控系统启动脚本
一键启动完整的监控和告警系统
"""
import threading
import time
from core.metrics_collector import metrics_collector
from core.alert_engine import alert_engine
from core.dashboard import MonitoringDashboard
from core.notification_handlers import setup_notification_handlers
def main():
"""主函数"""
print("🚀 启动性能监控系统...")
# 1. 启动指标收集
print("📊 启动指标收集...")
metrics_collector.start_collecting()
# 2. 设置通知处理器
print("📢 设置通知处理器...")
setup_notification_handlers()
# 3. 启动告警引擎
print("🚨 启动告警引擎...")
alert_engine.start_monitoring()
# 4. 启动监控仪表盘
print("📈 启动监控仪表盘...")
dashboard = MonitoringDashboard(metrics_collector, alert_engine)
# 在单独线程中启动仪表盘
dashboard_thread = threading.Thread(
target=lambda: dashboard.run(host='0.0.0.0', port=5000),
daemon=True
)
dashboard_thread.start()
print("✅ 监控系统启动完成!")
print("📊 监控仪表盘: http://localhost:5000")
print("🚨 告警系统已激活")
print("📈 指标收集已开始")
try:
# 保持主线程运行
while True:
time.sleep(60)
print(f"📊 系统运行中... 活跃告警: {len(alert_engine.get_active_alerts())}")
except KeyboardInterrupt:
print("\n🛑 正在停止监控系统...")
metrics_collector.stop_collecting()
alert_engine.stop_monitoring()
print("✅ 监控系统已停止")
if __name__ == "__main__":
main()最佳实践与经验总结
1. 监控指标设计原则
"""
监控指标设计的最佳实践
"""
monitoring_best_practices = {
"指标选择": {
"原则": "选择真正重要的指标,避免指标爆炸",
"建议": "遵循USE方法论(使用率、饱和度、错误)",
"示例": "CPU使用率、内存使用率、响应时间、错误率"
},
"告警设计": {
"原则": "告警应该可操作,避免告警疲劳",
"建议": "设置合理的阈值和持续时间",
"示例": "CPU使用率>80%持续5分钟才告警"
},
"数据保留": {
"原则": "根据业务需求设置合理的数据保留期",
"建议": "实时数据保留1天,聚合数据保留1年",
"示例": "1分钟粒度保留1天,1小时粒度保留1年"
},
"可视化设计": {
"原则": "图表应该直观易懂,突出重点",
"建议": "使用合适的图表类型,避免信息过载",
"示例": "趋势用折线图,分布用直方图,状态用仪表盘"
}
}2. 故障处理流程
"""
基于监控的故障处理流程
"""
incident_response_flow = {
"检测": "监控系统自动检测异常并触发告警",
"通知": "通过多种渠道通知相关人员",
"响应": "运维人员快速响应并开始处理",
"诊断": "利用监控数据进行问题诊断",
"修复": "根据诊断结果进行问题修复",
"验证": "验证问题是否解决",
"总结": "事后总结并优化监控规则"
}总结
性能监控与智能告警系统就像给测试装上了"千里眼"和"顺风耳",让我们能够实时洞察系统状态,及时发现和解决问题。通过这篇文章,我们深入了解了:
- 指标收集:全方位的性能数据采集体系
- 智能告警:基于规则和机器学习的告警引擎
- 通知机制:多渠道的告警通知系统
- 可视化展示:直观的监控仪表盘和图表
- 最佳实践:监控系统设计和运维经验
这套监控告警系统不仅提供了全面的可观测性,更重要的是建立了完整的问题发现和处理机制,让性能测试变得更加智能和可控。
下一篇文章,我们将对整个性能测试框架项目进行总结,分享项目开发的经验和最佳实践。
推荐阅读
