构建企业级代理IP自动化测试与监控平台:Python实战指南

最近和几个做数据采集和风控的朋友聊天,大家普遍反映一个痛点:市面上各种代理IP服务商都说自己“稳定高速”,但实际用起来效果参差不齐。团队花了不少钱采购,结果关键时刻掉链子,不是IP大量失效就是响应慢得让人抓狂。更头疼的是,缺乏一套系统化的方法来持续评估这些代理的质量,每次都是出了问题才手忙脚乱地临时测试。

这让我意识到,对于技术团队来说,单纯依赖服务商的口头承诺或一次性评测是远远不够的。我们需要的是一个能够持续运行、自动评估、智能预警的代理IP质量监控体系。这不仅能帮我们筛选出真正靠谱的服务商,还能在日常使用中及时发现性能波动,避免业务受到影响。

今天,我就把自己在多个项目中搭建这类系统的经验整理出来,分享一套完整的Python自动化测试方案。这套方案不绑定任何特定服务商,你可以用它来测试任何提供API接口的代理IP服务。我们会从最基础的API调用开始,一步步构建出包含有效性验证、性能统计、稳定性监控甚至可视化报表的完整工具链。所有代码都是模块化设计,你可以直接集成到现有系统中,或者根据团队需求进行定制。

1. 系统架构设计与核心模块规划

在开始写代码之前,我们先要理清楚整个系统需要哪些功能模块。一个完整的代理IP自动化测试平台,远不止“获取IP-测试连通性”这么简单。我们需要考虑的是全生命周期管理。

1.1 核心功能模块分解

我把整个系统划分为五个核心模块,每个模块都有明确的职责:

  1. IP获取与调度模块:负责从不同服务商的API获取IP,并实现智能调度策略
  2. 质量检测引擎:执行多种维度的测试,包括连通性、匿名度、地理位置等
  3. 性能数据收集器:记录每次测试的详细指标,建立历史数据库
  4. 数据分析与评分系统:基于历史数据计算每个IP的综合评分
  5. 监控与告警模块:实时监控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 生产环境部署

在实际生产环境中部署这套系统时,有几个关键点需要注意:

  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"]
  1. 多节点分布式测试:如果代理数量很大,可以考虑分布式部署。使用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)
  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),可以快速调整测试逻辑。这种可扩展性让系统能够随着业务需求一起成长,而不是每次变化都要重写代码。

如果你打算在生产环境使用,我建议先从简单的版本开始,只实现最核心的连通性测试和基础监控。运行一段时间收集数据,了解你们的实际使用模式,然后再逐步添加高级功能。毕竟,最适合的监控系统是那个真正被用起来的系统,而不是功能最全的系统。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐