
Locust分布式压测实战
大约 7 分钟
Locust分布式压测实战
🌐 单机压测不够劲?是时候组建"压测军团"了!
就像组建复仇者联盟一样,我们要让多台机器协同作战,发挥集群的威力。
今天分享分布式压测的实战经验,让你的压测能力翻倍提升!💪
🏗️ 分布式架构原理
Master-Worker模式
Locust的分布式架构采用经典的Master-Worker模式:
📱 Web界面
↓
🎯 Master节点
↙️ ↘️ ↘️
Worker1 Worker2 Worker3
↓ ↓ ↓
目标系统 目标系统 目标系统角色分工:
- Master节点:负责协调、统计汇总、Web界面
- Worker节点:负责实际发送请求、执行压测脚本
为什么需要分布式?
作为测试工程师,我遇到过这些场景:
- 单机性能瓶颈:一台机器CPU/网络带宽有限
- 大规模并发需求:需要模拟上万用户同时在线
- 地理位置分布:模拟不同地区用户访问
- 资源隔离:避免压测影响其他服务
🚀 快速搭建分布式环境
环境准备
假设我们有3台机器:
- Master机器:192.168.1.100
- Worker机器1:192.168.1.101
- Worker机器2:192.168.1.102
1. 准备压测脚本
在所有机器上创建相同的脚本文件 distributed_test.py:
from locust import HttpUser, task, between
import random
class DistributedUser(HttpUser):
wait_time = between(1, 3)
host = "https://httpbin.org"
def on_start(self):
"""用户启动时执行"""
# 获取当前Worker的标识
import socket
self.worker_id = socket.gethostname()
print(f"🤖 Worker {self.worker_id} 用户启动")
@task(3)
def get_request(self):
"""GET请求测试"""
self.client.get("/get", params={
"worker": self.worker_id,
"user": random.randint(1, 1000)
})
@task(2)
def post_request(self):
"""POST请求测试"""
self.client.post("/post", json={
"worker": self.worker_id,
"timestamp": time.time(),
"data": "distributed test"
})
@task(1)
def json_request(self):
"""JSON数据请求"""
response = self.client.get("/json")
if response.status_code == 200:
print(f"✅ Worker {self.worker_id} JSON请求成功")2. 启动Master节点
在Master机器(192.168.1.100)上执行:
# 启动Master节点
locust -f distributed_test.py --master
# 指定Web界面端口(可选)
locust -f distributed_test.py --master --web-port=8089
# 指定Master监听地址(可选)
locust -f distributed_test.py --master --master-bind-host=0.0.0.0 --master-bind-port=5557启动后会看到:
[2024-01-01 10:00:00,000] INFO/locust.main: Starting web interface at http://0.0.0.0:8089
[2024-01-01 10:00:00,000] INFO/locust.main: Starting Locust 2.x.x3. 启动Worker节点
在Worker机器上执行:
# Worker1 (192.168.1.101)
locust -f distributed_test.py --worker --master-host=192.168.1.100
# Worker2 (192.168.1.102)
locust -f distributed_test.py --worker --master-host=192.168.1.100
# 指定Master端口(如果修改了默认端口)
locust -f distributed_test.py --worker --master-host=192.168.1.100 --master-port=5557成功连接后会看到:
[2024-01-01 10:00:05,000] INFO/locust.runners: Connected to locust master4. 验证集群状态
打开浏览器访问:http://192.168.1.100:8089
在Web界面右上角可以看到Worker数量,比如:Workers: 2
🔧 高级配置技巧
环境变量配置
创建配置文件 config.py:
import os
class Config:
# Master配置
MASTER_HOST = os.getenv('LOCUST_MASTER_HOST', '192.168.1.100')
MASTER_PORT = int(os.getenv('LOCUST_MASTER_PORT', 5557))
WEB_PORT = int(os.getenv('LOCUST_WEB_PORT', 8089))
# 测试目标配置
TARGET_HOST = os.getenv('TARGET_HOST', 'https://httpbin.org')
# Worker配置
WORKER_ID = os.getenv('WORKER_ID', 'unknown')
# 在压测脚本中使用
from config import Config
class ConfigurableUser(HttpUser):
wait_time = between(1, 2)
host = Config.TARGET_HOST
def on_start(self):
self.worker_id = Config.WORKER_ID
print(f"🔧 Worker {self.worker_id} 配置加载完成")Docker化部署
创建 Dockerfile:
FROM python:3.9-slim
WORKDIR /app
# 安装依赖
COPY requirements.txt .
RUN pip install -r requirements.txt
# 复制脚本
COPY . .
# 暴露端口
EXPOSE 8089 5557
# 启动脚本
COPY entrypoint.sh .
RUN chmod +x entrypoint.sh
ENTRYPOINT ["./entrypoint.sh"]创建 entrypoint.sh:
#!/bin/bash
if [ "$LOCUST_MODE" = "master" ]; then
echo "🎯 启动Master节点..."
locust -f distributed_test.py --master --web-port=8089
elif [ "$LOCUST_MODE" = "worker" ]; then
echo "🤖 启动Worker节点..."
locust -f distributed_test.py --worker --master-host=$MASTER_HOST
else
echo "❌ 请设置LOCUST_MODE环境变量 (master/worker)"
exit 1
fi使用Docker Compose部署:
# docker-compose.yml
version: '3.8'
services:
locust-master:
build: .
ports:
- "8089:8089"
- "5557:5557"
environment:
- LOCUST_MODE=master
volumes:
- ./logs:/app/logs
locust-worker-1:
build: .
environment:
- LOCUST_MODE=worker
- MASTER_HOST=locust-master
- WORKER_ID=worker-1
depends_on:
- locust-master
locust-worker-2:
build: .
environment:
- LOCUST_MODE=worker
- MASTER_HOST=locust-master
- WORKER_ID=worker-2
depends_on:
- locust-master📊 数据同步与共享
Worker间数据共享
在分布式环境中,有时需要在Worker间共享数据:
from locust import HttpUser, task, events
from locust.runners import MasterRunner, WorkerRunner
import json
# 全局数据存储
shared_data = {
"user_tokens": [],
"test_data": []
}
def setup_shared_data(environment, msg, **kwargs):
"""接收Master分发的共享数据"""
global shared_data
shared_data.update(msg.data)
print(f"📦 收到共享数据: {len(shared_data['user_tokens'])} 个token")
@events.init.add_listener
def on_locust_init(environment, **kwargs):
"""初始化消息处理器"""
if not isinstance(environment.runner, MasterRunner):
environment.runner.register_message("shared_data", setup_shared_data)
@events.test_start.add_listener
def on_test_start(environment, **kwargs):
"""测试开始时分发数据"""
if isinstance(environment.runner, MasterRunner):
# 准备共享数据
tokens = [f"token_{i}" for i in range(1000)]
test_data = [{"id": i, "name": f"user_{i}"} for i in range(100)]
# 分发给所有Worker
data_to_share = {
"user_tokens": tokens,
"test_data": test_data
}
for worker in environment.runner.clients:
environment.runner.send_message("shared_data", data_to_share, worker)
class SharedDataUser(HttpUser):
wait_time = between(1, 2)
host = "https://httpbin.org"
@task
def use_shared_token(self):
"""使用共享的token"""
if shared_data["user_tokens"]:
token = shared_data["user_tokens"].pop(0)
self.client.get("/bearer", headers={
"Authorization": f"Bearer {token}"
})
# 使用完后放回队列(如果需要重复使用)
shared_data["user_tokens"].append(token)唯一数据分配
确保每个Worker使用不同的测试数据:
from locust import HttpUser, task, events
from locust.runners import WorkerRunner
import uuid
class UniqueDataUser(HttpUser):
wait_time = between(1, 2)
host = "https://httpbin.org"
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
# 为每个用户生成唯一标识
self.unique_id = str(uuid.uuid4())
# 基于Worker ID生成数据范围
if isinstance(self.environment.runner, WorkerRunner):
worker_index = hash(self.environment.runner.worker_index) % 1000
self.data_range_start = worker_index * 1000
self.data_range_end = (worker_index + 1) * 1000
else:
self.data_range_start = 0
self.data_range_end = 1000
@task
def unique_request(self):
"""使用唯一数据的请求"""
unique_value = self.data_range_start + hash(self.unique_id) % 1000
self.client.post("/post", json={
"user_id": self.unique_id,
"unique_value": unique_value,
"worker_range": f"{self.data_range_start}-{self.data_range_end}"
})🔍 监控与调试
分布式监控脚本
from locust import events
import psutil
import time
import json
class DistributedMonitor:
"""分布式环境监控"""
def __init__(self):
self.worker_stats = {}
self.start_time = time.time()
def collect_worker_stats(self):
"""收集Worker统计信息"""
import socket
worker_id = socket.gethostname()
stats = {
"worker_id": worker_id,
"timestamp": time.time(),
"cpu_percent": psutil.cpu_percent(),
"memory_percent": psutil.virtual_memory().percent,
"network_io": psutil.net_io_counters()._asdict(),
"uptime": time.time() - self.start_time
}
return stats
# 全局监控实例
monitor = DistributedMonitor()
@events.request_success.add_listener
def on_request_success(request_type, name, response_time, response_length, **kwargs):
"""请求成功时收集统计"""
stats = monitor.collect_worker_stats()
# 可以发送到监控系统或日志
if response_time > 1000: # 响应时间超过1秒
print(f"⚠️ 慢请求告警: {name} - {response_time}ms - Worker: {stats['worker_id']}")
@events.user_error.add_listener
def on_user_error(user_instance, exception, tb, **kwargs):
"""用户错误时记录详细信息"""
stats = monitor.collect_worker_stats()
error_info = {
"worker_id": stats["worker_id"],
"error": str(exception),
"timestamp": time.time()
}
print(f"❌ Worker错误: {json.dumps(error_info, ensure_ascii=False)}")健康检查脚本
#!/bin/bash
# health_check.sh - 检查分布式集群健康状态
MASTER_HOST="192.168.1.100"
MASTER_PORT="8089"
echo "🔍 检查Locust集群健康状态..."
# 检查Master节点
echo "检查Master节点..."
if curl -s "http://${MASTER_HOST}:${MASTER_PORT}/stats/requests" > /dev/null; then
echo "✅ Master节点正常"
else
echo "❌ Master节点异常"
exit 1
fi
# 获取Worker数量
WORKER_COUNT=$(curl -s "http://${MASTER_HOST}:${MASTER_PORT}/stats/requests" | jq '.workers | length')
echo "📊 当前Worker数量: ${WORKER_COUNT}"
if [ "$WORKER_COUNT" -lt 2 ]; then
echo "⚠️ Worker数量不足"
exit 1
fi
echo "🎉 集群状态正常"🎯 最佳实践
1. 网络优化
# 优化网络配置
class OptimizedDistributedUser(HttpUser):
wait_time = between(1, 2)
host = "https://httpbin.org"
def on_start(self):
"""优化网络设置"""
# 设置连接池
from requests.adapters import HTTPAdapter
adapter = HTTPAdapter(
pool_connections=20,
pool_maxsize=50,
pool_block=False
)
self.client.mount("http://", adapter)
self.client.mount("https://", adapter)
# 设置超时
self.client.timeout = (5, 30)2. 资源监控
# 监控脚本 monitor.sh
#!/bin/bash
while true; do
echo "=== $(date) ==="
echo "CPU使用率: $(top -bn1 | grep "Cpu(s)" | awk '{print $2}' | cut -d'%' -f1)"
echo "内存使用率: $(free | grep Mem | awk '{printf("%.2f%%", $3/$2 * 100.0)}')"
echo "网络连接数: $(netstat -an | grep ESTABLISHED | wc -l)"
echo "Locust进程: $(ps aux | grep locust | grep -v grep | wc -l)"
echo "---"
sleep 10
done3. 故障恢复
# 自动重连机制
import time
import subprocess
import sys
def auto_restart_worker():
"""Worker自动重启机制"""
max_retries = 5
retry_count = 0
while retry_count < max_retries:
try:
# 启动Worker
cmd = [
"locust",
"-f", "distributed_test.py",
"--worker",
"--master-host", "192.168.1.100"
]
process = subprocess.Popen(cmd)
process.wait()
except KeyboardInterrupt:
print("🛑 手动停止")
break
except Exception as e:
print(f"❌ Worker异常: {e}")
retry_count += 1
print(f"🔄 {retry_count}/{max_retries} 次重试...")
time.sleep(10)
print("💀 Worker重启次数超限,退出")
if __name__ == "__main__":
auto_restart_worker()🚀 性能调优建议
1. 合理分配资源
- Master节点:2-4核CPU,4-8GB内存
- Worker节点:根据并发需求配置,建议每1000并发用户配置2核CPU
2. 网络优化
- 使用千兆网络
- 调整TCP参数
- 监控网络带宽使用情况
3. 监控指标
重点关注:
- Worker连接状态
- 网络延迟
- 系统资源使用率
- 请求成功率
🎯 总结
分布式压测就像指挥一支军队,需要:
- 统一指挥:Master节点协调全局
- 分工明确:Worker节点各司其职
- 通信顺畅:网络配置要优化
- 监控到位:实时掌握集群状态
掌握了分布式压测,你就能应对更大规模的性能测试挑战!🏆
下一步可以学习性能调优和问题诊断,让你的压测技能更加全面。
