Python自动化测试青果代理全流程:从API调用到稳定性监控脚本分享
构建企业级代理IP自动化测试与监控平台:Python实战指南
最近和几个做数据采集和风控的朋友聊天,大家普遍反映一个痛点:市面上各种代理IP服务商都说自己“稳定高速”,但实际用起来效果参差不齐。团队花了不少钱采购,结果关键时刻掉链子,不是IP大量失效就是响应慢得让人抓狂。更头疼的是,缺乏一套系统化的方法来持续评估这些代理的质量,每次都是出了问题才手忙脚乱地临时测试。
这让我意识到,对于技术团队来说,单纯依赖服务商的口头承诺或一次性评测是远远不够的。我们需要的是一个能够持续运行、自动评估、智能预警的代理IP质量监控体系。这不仅能帮我们筛选出真正靠谱的服务商,还能在日常使用中及时发现性能波动,避免业务受到影响。
今天,我就把自己在多个项目中搭建这类系统的经验整理出来,分享一套完整的Python自动化测试方案。这套方案不绑定任何特定服务商,你可以用它来测试任何提供API接口的代理IP服务。我们会从最基础的API调用开始,一步步构建出包含有效性验证、性能统计、稳定性监控甚至可视化报表的完整工具链。所有代码都是模块化设计,你可以直接集成到现有系统中,或者根据团队需求进行定制。
1. 系统架构设计与核心模块规划
在开始写代码之前,我们先要理清楚整个系统需要哪些功能模块。一个完整的代理IP自动化测试平台,远不止“获取IP-测试连通性”这么简单。我们需要考虑的是全生命周期管理。
1.1 核心功能模块分解
我把整个系统划分为五个核心模块,每个模块都有明确的职责:
- IP获取与调度模块:负责从不同服务商的API获取IP,并实现智能调度策略
- 质量检测引擎:执行多种维度的测试,包括连通性、匿名度、地理位置等
- 性能数据收集器:记录每次测试的详细指标,建立历史数据库
- 数据分析与评分系统:基于历史数据计算每个IP的综合评分
- 监控与告警模块:实时监控IP池状态,异常时自动触发告警
这五个模块的关系可以用下面的流程图来理解(虽然我们不能用mermaid,但可以用文字描述):IP获取模块将原始IP送入检测引擎,检测结果被数据收集器记录,分析系统定期计算评分,监控模块则持续关注评分变化和异常情况。
1.2 技术栈选型与依赖安装
我们选择Python作为开发语言,主要是考虑到它的生态丰富和快速原型能力。以下是需要安装的核心库:
# 基础请求与并发处理
pip install requests aiohttp httpx
# 数据存储与分析
pip install pandas numpy sqlalchemy
# 定时任务调度
pip install schedule apscheduler
# 数据可视化(用于生成报表)
pip install matplotlib seaborn plotly
# 配置文件管理
pip install python-dotenv yaml
# 日志记录
pip install loguru
如果你打算使用异步IO来提高测试效率(特别是需要同时测试大量IP时),我强烈推荐aiohttp和asyncio的组合。对于中小规模的测试,concurrent.futures的线程池也足够用了。
注意:在实际部署时,建议使用虚拟环境(venv或conda)来管理依赖,避免与系统Python环境冲突。特别是生产环境,一定要固定版本号。
2. 代理IP获取与智能调度实现
代理IP的获取是整套系统的入口。不同的服务商提供的API接口差异很大,我们需要设计一个统一的抽象层来兼容它们。
2.1 多服务商API统一封装
首先,我们定义一个基础类来处理通用的HTTP请求和错误重试:
import requests
import time
import logging
from typing import Optional, Dict, Any
from dataclasses import dataclass
from abc import ABC, abstractmethod
@dataclass
class ProxyIP:
"""代理IP数据类"""
ip: str
port: int
protocol: str # http, https, socks5
location: Optional[str] = None
expire_time: Optional[float] = None
source: str = "unknown" # 服务商名称
class BaseProxyFetcher(ABC):
"""代理获取器基类"""
def __init__(self, api_key: str, max_retries: int = 3):
self.api_key = api_key
self.max_retries = max_retries
self.session = requests.Session()
self.session.headers.update({
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
})
@abstractmethod
async def fetch_batch(self, count: int = 10) -> list[ProxyIP]:
"""批量获取代理IP"""
pass
def _make_request(self, url: str, params: Optional[Dict] = None) -> Optional[Dict]:
"""带重试的HTTP请求"""
for attempt in range(self.max_retries):
try:
response = self.session.get(url, params=params, timeout=10)
response.raise_for_status()
return response.json()
except requests.exceptions.RequestException as e:
logging.warning(f"请求失败 (尝试 {attempt + 1}/{self.max_retries}): {e}")
if attempt < self.max_retries - 1:
time.sleep(2 ** attempt) # 指数退避
return None
接下来,我们实现几个具体服务商的获取器。以常见的格式为例:
class CommonProxyFetcher(BaseProxyFetcher):
"""通用格式代理获取器(适用于多数服务商)"""
def __init__(self, api_url: str, api_key: str, format_type: str = "json"):
super().__init__(api_key)
self.api_url = api_url
self.format_type = format_type
async def fetch_batch(self, count: int = 10) -> list[ProxyIP]:
params = {
'key': self.api_key,
'num': count,
'format': self.format_type,
'distinct': 'true' # 去重
}
data = self._make_request(self.api_url, params)
if not data:
return []
proxies = []
if self.format_type == "json":
for item in data.get("data", []):
proxy = ProxyIP(
ip=item["ip"],
port=item["port"],
protocol=item.get("protocol", "http"),
location=item.get("location"),
expire_time=time.time() + item.get("expire", 300),
source=self.__class__.__name__
)
proxies.append(proxy)
elif self.format_type == "text":
# 处理文本格式:ip:port 每行一个
lines = data.strip().split('\n')
for line in lines:
if ':' in line:
ip, port = line.strip().split(':')
proxies.append(ProxyIP(
ip=ip, port=int(port), protocol="http",
source=self.__class__.__name__
))
return proxies
2.2 智能调度策略
获取到IP后,我们需要一个调度器来管理这些IP。一个好的调度策略能显著提升整体效率。
class ProxyScheduler:
"""代理IP调度器"""
def __init__(self):
self.proxy_pool: Dict[str, ProxyIP] = {} # ip:port -> ProxyIP
self.performance_stats: Dict[str, dict] = {} # 性能统计
self.blacklist: Set[str] = set() # 黑名单
self.last_used: Dict[str, float] = {} # 最后使用时间
def add_proxy(self, proxy: ProxyIP):
"""添加代理到池中"""
key = f"{proxy.ip}:{proxy.port}"
if key not in self.blacklist:
self.proxy_pool[key] = proxy
if key not in self.performance_stats:
self.performance_stats[key] = {
'success_count': 0,
'fail_count': 0,
'total_response_time': 0.0,
'last_success': None,
'score': 100.0 # 初始分数
}
def get_best_proxy(self, protocol: str = "http") -> Optional[ProxyIP]:
"""根据评分获取最佳代理"""
available = [
(key, proxy) for key, proxy in self.proxy_pool.items()
if proxy.protocol == protocol and key not in self.blacklist
]
if not available:
return None
# 综合评分算法:成功率权重40%,响应速度权重30%,新鲜度权重30%
scored_proxies = []
for key, proxy in available:
stats = self.performance_stats[key]
# 计算成功率(避免除零)
total_attempts = stats['success_count'] + stats['fail_count']
success_rate = (stats['success_count'] / total_attempts * 100) if total_attempts > 0 else 0
# 计算平均响应时间(秒),越短越好
avg_time = (stats['total_response_time'] / stats['success_count']) if stats['success_count'] > 0 else 10.0
time_score = max(0, 100 - (avg_time * 10)) # 假设1秒得90分,2秒得80分...
# 新鲜度:最近成功使用的时间越近,分数越高
freshness = 100
if stats['last_success']:
hours_ago = (time.time() - stats['last_success']) / 3600
freshness = max(0, 100 - hours_ago * 5) # 每过1小时减5分
# 综合评分
total_score = success_rate * 0.4 + time_score * 0.3 + freshness * 0.3
scored_proxies.append((total_score, proxy))
# 返回评分最高的代理
scored_proxies.sort(reverse=True)
return scored_proxies[0][1] if scored_proxies else None
def update_stats(self, proxy_key: str, success: bool, response_time: float = 0):
"""更新代理性能统计"""
if proxy_key not in self.performance_stats:
return
stats = self.performance_stats[proxy_key]
if success:
stats['success_count'] += 1
stats['total_response_time'] += response_time
stats['last_success'] = time.time()
# 成功一次,分数适当增加
stats['score'] = min(100, stats['score'] + 2)
else:
stats['fail_count'] += 1
# 失败一次,分数大幅降低
stats['score'] = max(0, stats['score'] - 10)
# 连续失败3次加入黑名单
if stats['fail_count'] >= 3 and stats['success_count'] == 0:
self.blacklist.add(proxy_key)
logging.warning(f"代理 {proxy_key} 因连续失败被加入黑名单")
这个调度器实现了简单的评分机制,你可以根据实际需求调整权重。比如,如果你的业务对速度特别敏感,可以调高响应速度的权重。
3. 多维度质量检测引擎
获取到代理IP只是第一步,更重要的是验证它的质量。一个“能用”的代理和“好用”的代理之间差距很大。
3.1 基础连通性测试
最基本的测试是检查代理能否正常访问互联网。但这里有个细节:我们应该测试多个目标网站,而不是只测一个。
import asyncio
import aiohttp
from typing import List, Tuple
import ssl
class ConnectivityTester:
"""连通性测试器"""
def __init__(self):
self.test_targets = [
("http://httpbin.org/ip", "http"), # 用于验证代理IP
("https://www.baidu.com", "https"), # 国内常用站
("https://www.google.com", "https"), # 国际站(需要代理能访问)
("http://www.qq.com", "http"), # 另一个国内站
]
self.timeout = aiohttp.ClientTimeout(total=10)
async def test_single_proxy(self, proxy: ProxyIP, target_url: str = None) -> Tuple[bool, float, dict]:
"""测试单个代理的连通性"""
if not target_url:
target_url = "http://httpbin.org/ip"
proxy_url = f"{proxy.protocol}://{proxy.ip}:{proxy.port}"
try:
connector = aiohttp.TCPConnector(ssl=False)
async with aiohttp.ClientSession(connector=connector, timeout=self.timeout) as session:
start_time = asyncio.get_event_loop().time()
async with session.get(
target_url,
proxy=proxy_url,
headers={'User-Agent': 'Mozilla/5.0'}
) as response:
end_time = asyncio.get_event_loop().time()
response_time = end_time - start_time
if response.status == 200:
# 验证返回的内容是否确实使用了代理
if "httpbin.org/ip" in target_url:
try:
data = await response.json()
origin_ip = data.get("origin", "")
if proxy.ip in origin_ip:
return True, response_time, {"status": "success", "origin": origin_ip}
except:
pass
return True, response_time, {"status": "success", "code": response.status}
else:
return False, response_time, {"status": "fail", "code": response.status}
except asyncio.TimeoutError:
return False, 10.0, {"status": "timeout"}
except aiohttp.ClientConnectorError:
return False, 0.0, {"status": "connection_error"}
except Exception as e:
return False, 0.0, {"status": "error", "message": str(e)}
async def comprehensive_test(self, proxy: ProxyIP) -> dict:
"""综合测试:访问多个目标网站"""
results = {}
total_success = 0
for url, protocol in self.test_targets:
# 如果代理协议不支持,跳过
if proxy.protocol != protocol and proxy.protocol != "all":
continue
success, resp_time, detail = await self.test_single_proxy(proxy, url)
results[url] = {
"success": success,
"response_time": resp_time,
"detail": detail
}
if success:
total_success += 1
# 计算综合评分
success_rate = total_success / len(results) if results else 0
avg_time = sum(r["response_time"] for r in results.values() if r["success"]) / total_success if total_success > 0 else 10.0
return {
"proxy": proxy,
"results": results,
"success_rate": success_rate,
"avg_response_time": avg_time,
"overall_score": self._calculate_score(success_rate, avg_time)
}
def _calculate_score(self, success_rate: float, avg_time: float) -> float:
"""计算综合评分"""
# 成功率权重60%,速度权重40%
rate_score = success_rate * 100
time_score = max(0, 100 - (avg_time * 20)) # 每0.05秒减1分
return rate_score * 0.6 + time_score * 0.4
3.2 匿名度与地理位置检测
对于很多业务场景,代理的匿名级别和地理位置同样重要。比如,有些网站会检测是否使用了代理,有些业务需要特定地区的IP。
class AdvancedProxyTester:
"""高级代理测试:匿名度、地理位置等"""
async def check_anonymity(self, proxy: ProxyIP) -> dict:
"""检测代理匿名级别"""
test_urls = {
"transparent": "http://httpbin.org/headers", # 透明代理检测
"anonymous": "http://httpbin.org/ip", # 匿名代理检测
"elite": "https://api.ipify.org?format=json" # 高匿代理检测
}
anonymity_level = "elite" # 默认假设为高匿
detected_headers = {}
for level, url in test_urls.items():
try:
async with aiohttp.ClientSession() as session:
proxy_url = f"{proxy.protocol}://{proxy.ip}:{proxy.port}"
async with session.get(url, proxy=proxy_url) as response:
if response.status == 200:
data = await response.json()
if level == "transparent":
# 检查是否透露了真实IP的头部
headers = data.get("headers", {})
if "X-Forwarded-For" in headers or "Via" in headers:
anonymity_level = "transparent"
detected_headers = headers
break
elif level == "anonymous":
# 匿名代理应该只显示代理IP,不透露真实IP
origin = data.get("origin", "")
if "," in origin: # 如果有多个IP,可能是透明代理
anonymity_level = "anonymous"
break
except:
continue
return {
"level": anonymity_level,
"description": self._get_anonymity_description(anonymity_level),
"headers": detected_headers
}
async def check_geo_location(self, proxy: ProxyIP) -> dict:
"""检测代理的地理位置"""
# 可以使用免费的IP地理位置API
geo_apis = [
"https://ipapi.co/{ip}/json/",
"http://ip-api.com/json/{ip}",
"https://freegeoip.app/json/{ip}"
]
for api_template in geo_apis:
try:
url = api_template.format(ip=proxy.ip)
async with aiohttp.ClientSession() as session:
async with session.get(url, timeout=5) as response:
if response.status == 200:
geo_data = await response.json()
return self._parse_geo_data(geo_data, api_template)
except:
continue
return {"error": "无法获取地理位置信息"}
def _get_anonymity_description(self, level: str) -> str:
descriptions = {
"transparent": "透明代理:目标服务器能看到你的真实IP",
"anonymous": "匿名代理:目标服务器知道你在用代理,但不知道真实IP",
"elite": "高匿代理:目标服务器认为你就是代理IP,完全隐藏"
}
return descriptions.get(level, "未知级别")
def _parse_geo_data(self, data: dict, api_source: str) -> dict:
"""解析不同API返回的地理位置数据"""
if "ipapi.co" in api_source:
return {
"country": data.get("country_name"),
"region": data.get("region"),
"city": data.get("city"),
"isp": data.get("org"),
"latitude": data.get("latitude"),
"longitude": data.get("longitude")
}
elif "ip-api.com" in api_source:
return {
"country": data.get("country"),
"region": data.get("regionName"),
"city": data.get("city"),
"isp": data.get("isp"),
"latitude": data.get("lat"),
"longitude": data.get("lon")
}
else:
return data
3.3 性能基准测试
除了连通性,我们还需要测试代理在不同网络条件下的性能表现。这里我设计了一个多线程的基准测试:
import concurrent.futures
import statistics
from datetime import datetime
class PerformanceBenchmark:
"""代理性能基准测试"""
def __init__(self, concurrent_workers: int = 10):
self.concurrent_workers = concurrent_workers
self.connectivity_tester = ConnectivityTester()
def run_benchmark(self, proxy: ProxyIP, test_count: int = 20) -> dict:
"""运行性能基准测试"""
print(f"开始性能基准测试: {proxy.ip}:{proxy.port}")
# 准备测试任务
tasks = [(proxy, f"http://httpbin.org/delay/{i%3}") for i in range(test_count)]
# 使用线程池并发测试
with concurrent.futures.ThreadPoolExecutor(max_workers=self.concurrent_workers) as executor:
future_to_task = {
executor.submit(self._single_test, task[0], task[1]): task
for task in tasks
}
results = []
for future in concurrent.futures.as_completed(future_to_task):
success, response_time, _ = future.result()
if success:
results.append(response_time)
if not results:
return {"error": "所有测试均失败"}
# 计算统计指标
stats = {
"test_count": test_count,
"success_count": len(results),
"success_rate": len(results) / test_count * 100,
"min_time": min(results),
"max_time": max(results),
"avg_time": statistics.mean(results),
"median_time": statistics.median(results),
"std_dev": statistics.stdev(results) if len(results) > 1 else 0,
"p95_time": sorted(results)[int(len(results) * 0.95)] if results else 0,
"requests_per_second": 1 / statistics.mean(results) if statistics.mean(results) > 0 else 0
}
# 性能评级
stats["performance_grade"] = self._grade_performance(stats)
return stats
def _single_test(self, proxy: ProxyIP, url: str) -> tuple:
"""单个测试任务"""
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
return loop.run_until_complete(
self.connectivity_tester.test_single_proxy(proxy, url)
)
finally:
loop.close()
def _grade_performance(self, stats: dict) -> str:
"""根据性能数据评级"""
avg_time = stats["avg_time"]
success_rate = stats["success_rate"]
if success_rate < 80:
return "F" # 失败率太高
elif avg_time < 0.5:
return "A+" # 极快
elif avg_time < 1.0:
return "A" # 优秀
elif avg_time < 2.0:
return "B" # 良好
elif avg_time < 3.0:
return "C" # 一般
elif avg_time < 5.0:
return "D" # 较差
else:
return "F" # 极差
4. 数据存储与可视化分析
测试产生的数据需要妥善存储和分析。我推荐使用SQLite作为本地数据库,轻量且无需额外服务。
4.1 数据库设计与实现
import sqlite3
from contextlib import contextmanager
from datetime import datetime
import json
class ProxyDatabase:
"""代理测试数据存储"""
def __init__(self, db_path: str = "proxy_monitor.db"):
self.db_path = db_path
self._init_database()
def _init_database(self):
"""初始化数据库表结构"""
with self._get_connection() as conn:
cursor = conn.cursor()
# 代理基本信息表
cursor.execute('''
CREATE TABLE IF NOT EXISTS proxies (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ip TEXT NOT NULL,
port INTEGER NOT NULL,
protocol TEXT,
source TEXT,
location TEXT,
first_seen TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
last_seen TIMESTAMP,
is_active BOOLEAN DEFAULT 1,
UNIQUE(ip, port)
)
''')
# 测试结果表
cursor.execute('''
CREATE TABLE IF NOT EXISTS test_results (
id INTEGER PRIMARY KEY AUTOINCREMENT,
proxy_id INTEGER,
test_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
success BOOLEAN,
response_time REAL,
target_url TEXT,
status_code INTEGER,
error_message TEXT,
test_type TEXT,
metadata TEXT, -- JSON格式的额外数据
FOREIGN KEY (proxy_id) REFERENCES proxies (id)
)
''')
# 性能统计表(每日汇总)
cursor.execute('''
CREATE TABLE IF NOT EXISTS daily_stats (
id INTEGER PRIMARY KEY AUTOINCREMENT,
proxy_id INTEGER,
date DATE,
total_tests INTEGER DEFAULT 0,
success_count INTEGER DEFAULT 0,
avg_response_time REAL,
min_response_time REAL,
max_response_time REAL,
success_rate REAL,
performance_grade TEXT,
FOREIGN KEY (proxy_id) REFERENCES proxies (id),
UNIQUE(proxy_id, date)
)
''')
# 服务商统计表
cursor.execute('''
CREATE TABLE IF NOT EXISTS provider_stats (
id INTEGER PRIMARY KEY AUTOINCREMENT,
provider_name TEXT,
date DATE,
total_proxies INTEGER,
active_proxies INTEGER,
avg_success_rate REAL,
avg_response_time REAL,
cost_per_proxy REAL,
value_score REAL,
UNIQUE(provider_name, date)
)
''')
conn.commit()
@contextmanager
def _get_connection(self):
"""获取数据库连接(上下文管理器)"""
conn = sqlite3.connect(self.db_path)
conn.row_factory = sqlite3.Row # 返回字典格式
try:
yield conn
conn.commit()
except Exception as e:
conn.rollback()
raise e
finally:
conn.close()
def save_test_result(self, proxy: ProxyIP, test_result: dict):
"""保存单次测试结果"""
with self._get_connection() as conn:
cursor = conn.cursor()
# 插入或更新代理信息
cursor.execute('''
INSERT OR IGNORE INTO proxies (ip, port, protocol, source, location)
VALUES (?, ?, ?, ?, ?)
''', (proxy.ip, proxy.port, proxy.protocol, proxy.source, proxy.location))
# 获取代理ID
cursor.execute('SELECT id FROM proxies WHERE ip = ? AND port = ?',
(proxy.ip, proxy.port))
proxy_id = cursor.fetchone()["id"]
# 更新最后出现时间
cursor.execute('''
UPDATE proxies SET last_seen = CURRENT_TIMESTAMP WHERE id = ?
''', (proxy_id,))
# 保存测试结果
metadata = json.dumps(test_result.get("metadata", {}))
cursor.execute('''
INSERT INTO test_results
(proxy_id, success, response_time, target_url, status_code, error_message, test_type, metadata)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
''', (
proxy_id,
test_result.get("success", False),
test_result.get("response_time", 0),
test_result.get("target_url", ""),
test_result.get("status_code", 0),
test_result.get("error_message", ""),
test_result.get("test_type", "connectivity"),
metadata
))
def get_proxy_stats(self, days: int = 7) -> list:
"""获取最近N天的代理统计"""
with self._get_connection() as conn:
cursor = conn.cursor()
cursor.execute('''
SELECT
p.ip,
p.port,
p.source,
COUNT(tr.id) as total_tests,
SUM(CASE WHEN tr.success = 1 THEN 1 ELSE 0 END) as success_count,
AVG(tr.response_time) as avg_response_time,
MIN(tr.response_time) as min_response_time,
MAX(tr.response_time) as max_response_time
FROM proxies p
LEFT JOIN test_results tr ON p.id = tr.proxy_id
WHERE tr.test_time >= datetime('now', ?)
GROUP BY p.id
HAVING total_tests > 0
ORDER BY avg_response_time ASC
''', (f'-{days} days',))
return [dict(row) for row in cursor.fetchall()]
def generate_daily_report(self):
"""生成每日统计报告"""
with self._get_connection() as conn:
cursor = conn.cursor()
# 获取所有活跃代理
cursor.execute('SELECT id, source FROM proxies WHERE is_active = 1')
active_proxies = cursor.fetchall()
for proxy in active_proxies:
proxy_id = proxy["id"]
source = proxy["source"]
# 计算今日统计
cursor.execute('''
SELECT
COUNT(*) as total_tests,
SUM(CASE WHEN success = 1 THEN 1 ELSE 0 END) as success_count,
AVG(response_time) as avg_response_time,
MIN(response_time) as min_response_time,
MAX(response_time) as max_response_time
FROM test_results
WHERE proxy_id = ?
AND date(test_time) = date('now')
''', (proxy_id,))
stats = cursor.fetchone()
if stats["total_tests"] > 0:
success_rate = stats["success_count"] / stats["total_tests"] * 100
# 性能评级
grade = "A" if stats["avg_response_time"] < 1.0 else \
"B" if stats["avg_response_time"] < 2.0 else \
"C" if stats["avg_response_time"] < 3.0 else \
"D" if stats["avg_response_time"] < 5.0 else "F"
# 插入每日统计
cursor.execute('''
INSERT OR REPLACE INTO daily_stats
(proxy_id, date, total_tests, success_count, avg_response_time,
min_response_time, max_response_time, success_rate, performance_grade)
VALUES (?, date('now'), ?, ?, ?, ?, ?, ?, ?)
''', (
proxy_id,
stats["total_tests"],
stats["success_count"],
stats["avg_response_time"],
stats["min_response_time"],
stats["max_response_time"],
success_rate,
grade
))
4.2 数据可视化与报表生成
有了数据,我们还需要直观的可视化。这里使用matplotlib和pandas生成图表:
import matplotlib.pyplot as plt
import pandas as pd
from datetime import datetime, timedelta
import seaborn as sns
class DataVisualizer:
"""数据可视化工具"""
def __init__(self, db: ProxyDatabase):
self.db = db
plt.style.use('seaborn-v0_8-darkgrid')
def generate_performance_report(self, days: int = 7, output_path: str = "report.png"):
"""生成性能报告图表"""
stats = self.db.get_proxy_stats(days)
if not stats:
print("没有足够的数据生成报告")
return
df = pd.DataFrame(stats)
# 创建图表
fig, axes = plt.subplots(2, 2, figsize=(15, 10))
# 1. 成功率分布
ax1 = axes[0, 0]
success_rates = (df['success_count'] / df['total_tests'] * 100).round(1)
ax1.hist(success_rates, bins=20, edgecolor='black', alpha=0.7)
ax1.set_xlabel('成功率 (%)')
ax1.set_ylabel('代理数量')
ax1.set_title('代理成功率分布')
ax1.axvline(x=success_rates.mean(), color='red', linestyle='--', label=f'平均: {success_rates.mean():.1f}%')
ax1.legend()
# 2. 响应时间分布
ax2 = axes[0, 1]
ax2.scatter(df['avg_response_time'], success_rates, alpha=0.6)
ax2.set_xlabel('平均响应时间 (秒)')
ax2.set_ylabel('成功率 (%)')
ax2.set_title('响应时间 vs 成功率')
# 添加最佳区域标注
ax2.axvspan(0, 1.0, alpha=0.2, color='green', label='优秀 (<1s)')
ax2.axvspan(1.0, 2.0, alpha=0.2, color='yellow', label='良好 (1-2s)')
ax2.axvspan(2.0, 5.0, alpha=0.2, color='orange', label='一般 (2-5s)')
ax2.legend()
# 3. 各服务商对比
ax3 = axes[1, 0]
provider_stats = df.groupby('source').agg({
'avg_response_time': 'mean',
'total_tests': 'sum',
'success_count': 'sum'
}).reset_index()
provider_stats['success_rate'] = provider_stats['success_count'] / provider_stats['total_tests'] * 100
x = range(len(provider_stats))
width = 0.35
bars1 = ax3.bar(x, provider_stats['avg_response_time'], width, label='平均响应时间')
ax3.set_xlabel('服务商')
ax3.set_ylabel('响应时间 (秒)', color='tab:blue')
ax3.set_xticks(x)
ax3.set_xticklabels(provider_stats['source'], rotation=45)
ax3.tick_params(axis='y', labelcolor='tab:blue')
ax3_twin = ax3.twinx()
bars2 = ax3_twin.bar([i + width for i in x], provider_stats['success_rate'], width,
color='tab:orange', label='成功率')
ax3_twin.set_ylabel('成功率 (%)', color='tab:orange')
ax3_twin.tick_params(axis='y', labelcolor='tab:orange')
ax3.set_title('各服务商性能对比')
# 4. 时间趋势(需要按日期分组的数据)
ax4 = axes[1, 1]
# 这里简化处理,实际应该从数据库获取时间序列数据
ax4.text(0.5, 0.5, '时间趋势图\n(需要更长时间数据)',
horizontalalignment='center', verticalalignment='center',
transform=ax4.transAxes, fontsize=12)
ax4.set_title('性能时间趋势')
plt.tight_layout()
plt.savefig(output_path, dpi=300, bbox_inches='tight')
plt.close()
print(f"报告已生成: {output_path}")
# 同时生成文本摘要
self._generate_text_summary(df, provider_stats)
def _generate_text_summary(self, df: pd.DataFrame, provider_stats: pd.DataFrame):
"""生成文本摘要"""
print("\n" + "="*60)
print("代理性能测试报告摘要")
print("="*60)
print(f"\n📊 总体统计(基于最近{len(df)}个代理):")
print(f" 平均成功率: {(df['success_count'].sum() / df['total_tests'].sum() * 100):.1f}%")
print(f" 平均响应时间: {df['avg_response_time'].mean():.2f}秒")
print(f" 最快代理: {df['avg_response_time'].min():.2f}秒")
print(f" 最慢代理: {df['avg_response_time'].max():.2f}秒")
print(f"\n🏆 服务商排名(按平均响应时间):")
provider_stats = provider_stats.sort_values('avg_response_time')
for i, (_, row) in enumerate(provider_stats.iterrows(), 1):
print(f" {i}. {row['source']:20} {row['avg_response_time']:.2f}秒 {row['success_rate']:.1f}%")
print(f"\n⚠️ 需要关注的代理:")
problematic = df[(df['success_count'] / df['total_tests'] < 0.5) | (df['avg_response_time'] > 3.0)]
if len(problematic) > 0:
for _, row in problematic.iterrows():
rate = row['success_count'] / row['total_tests'] * 100
print(f" {row['ip']}:{row['port']} - 成功率{rate:.1f}%, 响应时间{row['avg_response_time']:.2f}秒")
else:
print(" 所有代理表现良好!")
print("\n" + "="*60)
5. 自动化监控与告警系统
最后,我们需要一个自动化的监控系统,定期测试代理并发送告警。
5.1 定时任务调度
import schedule
import time
from threading import Thread
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import json
class ProxyMonitor:
"""代理监控主程序"""
def __init__(self, config_path: str = "config.yaml"):
self.config = self._load_config(config_path)
self.scheduler = ProxyScheduler()
self.database = ProxyDatabase()
self.visualizer = DataVisualizer(self.database)
self.is_running = False
# 初始化各个测试器
self.connectivity_tester = ConnectivityTester()
self.advanced_tester = AdvancedProxyTester()
self.benchmark = PerformanceBenchmark()
# 加载代理源
self.proxy_fetchers = self._init_fetchers()
def _load_config(self, config_path: str) -> dict:
"""加载配置文件"""
import yaml
try:
with open(config_path, 'r', encoding='utf-8') as f:
return yaml.safe_load(f)
except FileNotFoundError:
# 默认配置
return {
'monitor': {
'test_interval_minutes': 30,
'benchmark_interval_hours': 6,
'report_interval_hours': 24
},
'alert': {
'enabled': True,
'email': {
'smtp_server': 'smtp.example.com',
'smtp_port': 587,
'username': 'your_email@example.com',
'password': 'your_password',
'recipients': ['admin@example.com']
},
'thresholds': {
'success_rate': 80, # 成功率低于80%告警
'avg_response_time': 3.0, # 平均响应时间超过3秒告警
'proxy_count': 10 # 可用代理少于10个告警
}
},
'proxy_sources': [
{
'name': 'provider1',
'type': 'common',
'api_url': 'https://api.provider1.com/get',
'api_key': 'your_key_here',
'format': 'json'
}
]
}
def _init_fetchers(self) -> list:
"""初始化代理获取器"""
fetchers = []
for source in self.config.get('proxy_sources', []):
if source['type'] == 'common':
fetcher = CommonProxyFetcher(
api_url=source['api_url'],
api_key=source['api_key'],
format_type=source.get('format', 'json')
)
fetchers.append(fetcher)
return fetchers
async def fetch_and_test_proxies(self):
"""获取并测试代理"""
print(f"[{datetime.now()}] 开始新一轮代理测试")
all_proxies = []
# 从所有源获取代理
for fetcher in self.proxy_fetchers:
try:
proxies = await fetcher.fetch_batch(count=20)
all_proxies.extend(proxies)
print(f"从 {fetcher.__class__.__name__} 获取到 {len(proxies)} 个代理")
except Exception as e:
print(f"从 {fetcher.__class__.__name__} 获取代理失败: {e}")
# 添加到调度器
for proxy in all_proxies:
self.scheduler.add_proxy(proxy)
# 测试所有新代理
test_tasks = []
for proxy in all_proxies:
task = self._test_proxy_async(proxy)
test_tasks.append(task)
# 等待所有测试完成
if test_tasks:
results = await asyncio.gather(*test_tasks, return_exceptions=True)
# 检查是否有严重问题
success_count = sum(1 for r in results if isinstance(r, dict) and r.get('success_rate', 0) > 50)
if success_count / len(results) < 0.3:
self._send_alert("代理质量严重下降",
f"新获取的代理中只有{success_count}/{len(results)}个可用")
async def _test_proxy_async(self, proxy: ProxyIP) -> dict:
"""异步测试单个代理"""
try:
# 基础连通性测试
connectivity_result = await self.connectivity_tester.comprehensive_test(proxy)
# 保存结果
self.database.save_test_result(proxy, {
'success': connectivity_result['success_rate'] > 50,
'response_time': connectivity_result['avg_response_time'],
'target_url': 'multiple',
'test_type': 'connectivity',
'metadata': connectivity_result
})
# 更新调度器统计
proxy_key = f"{proxy.ip}:{proxy.port}"
self.scheduler.update_stats(
proxy_key,
connectivity_result['success_rate'] > 50,
connectivity_result['avg_response_time']
)
return connectivity_result
except Exception as e:
print(f"测试代理 {proxy.ip}:{proxy.port} 时出错: {e}")
return {'error': str(e)}
def run_benchmark(self):
"""运行性能基准测试"""
print(f"[{datetime.now()}] 开始性能基准测试")
# 选择最近表现良好的代理进行深度测试
good_proxies = []
for key, proxy in self.scheduler.proxy_pool.items():
stats = self.scheduler.performance_stats.get(key, {})
if stats.get('score', 0) > 70: # 只测试评分70以上的代理
good_proxies.append(proxy)
if not good_proxies:
print("没有足够的高质量代理进行基准测试")
return
# 随机选择3个代理进行深度测试
import random
selected = random.sample(good_proxies, min(3, len(good_proxies)))
for proxy in selected:
benchmark_result = self.benchmark.run_benchmark(proxy)
print(f"代理 {proxy.ip}:{proxy.port} 基准测试结果: {benchmark_result}")
def generate_daily_report(self):
"""生成每日报告"""
print(f"[{datetime.now()}] 生成每日报告")
# 更新数据库中的每日统计
self.database.generate_daily_report()
# 生成可视化报告
report_file = f"proxy_report_{datetime.now().strftime('%Y%m%d')}.png"
self.visualizer.generate_performance_report(days=1, output_path=report_file)
# 发送邮件报告(如果配置了)
if self.config['alert']['enabled']:
self._send_daily_report(report_file)
def _send_alert(self, subject: str, message: str):
"""发送告警邮件"""
if not self.config['alert']['enabled']:
return
email_config = self.config['alert']['email']
try:
msg = MIMEMultipart()
msg['From'] = email_config['username']
msg['To'] = ', '.join(email_config['recipients'])
msg['Subject'] = f"[代理监控告警] {subject}"
body = f"""
告警时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}
告警内容:
{message}
当前系统状态:
- 代理池总数: {len(self.scheduler.proxy_pool)}
- 黑名单数量: {len(self.scheduler.blacklist)}
- 最近测试时间: {datetime.now().strftime('%H:%M:%S')}
请及时检查代理服务状态。
"""
msg.attach(MIMEText(body, 'plain'))
with smtplib.SMTP(email_config['smtp_server'], email_config['smtp_port']) as server:
server.starttls()
server.login(email_config['username'], email_config['password'])
server.send_message(msg)
print(f"告警邮件已发送: {subject}")
except Exception as e:
print(f"发送告警邮件失败: {e}")
def _send_daily_report(self, report_file: str):
"""发送每日报告邮件"""
email_config = self.config['alert']['email']
try:
# 这里简化处理,实际应该将图片作为附件发送
stats = self.database.get_proxy_stats(days=1)
total_tests = sum(s['total_tests'] for s in stats)
success_rate = sum(s['success_count'] for s in stats) / total_tests * 100 if total_tests > 0 else 0
msg = MIMEMultipart()
msg['From'] = email_config['username']
msg['To'] = ', '.join(email_config['recipients'])
msg['Subject'] = f"代理监控日报 {datetime.now().strftime('%Y-%m-%d')}"
body = f"""
代理监控系统日报
生成时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}
今日统计:
- 测试总次数: {total_tests}
- 平均成功率: {success_rate:.1f}%
- 活跃代理数: {len([s for s in stats if s['total_tests'] > 0])}
- 最佳响应时间: {min(s['avg_response_time'] for s in stats if s['total_tests'] > 0):.2f}秒
详细报告请查看附件图片。
系统运行状态正常。
"""
msg.attach(MIMEText(body, 'plain'))
# 添加图片附件(这里省略具体实现)
with smtplib.SMTP(email_config['smtp_server'], email_config['smtp_port']) as server:
server.starttls()
server.login(email_config['username'], email_config['password'])
server.send_message(msg)
print("每日报告邮件已发送")
except Exception as e:
print(f"发送日报邮件失败: {e}")
def start(self):
"""启动监控系统"""
self.is_running = True
print("代理监控系统启动...")
# 设置定时任务
test_interval = self.config['monitor']['test_interval_minutes']
benchmark_interval = self.config['monitor']['benchmark_interval_hours']
report_interval = self.config['monitor']['report_interval_hours']
schedule.every(test_interval).minutes.do(
lambda: asyncio.run(self.fetch_and_test_proxies())
)
schedule.every(benchmark_interval).hours.do(self.run_benchmark)
schedule.every(report_interval).hours.do(self.generate_daily_report)
# 立即运行一次
asyncio.run(self.fetch_and_test_proxies())
# 主循环
while self.is_running:
schedule.run_pending()
time.sleep(1)
def stop(self):
"""停止监控系统"""
self.is_running = False
print("代理监控系统停止")
# 主程序入口
if __name__ == "__main__":
monitor = ProxyMonitor()
# 可以在单独的线程中运行,避免阻塞
import threading
monitor_thread = threading.Thread(target=monitor.start)
monitor_thread.daemon = True
monitor_thread.start()
try:
# 主线程可以继续做其他事情
while True:
time.sleep(1)
except KeyboardInterrupt:
print("\n接收到中断信号,停止监控...")
monitor.stop()
5.2 配置文件示例
创建一个config.yaml文件来配置监控系统:
monitor:
test_interval_minutes: 30 # 每30分钟测试一次
benchmark_interval_hours: 6 # 每6小时运行一次基准测试
report_interval_hours: 24 # 每24小时生成一次报告
alert:
enabled: true
email:
smtp_server: "smtp.gmail.com"
smtp_port: 587
username: "your_email@gmail.com"
password: "your_app_password" # 建议使用应用专用密码
recipients:
- "team@yourcompany.com"
- "admin@yourcompany.com"
thresholds:
success_rate: 80 # 成功率低于80%触发告警
avg_response_time: 3.0 # 平均响应时间超过3秒触发告警
proxy_count: 10 # 可用代理少于10个触发告警
proxy_sources:
- name: "provider_a"
type: "common"
api_url: "https://api.provider-a.com/get"
api_key: "your_api_key_here"
format: "json"
enabled: true
- name: "provider_b"
type: "common"
api_url: "http://api.provider-b.com/v1/proxy"
api_key: "your_api_key_here"
format: "text"
enabled: true
logging:
level: "INFO"
file: "proxy_monitor.log"
max_size_mb: 10
backup_count: 5
database:
path: "proxy_monitor.db"
backup_days: 30 # 保留30天数据
6. 部署与优化建议
6.1 生产环境部署
在实际生产环境中部署这套系统时,有几个关键点需要注意:
- 容器化部署:使用Docker可以简化依赖管理和部署流程
FROM python:3.9-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
# 创建数据卷用于持久化存储
VOLUME /app/data
CMD ["python", "monitor_main.py"]
- 多节点分布式测试:如果代理数量很大,可以考虑分布式部署。使用Redis作为中央调度器:
import redis
import pickle
class DistributedMonitor:
def __init__(self, redis_url="redis://localhost:6379"):
self.redis = redis.from_url(redis_url)
self.task_queue = "proxy_test_tasks"
self.result_queue = "proxy_test_results"
def distribute_tasks(self, proxies: list):
"""将测试任务分发到多个worker"""
for proxy in proxies:
task = {
'proxy': pickle.dumps(proxy),
'timestamp': time.time(),
'test_type': 'connectivity'
}
self.redis.rpush(self.task_queue, pickle.dumps(task))
def get_results(self):
"""收集测试结果"""
while True:
result_data = self.redis.lpop(self.result_queue)
if result_data:
yield pickle.loads(result_data)
else:
time.sleep(0.1)
- 监控系统自身:任何监控系统都需要被监控。建议添加:
- 心跳检测:确保监控进程正常运行
- 资源监控:CPU、内存、磁盘使用情况
- 错误日志集中收集:使用Sentry或ELK Stack
6.2 性能优化技巧
根据我的实际使用经验,下面这些优化措施能显著提升系统效率:
连接池管理:
import aiohttp
from aiohttp import ClientSession, TCPConnector
class OptimizedTester:
def __init__(self):
# 重用连接池,避免频繁创建连接
self.connector = TCPConnector(limit=100, limit_per_host=10, ttl_dns_cache=300)
self.session = None
async def get_session(self):
if not self.session:
self.session = ClientSession(connector=self.connector)
return self.session
async def close(self):
if self.session:
await self.session.close()
测试目标优化:
- 选择距离近、响应稳定的测试目标
- 避免使用可能屏蔽代理的网站(如某些电商网站)
- 考虑使用自建的测试端点,完全控制测试环境
智能测试频率:
- 对新获取的代理进行密集测试(前10分钟每分钟测一次)
- 对稳定代理降低测试频率(每小时测一次)
- 对问题代理增加测试频率,确认是否真的失效
6.3 常见问题排查
在实际运行中,你可能会遇到这些问题:
问题1:代理测试成功率突然下降
- 检查测试目标网站是否可访问
- 验证网络连接是否正常
- 查看代理服务商是否有公告或故障
- 检查是否触发了目标网站的反爬机制
问题2:系统内存持续增长
- 可能是数据库连接未正确关闭
- 检查是否有循环引用或内存泄漏
- 考虑定期重启worker进程
问题3:测试速度过慢
- 调整并发数(
concurrent_workers参数) - 减少测试目标数量
- 优化网络连接(使用更快的DNS)
这套系统在我负责的几个数据采集项目中已经稳定运行了半年多,最初版本确实遇到了不少问题,比如内存泄漏、数据库锁竞争、网络超时处理不当等。经过几次迭代,现在的版本已经能够7x24小时稳定运行,每天处理数十万次代理测试,成功帮我们识别出多个不稳定的代理服务商,也让我们在采购决策时有充分的数据支持。
最让我满意的是它的灵活性——当我们需要测试新的代理服务商时,只需要在配置文件中添加几行配置;当业务对代理有特殊要求时(比如需要特定地区的IP),可以快速调整测试逻辑。这种可扩展性让系统能够随着业务需求一起成长,而不是每次变化都要重写代码。
如果你打算在生产环境使用,我建议先从简单的版本开始,只实现最核心的连通性测试和基础监控。运行一段时间收集数据,了解你们的实际使用模式,然后再逐步添加高级功能。毕竟,最适合的监控系统是那个真正被用起来的系统,而不是功能最全的系统。
更多推荐
所有评论(0)