设计目标

自建代理池解决商业代理 API 调用限频问题:
- 聚合多个代理来源(文件/API/数据库)
- 健康检测淘汰失效代理
- 轮转分发,支持并发请求
- 对外提供 HTTP API,业务方直接拉取可用代理

架构

proxy_pool/
├── app.py          # Flask HTTP API 服务
├── pool.py         # 代理池核心(deque + 健康检测)
├── sources.py      # 代理来源(文件/第三方API)
├── validator.py    # IP 合法性验证(过滤内网地址)
└── config.py       # 配置

代理池核心

from collections import deque
import threading, time, requests, logging

logger = logging.getLogger(__name__)

class ProxyPool:
    def __init__(self, sources: list, check_interval: int = 300):
        self._pool: deque = deque()
        self._lock = threading.Lock()
        self._sources = sources
        self._check_interval = check_interval
        self._backoff_base = 1.5
        self._backoff_max = 300  # 最长等待5分钟

    def start(self):
        # 启动后台刷新线程
        t = threading.Thread(target=self._refresh_loop, daemon=True)
        t.start()

    def _refresh_loop(self):
        fail_count = 0
        while True:
            try:
                self._refresh()
                fail_count = 0
                time.sleep(self._check_interval)
            except Exception as e:
                fail_count += 1
                wait = min(self._backoff_base ** fail_count, self._backoff_max)
                logger.error(f"刷新失败 #{fail_count}{wait:.0f}s 后重试: {e}")
                time.sleep(wait)

    def _refresh(self):
        new_proxies = []
        for source in self._sources:
            try:
                proxies = source.fetch()
                valid = [p for p in proxies if validate_ip(p)]
                new_proxies.extend(valid)
            except Exception as e:
                logger.warning(f"来源 {source} 失败: {e}")

        # 健康检测
        alive = []
        for proxy in new_proxies:
            if self._check_proxy(proxy):
                alive.append(proxy)

        with self._lock:
            self._pool.clear()
            self._pool.extend(alive)
        logger.info(f"代理池刷新完成,有效代理: {len(alive)}/{len(new_proxies)}")

    def _check_proxy(self, proxy: str, timeout: int = 8) -> bool:
        try:
            r = requests.get(
                "https://httpbin.org/ip",
                proxies={"http": f"http://{proxy}", "https": f"http://{proxy}"},
                timeout=timeout
            )
            return r.status_code == 200
        except Exception:
            return False

    def get(self) -> str | None:
        # 轮转获取一个代理
        with self._lock:
            if not self._pool:
                return None
            proxy = self._pool[0]
            self._pool.rotate(-1)  # 将头部移到尾部,实现轮转
            return proxy

    def remove(self, proxy: str):
        # 移除失效代理
        with self._lock:
            try:
                self._pool.remove(proxy)
            except ValueError:
                pass

    @property
    def size(self) -> int:
        return len(self._pool)

IP 内网过滤

import ipaddress, re

PRIVATE_NETS = [
    ipaddress.ip_network("10.0.0.0/8"),
    ipaddress.ip_network("172.16.0.0/12"),
    ipaddress.ip_network("192.168.0.0/16"),
    ipaddress.ip_network("127.0.0.0/8"),
    ipaddress.ip_network("::1/128"),
]
PROXY_RE = re.compile(r"^(\d{1,3}(?:\.\d{1,3}){3}):(\d{2,5})$")

def validate_ip(proxy: str) -> bool:
    m = PROXY_RE.match(proxy.strip())
    if not m:
        return False
    ip_str, port_str = m.group(1), m.group(2)
    port = int(port_str)
    if not (1 <= port <= 65535):
        return False
    try:
        ip = ipaddress.ip_address(ip_str)
        for net in PRIVATE_NETS:
            if ip in net:
                return False
        return True
    except ValueError:
        return False

Flask API 接口

from flask import Flask, jsonify, request, abort
from functools import wraps
import time

app = Flask(__name__)
pool = ProxyPool(sources=[FileSource("proxies.txt")])

# 简单速率限制(生产环境用 Flask-Limiter)
_rate_limit = {}
RATE_LIMIT = 60  # 每分钟最多60次

def rate_limit(f):
    @wraps(f)
    def wrapper(*args, **kwargs):
        ip = request.remote_addr
        now = time.time()
        calls = [t for t in _rate_limit.get(ip, []) if now - t < 60]
        if len(calls) >= RATE_LIMIT:
            abort(429)
        calls.append(now)
        _rate_limit[ip] = calls
        return f(*args, **kwargs)
    return wrapper

@app.route("/proxy/get")
@rate_limit
def get_proxy():
    proxy = pool.get()
    if not proxy:
        return jsonify({"error": "代理池为空"}), 503
    return jsonify({"proxy": proxy, "size": pool.size})

@app.route("/proxy/remove", methods=["POST"])
@rate_limit
def remove_proxy():
    proxy = request.json.get("proxy")
    if proxy:
        pool.remove(proxy)
    return jsonify({"ok": True})

@app.route("/proxy/stats")
def stats():
    return jsonify({"size": pool.size})

if __name__ == "__main__":
    pool.start()
    app.run(host="0.0.0.0", port=5000)

使用方式

# 业务方调用示例
import requests

def get_proxy_from_pool() -> str:
    r = requests.get("http://localhost:5000/proxy/get")
    return r.json()["proxy"]

def report_bad_proxy(proxy: str):
    requests.post("http://localhost:5000/proxy/remove", json={"proxy": proxy})

# 带失效重试的请求
def fetch_with_proxy(url: str, max_retries: int = 3):
    for _ in range(max_retries):
        proxy = get_proxy_from_pool()
        try:
            r = requests.get(url,
                proxies={"http": f"http://{proxy}", "https": f"http://{proxy}"},
                timeout=10)
            return r
        except Exception:
            report_bad_proxy(proxy)
    raise RuntimeError("所有代理均失效")