
分布式压测架构与集群管理
大约 16 分钟
分布式压测架构与集群管理
前言:从单兵作战到集团军作战
还记得我第一次遇到需要模拟10万并发用户的需求时,就像一个人要搬一座山一样无助。单台机器的资源限制让我意识到,有些战斗不是一个人能打的,需要组建一支军队。
分布式压测就像指挥一场大型战役,需要统一的指挥中心、协调的作战单位、高效的通信机制。今天我们就来探讨如何构建一个强大而稳定的分布式压测集群,让性能测试能力突破单机限制。
分布式架构设计理念
设计目标
"""
分布式压测架构的设计目标
就像设计一支现代化军队,既要火力强大,又要指挥灵活
"""
distributed_goals = {
"可扩展性": "支持动态扩容,理论上无并发上限",
"高可用性": "单点故障不影响整体测试",
"负载均衡": "合理分配压测任务,避免热点",
"实时监控": "全局视角监控集群状态",
"故障恢复": "自动检测和处理节点故障",
"资源优化": "最大化利用集群资源"
}架构概览
"""
分布式压测架构图
Master-Worker模式的现代化实现
"""
architecture_overview = {
"Master节点": {
"职责": "集群协调、任务分发、结果汇总",
"组件": ["Web UI", "任务调度器", "状态管理器", "结果收集器"],
"特点": "单点管理,可以有备份"
},
"Worker节点": {
"职责": "执行压测任务,上报结果",
"组件": ["任务执行器", "资源监控器", "心跳管理器"],
"特点": "无状态,可动态扩缩容"
},
"负载均衡器": {
"职责": "分发用户请求,避免单点过载",
"组件": ["请求分发器", "健康检查器", "故障转移器"],
"特点": "透明代理,高可用"
},
"监控系统": {
"职责": "实时监控集群状态和性能指标",
"组件": ["指标收集器", "告警管理器", "可视化面板"],
"特点": "全局视角,实时反馈"
}
}Master节点设计:集群的"大脑"
1. 集群管理器
# core/cluster_manager.py
"""
集群管理器 - 分布式系统的指挥中心
负责协调整个集群的运行
"""
import time
import threading
import json
from typing import Dict, List, Any, Optional
from dataclasses import dataclass, asdict
from datetime import datetime, timedelta
import socket
import requests
@dataclass
class WorkerNode:
"""Worker节点信息"""
node_id: str
host: str
port: int
status: str # 'active', 'inactive', 'error'
last_heartbeat: datetime
cpu_usage: float
memory_usage: float
active_users: int
total_requests: int
def is_healthy(self, timeout: int = 30) -> bool:
"""检查节点是否健康"""
if self.status != 'active':
return False
time_diff = datetime.now() - self.last_heartbeat
return time_diff.total_seconds() < timeout
def get_load_score(self) -> float:
"""计算节点负载评分(越低越好)"""
cpu_score = self.cpu_usage / 100.0
memory_score = self.memory_usage / 100.0
user_score = self.active_users / 1000.0 # 假设1000为满载
return (cpu_score * 0.4 + memory_score * 0.3 + user_score * 0.3)
class ClusterManager:
"""
集群管理器
就像一个智能的指挥官,
能够协调整个集群的运行
"""
def __init__(self, master_host: str = "0.0.0.0", master_port: int = 5557):
self.master_host = master_host
self.master_port = master_port
self.workers: Dict[str, WorkerNode] = {}
self.is_running = False
# 线程管理
self.heartbeat_thread = None
self.cleanup_thread = None
# 配置参数
self.heartbeat_interval = 10 # 心跳间隔(秒)
self.worker_timeout = 30 # Worker超时时间(秒)
self.cleanup_interval = 60 # 清理间隔(秒)
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger(f"ClusterManager")
logger.setLevel(logging.INFO)
return logger
def start(self):
"""启动集群管理器"""
if self.is_running:
return
self.is_running = True
# 启动心跳监控线程
self.heartbeat_thread = threading.Thread(target=self._heartbeat_monitor, daemon=True)
self.heartbeat_thread.start()
# 启动清理线程
self.cleanup_thread = threading.Thread(target=self._cleanup_monitor, daemon=True)
self.cleanup_thread.start()
self.logger.info(f"🚀 集群管理器启动: {self.master_host}:{self.master_port}")
def stop(self):
"""停止集群管理器"""
self.is_running = False
self.logger.info("🛑 集群管理器停止")
def register_worker(self, worker_info: Dict[str, Any]) -> bool:
"""注册Worker节点"""
try:
node_id = worker_info['node_id']
worker = WorkerNode(
node_id=node_id,
host=worker_info['host'],
port=worker_info['port'],
status='active',
last_heartbeat=datetime.now(),
cpu_usage=worker_info.get('cpu_usage', 0),
memory_usage=worker_info.get('memory_usage', 0),
active_users=worker_info.get('active_users', 0),
total_requests=worker_info.get('total_requests', 0)
)
self.workers[node_id] = worker
self.logger.info(f"✅ Worker节点注册成功: {node_id} ({worker.host}:{worker.port})")
return True
except Exception as e:
self.logger.error(f"❌ Worker节点注册失败: {e}")
return False
def update_worker_status(self, node_id: str, status_info: Dict[str, Any]) -> bool:
"""更新Worker状态"""
if node_id not in self.workers:
return False
try:
worker = self.workers[node_id]
worker.last_heartbeat = datetime.now()
worker.cpu_usage = status_info.get('cpu_usage', worker.cpu_usage)
worker.memory_usage = status_info.get('memory_usage', worker.memory_usage)
worker.active_users = status_info.get('active_users', worker.active_users)
worker.total_requests = status_info.get('total_requests', worker.total_requests)
worker.status = 'active'
return True
except Exception as e:
self.logger.error(f"❌ 更新Worker状态失败: {e}")
return False
def get_healthy_workers(self) -> List[WorkerNode]:
"""获取健康的Worker节点"""
return [worker for worker in self.workers.values() if worker.is_healthy()]
def get_best_worker(self) -> Optional[WorkerNode]:
"""获取负载最低的Worker节点"""
healthy_workers = self.get_healthy_workers()
if not healthy_workers:
return None
return min(healthy_workers, key=lambda w: w.get_load_score())
def distribute_users(self, total_users: int) -> Dict[str, int]:
"""分配用户到各个Worker节点"""
healthy_workers = self.get_healthy_workers()
if not healthy_workers:
return {}
# 根据节点负载能力分配用户
distribution = {}
# 计算总负载能力(负载越低,能力越强)
total_capacity = sum(1.0 / max(worker.get_load_score(), 0.1) for worker in healthy_workers)
for worker in healthy_workers:
capacity_ratio = (1.0 / max(worker.get_load_score(), 0.1)) / total_capacity
user_count = int(total_users * capacity_ratio)
distribution[worker.node_id] = user_count
# 处理余数
remaining_users = total_users - sum(distribution.values())
if remaining_users > 0:
best_worker = self.get_best_worker()
if best_worker:
distribution[best_worker.node_id] += remaining_users
return distribution
def _heartbeat_monitor(self):
"""心跳监控线程"""
while self.is_running:
try:
current_time = datetime.now()
for node_id, worker in list(self.workers.items()):
time_diff = current_time - worker.last_heartbeat
if time_diff.total_seconds() > self.worker_timeout:
worker.status = 'inactive'
self.logger.warning(f"⚠️ Worker节点超时: {node_id}")
time.sleep(self.heartbeat_interval)
except Exception as e:
self.logger.error(f"❌ 心跳监控异常: {e}")
def _cleanup_monitor(self):
"""清理监控线程"""
while self.is_running:
try:
current_time = datetime.now()
cleanup_threshold = current_time - timedelta(minutes=10)
# 清理长时间未活跃的节点
inactive_nodes = []
for node_id, worker in self.workers.items():
if worker.last_heartbeat < cleanup_threshold:
inactive_nodes.append(node_id)
for node_id in inactive_nodes:
del self.workers[node_id]
self.logger.info(f"🗑️ 清理非活跃节点: {node_id}")
time.sleep(self.cleanup_interval)
except Exception as e:
self.logger.error(f"❌ 清理监控异常: {e}")
def get_cluster_status(self) -> Dict[str, Any]:
"""获取集群状态"""
healthy_workers = self.get_healthy_workers()
total_users = sum(worker.active_users for worker in healthy_workers)
total_requests = sum(worker.total_requests for worker in healthy_workers)
avg_cpu = sum(worker.cpu_usage for worker in healthy_workers) / max(len(healthy_workers), 1)
avg_memory = sum(worker.memory_usage for worker in healthy_workers) / max(len(healthy_workers), 1)
return {
"cluster_info": {
"total_workers": len(self.workers),
"healthy_workers": len(healthy_workers),
"total_active_users": total_users,
"total_requests": total_requests,
"avg_cpu_usage": round(avg_cpu, 2),
"avg_memory_usage": round(avg_memory, 2)
},
"workers": [asdict(worker) for worker in self.workers.values()],
"timestamp": datetime.now().isoformat()
}
# 全局集群管理器实例
cluster_manager = ClusterManager()Worker节点设计:集群的"战士"
1. Worker节点实现
# core/worker_node.py
"""
Worker节点 - 分布式系统的执行单元
负责执行具体的压测任务
"""
import time
import threading
import psutil
import requests
import socket
from typing import Dict, Any, Optional
from datetime import datetime
import uuid
class WorkerNode:
"""
Worker节点
就像一个勤劳的工人,
专注于执行分配给它的任务
"""
def __init__(self, master_host: str, master_port: int = 5557):
self.master_host = master_host
self.master_port = master_port
self.node_id = self._generate_node_id()
self.host = self._get_local_ip()
self.port = 5558 # Worker默认端口
# 状态信息
self.is_running = False
self.active_users = 0
self.total_requests = 0
self.start_time = None
# 线程管理
self.heartbeat_thread = None
self.monitor_thread = None
# 配置参数
self.heartbeat_interval = 5 # 心跳间隔(秒)
self.monitor_interval = 10 # 监控间隔(秒)
self.logger = self._setup_logger()
def _generate_node_id(self) -> str:
"""生成节点ID"""
hostname = socket.gethostname()
timestamp = int(time.time())
random_id = str(uuid.uuid4())[:8]
return f"{hostname}_{timestamp}_{random_id}"
def _get_local_ip(self) -> str:
"""获取本机IP地址"""
try:
# 连接到一个不存在的地址来获取本机IP
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
s.connect(("8.8.8.8", 80))
ip = s.getsockname()[0]
s.close()
return ip
except:
return "127.0.0.1"
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger(f"WorkerNode-{self.node_id}")
logger.setLevel(logging.INFO)
return logger
def start(self):
"""启动Worker节点"""
if self.is_running:
return
self.is_running = True
self.start_time = datetime.now()
# 注册到Master节点
if self._register_to_master():
# 启动心跳线程
self.heartbeat_thread = threading.Thread(target=self._heartbeat_loop, daemon=True)
self.heartbeat_thread.start()
# 启动监控线程
self.monitor_thread = threading.Thread(target=self._monitor_loop, daemon=True)
self.monitor_thread.start()
self.logger.info(f"🚀 Worker节点启动成功: {self.node_id}")
else:
self.is_running = False
self.logger.error("❌ Worker节点启动失败:无法注册到Master")
def stop(self):
"""停止Worker节点"""
self.is_running = False
self._unregister_from_master()
self.logger.info(f"🛑 Worker节点停止: {self.node_id}")
def _register_to_master(self) -> bool:
"""注册到Master节点"""
try:
registration_data = {
"node_id": self.node_id,
"host": self.host,
"port": self.port,
"cpu_usage": self._get_cpu_usage(),
"memory_usage": self._get_memory_usage(),
"active_users": self.active_users,
"total_requests": self.total_requests
}
response = requests.post(
f"http://{self.master_host}:{self.master_port}/register_worker",
json=registration_data,
timeout=10
)
if response.status_code == 200:
self.logger.info(f"✅ 成功注册到Master: {self.master_host}:{self.master_port}")
return True
else:
self.logger.error(f"❌ 注册失败: {response.status_code}")
return False
except Exception as e:
self.logger.error(f"❌ 注册异常: {e}")
return False
def _unregister_from_master(self):
"""从Master节点注销"""
try:
response = requests.post(
f"http://{self.master_host}:{self.master_port}/unregister_worker",
json={"node_id": self.node_id},
timeout=5
)
if response.status_code == 200:
self.logger.info("✅ 成功从Master注销")
except Exception as e:
self.logger.warning(f"⚠️ 注销异常: {e}")
def _heartbeat_loop(self):
"""心跳循环"""
while self.is_running:
try:
self._send_heartbeat()
time.sleep(self.heartbeat_interval)
except Exception as e:
self.logger.error(f"❌ 心跳异常: {e}")
def _send_heartbeat(self):
"""发送心跳"""
try:
heartbeat_data = {
"node_id": self.node_id,
"cpu_usage": self._get_cpu_usage(),
"memory_usage": self._get_memory_usage(),
"active_users": self.active_users,
"total_requests": self.total_requests,
"timestamp": datetime.now().isoformat()
}
response = requests.post(
f"http://{self.master_host}:{self.master_port}/worker_heartbeat",
json=heartbeat_data,
timeout=5
)
if response.status_code != 200:
self.logger.warning(f"⚠️ 心跳响应异常: {response.status_code}")
except Exception as e:
self.logger.error(f"❌ 发送心跳失败: {e}")
def _monitor_loop(self):
"""监控循环"""
while self.is_running:
try:
self._collect_metrics()
time.sleep(self.monitor_interval)
except Exception as e:
self.logger.error(f"❌ 监控异常: {e}")
def _collect_metrics(self):
"""收集性能指标"""
try:
cpu_usage = self._get_cpu_usage()
memory_usage = self._get_memory_usage()
# 记录性能指标
self.logger.debug(f"📊 性能指标 - CPU: {cpu_usage}%, 内存: {memory_usage}%, 用户: {self.active_users}")
# 如果资源使用率过高,发出警告
if cpu_usage > 90:
self.logger.warning(f"⚠️ CPU使用率过高: {cpu_usage}%")
if memory_usage > 90:
self.logger.warning(f"⚠️ 内存使用率过高: {memory_usage}%")
except Exception as e:
self.logger.error(f"❌ 收集指标失败: {e}")
def _get_cpu_usage(self) -> float:
"""获取CPU使用率"""
try:
return psutil.cpu_percent(interval=1)
except:
return 0.0
def _get_memory_usage(self) -> float:
"""获取内存使用率"""
try:
return psutil.virtual_memory().percent
except:
return 0.0
def update_user_count(self, count: int):
"""更新活跃用户数"""
self.active_users = count
def increment_request_count(self, count: int = 1):
"""增加请求计数"""
self.total_requests += count
def get_status(self) -> Dict[str, Any]:
"""获取节点状态"""
uptime = (datetime.now() - self.start_time).total_seconds() if self.start_time else 0
return {
"node_id": self.node_id,
"host": self.host,
"port": self.port,
"status": "running" if self.is_running else "stopped",
"uptime": uptime,
"cpu_usage": self._get_cpu_usage(),
"memory_usage": self._get_memory_usage(),
"active_users": self.active_users,
"total_requests": self.total_requests,
"start_time": self.start_time.isoformat() if self.start_time else None
}
# 使用示例
def create_worker_node(master_host: str) -> WorkerNode:
"""创建Worker节点"""
worker = WorkerNode(master_host)
worker.start()
return worker2. 负载均衡策略
# core/load_balancer.py
"""
负载均衡器 - 智能的任务分发器
确保压测任务合理分配到各个节点
"""
from typing import List, Dict, Any, Optional
from enum import Enum
import random
import hashlib
class LoadBalanceStrategy(Enum):
"""负载均衡策略"""
ROUND_ROBIN = "round_robin" # 轮询
LEAST_CONNECTIONS = "least_conn" # 最少连接
WEIGHTED_ROUND_ROBIN = "weighted_rr" # 加权轮询
LEAST_RESPONSE_TIME = "least_rt" # 最短响应时间
HASH = "hash" # 哈希
RANDOM = "random" # 随机
class LoadBalancer:
"""
负载均衡器
就像一个智能的交通指挥员,
确保每条道路都不会拥堵
"""
def __init__(self, strategy: LoadBalanceStrategy = LoadBalanceStrategy.LEAST_CONNECTIONS):
self.strategy = strategy
self.round_robin_index = 0
self.worker_weights: Dict[str, float] = {}
self.worker_response_times: Dict[str, float] = {}
def select_worker(self, workers: List[WorkerNode], request_info: Dict[str, Any] = None) -> Optional[WorkerNode]:
"""
选择Worker节点
Args:
workers: 可用的Worker节点列表
request_info: 请求信息(用于哈希策略)
Returns:
选中的Worker节点
"""
if not workers:
return None
if self.strategy == LoadBalanceStrategy.ROUND_ROBIN:
return self._round_robin_select(workers)
elif self.strategy == LoadBalanceStrategy.LEAST_CONNECTIONS:
return self._least_connections_select(workers)
elif self.strategy == LoadBalanceStrategy.WEIGHTED_ROUND_ROBIN:
return self._weighted_round_robin_select(workers)
elif self.strategy == LoadBalanceStrategy.LEAST_RESPONSE_TIME:
return self._least_response_time_select(workers)
elif self.strategy == LoadBalanceStrategy.HASH:
return self._hash_select(workers, request_info)
elif self.strategy == LoadBalanceStrategy.RANDOM:
return self._random_select(workers)
else:
return workers[0] # 默认选择第一个
def _round_robin_select(self, workers: List[WorkerNode]) -> WorkerNode:
"""轮询选择"""
worker = workers[self.round_robin_index % len(workers)]
self.round_robin_index += 1
return worker
def _least_connections_select(self, workers: List[WorkerNode]) -> WorkerNode:
"""最少连接选择"""
return min(workers, key=lambda w: w.active_users)
def _weighted_round_robin_select(self, workers: List[WorkerNode]) -> WorkerNode:
"""加权轮询选择"""
# 根据节点性能计算权重
total_weight = 0
for worker in workers:
weight = self._calculate_weight(worker)
self.worker_weights[worker.node_id] = weight
total_weight += weight
if total_weight == 0:
return workers[0]
# 根据权重选择
random_value = random.uniform(0, total_weight)
current_weight = 0
for worker in workers:
current_weight += self.worker_weights.get(worker.node_id, 1)
if random_value <= current_weight:
return worker
return workers[-1]
def _least_response_time_select(self, workers: List[WorkerNode]) -> WorkerNode:
"""最短响应时间选择"""
# 如果没有响应时间数据,使用最少连接策略
if not self.worker_response_times:
return self._least_connections_select(workers)
return min(workers, key=lambda w: self.worker_response_times.get(w.node_id, float('inf')))
def _hash_select(self, workers: List[WorkerNode], request_info: Dict[str, Any]) -> WorkerNode:
"""哈希选择(一致性哈希)"""
if not request_info or 'user_id' not in request_info:
return self._round_robin_select(workers)
user_id = str(request_info['user_id'])
hash_value = int(hashlib.md5(user_id.encode()).hexdigest(), 16)
index = hash_value % len(workers)
return workers[index]
def _random_select(self, workers: List[WorkerNode]) -> WorkerNode:
"""随机选择"""
return random.choice(workers)
def _calculate_weight(self, worker: WorkerNode) -> float:
"""计算节点权重"""
# 权重 = 1 / (CPU使用率 * 0.4 + 内存使用率 * 0.3 + 连接数比例 * 0.3)
cpu_factor = max(worker.cpu_usage / 100.0, 0.1)
memory_factor = max(worker.memory_usage / 100.0, 0.1)
connection_factor = max(worker.active_users / 1000.0, 0.1) # 假设1000为满载
load_score = cpu_factor * 0.4 + memory_factor * 0.3 + connection_factor * 0.3
return 1.0 / load_score
def update_response_time(self, worker_id: str, response_time: float):
"""更新Worker响应时间"""
self.worker_response_times[worker_id] = response_time
def get_strategy_info(self) -> Dict[str, Any]:
"""获取策略信息"""
return {
"strategy": self.strategy.value,
"round_robin_index": self.round_robin_index,
"worker_weights": self.worker_weights.copy(),
"worker_response_times": self.worker_response_times.copy()
}
# 全局负载均衡器实例
load_balancer = LoadBalancer()集群部署与管理
1. Docker容器化部署
# Dockerfile.master - Master节点镜像
FROM python:3.9-slim
WORKDIR /app
# 安装系统依赖
RUN apt-get update && apt-get install -y \
gcc \
&& rm -rf /var/lib/apt/lists/*
# 复制依赖文件
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制应用代码
COPY . .
# 暴露端口
EXPOSE 8089 5557
# 启动Master节点
CMD ["python", "-m", "locust", "--master", "--web-host", "0.0.0.0", "--web-port", "8089", "--master-bind-port", "5557"]# Dockerfile.worker - Worker节点镜像
FROM python:3.9-slim
WORKDIR /app
# 安装系统依赖
RUN apt-get update && apt-get install -y \
gcc \
&& rm -rf /var/lib/apt/lists/*
# 复制依赖文件
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制应用代码
COPY . .
# 启动Worker节点
CMD ["python", "-m", "locust", "--worker", "--master-host", "${MASTER_HOST:-master}"]# docker-compose.yml - 集群编排
version: '3.8'
services:
master:
build:
context: .
dockerfile: Dockerfile.master
ports:
- "8089:8089"
- "5557:5557"
environment:
- LOCUST_HOST=https://httpbin.org
volumes:
- ./logs:/app/logs
- ./output:/app/output
networks:
- locust-network
worker1:
build:
context: .
dockerfile: Dockerfile.worker
environment:
- MASTER_HOST=master
- LOCUST_HOST=https://httpbin.org
depends_on:
- master
networks:
- locust-network
worker2:
build:
context: .
dockerfile: Dockerfile.worker
environment:
- MASTER_HOST=master
- LOCUST_HOST=https://httpbin.org
depends_on:
- master
networks:
- locust-network
worker3:
build:
context: .
dockerfile: Dockerfile.worker
environment:
- MASTER_HOST=master
- LOCUST_HOST=https://httpbin.org
depends_on:
- master
networks:
- locust-network
# 监控服务
prometheus:
image: prom/prometheus:latest
ports:
- "9090:9090"
volumes:
- ./monitoring/prometheus.yml:/etc/prometheus/prometheus.yml
networks:
- locust-network
grafana:
image: grafana/grafana:latest
ports:
- "3000:3000"
environment:
- GF_SECURITY_ADMIN_PASSWORD=admin
volumes:
- grafana-storage:/var/lib/grafana
- ./monitoring/grafana/dashboards:/etc/grafana/provisioning/dashboards
- ./monitoring/grafana/datasources:/etc/grafana/provisioning/datasources
networks:
- locust-network
networks:
locust-network:
driver: bridge
volumes:
grafana-storage:2. Kubernetes部署
# k8s/namespace.yaml
apiVersion: v1
kind: Namespace
metadata:
name: locust-cluster
---
# k8s/master-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: locust-master
namespace: locust-cluster
spec:
replicas: 1
selector:
matchLabels:
app: locust-master
template:
metadata:
labels:
app: locust-master
spec:
containers:
- name: locust-master
image: locust-framework:latest
ports:
- containerPort: 8089
- containerPort: 5557
env:
- name: LOCUST_HOST
value: "https://httpbin.org"
command: ["python", "-m", "locust"]
args: ["--master", "--web-host", "0.0.0.0", "--web-port", "8089", "--master-bind-port", "5557"]
resources:
requests:
memory: "512Mi"
cpu: "500m"
limits:
memory: "1Gi"
cpu: "1000m"
---
# k8s/master-service.yaml
apiVersion: v1
kind: Service
metadata:
name: locust-master-service
namespace: locust-cluster
spec:
selector:
app: locust-master
ports:
- name: web
port: 8089
targetPort: 8089
- name: master
port: 5557
targetPort: 5557
type: LoadBalancer
---
# k8s/worker-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: locust-worker
namespace: locust-cluster
spec:
replicas: 5 # 可以根据需要调整Worker数量
selector:
matchLabels:
app: locust-worker
template:
metadata:
labels:
app: locust-worker
spec:
containers:
- name: locust-worker
image: locust-framework:latest
env:
- name: LOCUST_HOST
value: "https://httpbin.org"
- name: MASTER_HOST
value: "locust-master-service"
command: ["python", "-m", "locust"]
args: ["--worker", "--master-host", "$(MASTER_HOST)"]
resources:
requests:
memory: "256Mi"
cpu: "250m"
limits:
memory: "512Mi"
cpu: "500m"
---
# k8s/hpa.yaml - 水平自动扩缩容
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: locust-worker-hpa
namespace: locust-cluster
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: locust-worker
minReplicas: 2
maxReplicas: 20
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 803. 集群监控系统
# core/cluster_monitor.py
"""
集群监控系统 - 集群的"健康管家"
实时监控集群状态和性能指标
"""
import time
import threading
import json
from typing import Dict, List, Any
from datetime import datetime, timedelta
from dataclasses import dataclass
import requests
@dataclass
class ClusterMetrics:
"""集群指标"""
timestamp: datetime
total_workers: int
healthy_workers: int
total_users: int
total_rps: float
avg_response_time: float
error_rate: float
cpu_usage: float
memory_usage: float
class ClusterMonitor:
"""
集群监控器
就像一个全天候的监控中心,
时刻关注集群的健康状态
"""
def __init__(self, cluster_manager):
self.cluster_manager = cluster_manager
self.metrics_history: List[ClusterMetrics] = []
self.alerts: List[Dict[str, Any]] = []
self.is_monitoring = False
self.monitor_thread = None
self.monitor_interval = 10 # 监控间隔(秒)
# 告警阈值
self.alert_thresholds = {
'cpu_usage': 80.0,
'memory_usage': 85.0,
'error_rate': 5.0,
'response_time': 5000.0, # 毫秒
'worker_failure_rate': 20.0 # 百分比
}
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志"""
import logging
logger = logging.getLogger("ClusterMonitor")
logger.setLevel(logging.INFO)
return logger
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._collect_metrics()
if metrics:
self.metrics_history.append(metrics)
# 保留最近1小时的数据
cutoff_time = datetime.now() - timedelta(hours=1)
self.metrics_history = [m for m in self.metrics_history if m.timestamp > cutoff_time]
# 检查告警
self._check_alerts(metrics)
time.sleep(self.monitor_interval)
except Exception as e:
self.logger.error(f"❌ 监控异常: {e}")
def _collect_metrics(self) -> ClusterMetrics:
"""收集集群指标"""
try:
cluster_status = self.cluster_manager.get_cluster_status()
cluster_info = cluster_status['cluster_info']
return ClusterMetrics(
timestamp=datetime.now(),
total_workers=cluster_info['total_workers'],
healthy_workers=cluster_info['healthy_workers'],
total_users=cluster_info['total_active_users'],
total_rps=cluster_info.get('total_rps', 0),
avg_response_time=cluster_info.get('avg_response_time', 0),
error_rate=cluster_info.get('error_rate', 0),
cpu_usage=cluster_info['avg_cpu_usage'],
memory_usage=cluster_info['avg_memory_usage']
)
except Exception as e:
self.logger.error(f"❌ 收集指标失败: {e}")
return None
def _check_alerts(self, metrics: ClusterMetrics):
"""检查告警条件"""
alerts = []
# CPU使用率告警
if metrics.cpu_usage > self.alert_thresholds['cpu_usage']:
alerts.append({
'type': 'cpu_high',
'level': 'warning',
'message': f'集群CPU使用率过高: {metrics.cpu_usage:.1f}%',
'value': metrics.cpu_usage,
'threshold': self.alert_thresholds['cpu_usage']
})
# 内存使用率告警
if metrics.memory_usage > self.alert_thresholds['memory_usage']:
alerts.append({
'type': 'memory_high',
'level': 'warning',
'message': f'集群内存使用率过高: {metrics.memory_usage:.1f}%',
'value': metrics.memory_usage,
'threshold': self.alert_thresholds['memory_usage']
})
# 错误率告警
if metrics.error_rate > self.alert_thresholds['error_rate']:
alerts.append({
'type': 'error_rate_high',
'level': 'critical',
'message': f'错误率过高: {metrics.error_rate:.1f}%',
'value': metrics.error_rate,
'threshold': self.alert_thresholds['error_rate']
})
# 响应时间告警
if metrics.avg_response_time > self.alert_thresholds['response_time']:
alerts.append({
'type': 'response_time_high',
'level': 'warning',
'message': f'平均响应时间过长: {metrics.avg_response_time:.1f}ms',
'value': metrics.avg_response_time,
'threshold': self.alert_thresholds['response_time']
})
# Worker故障率告警
if metrics.total_workers > 0:
failure_rate = (metrics.total_workers - metrics.healthy_workers) / metrics.total_workers * 100
if failure_rate > self.alert_thresholds['worker_failure_rate']:
alerts.append({
'type': 'worker_failure_high',
'level': 'critical',
'message': f'Worker故障率过高: {failure_rate:.1f}%',
'value': failure_rate,
'threshold': self.alert_thresholds['worker_failure_rate']
})
# 处理告警
for alert in alerts:
self._handle_alert(alert)
def _handle_alert(self, alert: Dict[str, Any]):
"""处理告警"""
alert['timestamp'] = datetime.now().isoformat()
self.alerts.append(alert)
# 记录告警日志
level_emoji = {'warning': '⚠️', 'critical': '🚨', 'info': 'ℹ️'}
emoji = level_emoji.get(alert['level'], '📢')
self.logger.warning(f"{emoji} {alert['message']}")
# 这里可以添加更多告警处理逻辑,如:
# - 发送邮件通知
# - 发送钉钉/企业微信消息
# - 调用Webhook
# - 自动扩容
# 保留最近100条告警
if len(self.alerts) > 100:
self.alerts = self.alerts[-100:]
def get_metrics_summary(self, duration_minutes: int = 30) -> Dict[str, Any]:
"""获取指标摘要"""
cutoff_time = datetime.now() - timedelta(minutes=duration_minutes)
recent_metrics = [m for m in self.metrics_history if m.timestamp > cutoff_time]
if not recent_metrics:
return {}
# 计算统计值
cpu_values = [m.cpu_usage for m in recent_metrics]
memory_values = [m.memory_usage for m in recent_metrics]
response_times = [m.avg_response_time for m in recent_metrics]
error_rates = [m.error_rate for m in recent_metrics]
return {
'duration_minutes': duration_minutes,
'data_points': len(recent_metrics),
'cpu_usage': {
'avg': sum(cpu_values) / len(cpu_values),
'max': max(cpu_values),
'min': min(cpu_values)
},
'memory_usage': {
'avg': sum(memory_values) / len(memory_values),
'max': max(memory_values),
'min': min(memory_values)
},
'response_time': {
'avg': sum(response_times) / len(response_times),
'max': max(response_times),
'min': min(response_times)
},
'error_rate': {
'avg': sum(error_rates) / len(error_rates),
'max': max(error_rates),
'min': min(error_rates)
},
'latest_metrics': recent_metrics[-1].__dict__ if recent_metrics else None
}
def get_recent_alerts(self, count: int = 10) -> List[Dict[str, Any]]:
"""获取最近的告警"""
return self.alerts[-count:] if self.alerts else []
# 全局集群监控器实例
cluster_monitor = ClusterMonitor(cluster_manager)实战部署指南
1. 快速部署脚本
#!/bin/bash
# deploy.sh - 一键部署脚本
set -e
echo "🚀 开始部署Locust分布式集群..."
# 检查Docker环境
if ! command -v docker &> /dev/null; then
echo "❌ Docker未安装,请先安装Docker"
exit 1
fi
if ! command -v docker-compose &> /dev/null; then
echo "❌ Docker Compose未安装,请先安装Docker Compose"
exit 1
fi
# 构建镜像
echo "📦 构建Docker镜像..."
docker build -t locust-framework:latest .
# 启动集群
echo "🔧 启动集群服务..."
docker-compose up -d
# 等待服务启动
echo "⏳ 等待服务启动..."
sleep 30
# 检查服务状态
echo "🔍 检查服务状态..."
docker-compose ps
# 显示访问信息
echo "✅ 部署完成!"
echo "📊 Locust Web界面: http://localhost:8089"
echo "📈 Grafana监控面板: http://localhost:3000 (admin/admin)"
echo "🔧 Prometheus: http://localhost:9090"
echo "🎯 开始压测:"
echo "1. 访问 http://localhost:8089"
echo "2. 设置用户数和启动速率"
echo "3. 点击 'Start swarming' 开始测试"2. 扩容脚本
#!/bin/bash
# scale.sh - 集群扩容脚本
WORKER_COUNT=${1:-5}
echo "📈 扩容Worker节点到 $WORKER_COUNT 个..."
# Docker Compose扩容
docker-compose up -d --scale worker=$WORKER_COUNT
echo "✅ 扩容完成!当前Worker数量: $WORKER_COUNT"
# 显示当前状态
docker-compose ps | grep worker总结
分布式压测架构就像建造一支现代化的军队,需要统一的指挥、协调的作战、高效的通信。通过这篇文章,我们深入了解了:
- 架构设计:Master-Worker模式的现代化实现
- 集群管理:智能的节点管理和任务分发
- 负载均衡:多种策略确保资源合理利用
- 容器化部署:Docker和Kubernetes的部署方案
- 监控告警:全方位的集群健康监控
这套分布式压测架构不仅能够突破单机限制,更重要的是提供了企业级的可靠性和可扩展性,让大规模性能测试变得简单而可控。
下一篇文章,我们将探讨性能监控与智能告警系统,看看如何构建完善的可观测性体系。
推荐阅读
