ProxyScrape Python 教程:使用 httpx 和 SOCKS5 构建异步轮换代理池
一个面向 ProxyScrape 的高级 Python 模式:将你的 HTTP 和 SOCKS5 端点建模为代理池,在并发限制下轮换使用,失败时退避,并跟踪每个端点的健康状态。
概述
这是一篇高级 Python 教程。你不会通过单个代理发送一个请求,而是构建一个小型抓取器,在套餐中的 ProxyScrape 端点之间轮换,限制并发,使用带抖动的退避重试,并将失败的端点暂时置于冷却状态。
你将:
- 安装支持 SOCKS5 的
httpx。 - 将 ProxyScrape 凭据保存在环境变量中。
- 将你的 HTTP 和 SOCKS5 端点建模为代理池。
- 轮换、限流、重试,并报告每个端点的健康状态。
- 在扩大规模前进行冒烟测试。
ProxyScrape 提供住宅、数据中心和移动代理,支持 HTTP 和 SOCKS5,并按量付费或每月固定费用计费。你选择哪种代理会影响成本和延迟,而不是客户端代码——下面的代码将端点视为配置。
在编写代码前选择代理类型
| 工作负载 | ProxyScrape 类型 | 为什么适合 |
|---|---|---|
| 高流量、低阻力目标(文档、站点地图、公共 API) | 数据中心 | 最快的选项;长期稳定运行适合定时任务 |
| 受机器人过滤的页面或特定地区内容 | 住宅 | 拥有真实用户来源的大型池(9700 万+ IP) |
| 登录流程和激进的反机器人目标 | 移动 | 信任度最高;使用低并发 |
前置条件
- Python 3.10 或更高版本。
- 一个带有凭据的 ProxyScrape 套餐,以及
http和/或socks5的主机和端口。 - 你的机器能够对 ProxyScrape 主机和端口进行出站访问。
步骤
第 1 步 — 创建项目并安装依赖
mkdir proxyscrape-async
cd proxyscrape-async
python -m venv .venv
source .venv/bin/activate
pip install "httpx[socks]"
[socks] 额外项会安装 socksio。没有它,第一次配置 socks5:// 端点时 httpx 会报错。httpx 0.26 及更高版本使用 proxy= 参数(更早的版本使用 proxies=),因此如果示例不匹配,请用 pip show httpx 检查你的版本。
第 2 步 — 将凭据和端点保存在环境变量中
# pool.py
import os
from dataclasses import dataclass
from urllib.parse import quote
@dataclass(frozen=True)
class ProxyEndpoint:
label: str
url: str
def endpoint_url(scheme: str, host: str, port: str, username: str, password: str) -> str:
user = quote(username, safe="")
secret = quote(password, safe="")
return f"{scheme}://{user}:{secret}@{host}:{port}"
HOST = os.environ["PROXYSCRAPE_HOST"]
USERNAME = os.environ["PROXYSCRAPE_USERNAME"]
PASSWORD = os.environ["PROXYSCRAPE_PASSWORD"]
POOL: list[ProxyEndpoint] = []
if http_port := os.environ.get("PROXYSCRAPE_HTTP_PORT"):
POOL.append(ProxyEndpoint("http-1", endpoint_url("http", HOST, http_port, USERNAME, PASSWORD)))
if socks_port := os.environ.get("PROXYSCRAPE_SOCKS5_PORT"):
POOL.append(ProxyEndpoint("socks5-1", endpoint_url("socks5", HOST, socks_port, USERNAME, PASSWORD)))
if not POOL:
raise SystemExit("Set PROXYSCRAPE_HTTP_PORT and/or PROXYSCRAPE_SOCKS5_PORT")
在 shell 中导出这些值,或在 CI 中从密钥管理器注入:
export PROXYSCRAPE_HOST="host-from-your-plan"
export PROXYSCRAPE_USERNAME="your-username"
export PROXYSCRAPE_PASSWORD="your-password"
export PROXYSCRAPE_HTTP_PORT="port-from-your-plan"
export PROXYSCRAPE_SOCKS5_PORT="port-from-your-plan"
通过 quote() 进行百分号编码很重要:包含 @、: 或 / 的密码否则会破坏 URL 解析。切勿提交这些值。
第 3 步 — 在扩大规模前验证每个端点
按顺序运行此检查清单。如果某一步失败,请在接触异步代码之前修复它。
- 检查 HTTP/HTTPS 端点。
- 检查 SOCKS5 端点。
- 确认返回的 IP 不是你自己的 IP。
curl -sS -x "http://$PROXYSCRAPE_USERNAME:$PROXYSCRAPE_PASSWORD@$PROXYSCRAPE_HOST:$PROXYSCRAPE_HTTP_PORT" "https://api.ipify.org?format=json"
curl -sS -x "socks5h://$PROXYSCRAPE_USERNAME:$PROXYSCRAPE_PASSWORD@$PROXYSCRAPE_HOST:$PROXYSCRAPE_SOCKS5_PORT" "https://api.ipify.org?format=json"
socks5h:// 告诉 curl 让代理解析 DNS。如果你的 Python 客户端同时提供 socks5:// 和 socks5h://,当本地 DNS 解析失败,或你希望 DNS 查询与请求一起离开你的机器时,优先使用带 h 的变体。
第 4 步 — 构建轮换的并发抓取器
# fetch_async.py
from __future__ import annotations
import asyncio
import itertools
import random
import time
import httpx
from pool import POOL, ProxyEndpoint
class RotatingFetcher:
"""Round-robin fetcher with concurrency limits, backoff, and cooldowns."""
def __init__(
self,
pool: list[ProxyEndpoint],
concurrency: int = 8,
max_attempts: int = 4,
timeout: float = 20.0,
cooldown_seconds: float = 45.0,
) -> None:
if not pool:
raise ValueError("Proxy pool is empty")
self._pool = pool
self._cycle = itertools.cycle(pool)
self._semaphore = asyncio.Semaphore(concurrency)
self._max_attempts = max_attempts
self._cooldown_seconds = cooldown_seconds
self._failures: dict[str, int] = {}
self._cooldown_until: dict[str, float] = {}
self._clients = {
item.label: httpx.AsyncClient(
proxy=item.url,
timeout=timeout,
follow_redirects=True,
headers={"User-Agent": "Mozilla/5.0 (compatible; ProxyScoutTutorial/1.0)"},
)
for item in pool
}
def _pick(self) -> ProxyEndpoint:
now = time.monotonic()
for _ in range(len(self._pool) * 2):
candidate = next(self._cycle)
if self._cooldown_until.get(candidate.label, 0.0) <= now:
return candidate
return min(self._pool, key=lambda item: self._cooldown_until.get(item.label, 0.0))
def _record_failure(self, label: str) -> None:
self._failures[label] = self._failures.get(label, 0) + 1
penalty = self._cooldown_seconds * self._failures[label]
self._cooldown_until[label] = time.monotonic() + penalty
def _record_success(self, label: str) -> None:
self._failures[label] = 0
self._cooldown_until[label] = 0.0
@staticmethod
def _backoff(attempt: int) -> float:
return min((2 ** attempt) + random.uniform(0, 0.5), 30.0)
async def fetch(self, url: str) -> httpx.Response:
last_error: Exception | None = None
for attempt in range(1, self._max_attempts + 1):
item = self._pick()
try:
async with self._semaphore:
response = await self._clients[item.label].get(url)
if response.status_code in {407, 429} or response.status_code >= 500:
raise httpx.HTTPStatusError(
f"retryable status {response.status_code}",
request=response.request,
response=response,
)
self._record_success(item.label)
return response
except Exception as exc:
last_error = exc
self._record_failure(item.label)
if attempt < self._max_attempts:
await asyncio.sleep(self._backoff(attempt))
raise RuntimeError(f"All {self._max_attempts} attempts failed for {url}") from last_error
def report(self) -> str:
active = {label: count for label, count in self._failures.items() if count}
if not active:
return "no endpoints have failed yet"
return ", ".join(f"{label}={count}" for label, count in sorted(active.items()))
async def aclose(self) -> None:
await asyncio.gather(*(client.aclose() for client in self._clients.values()))
值得保留的设计说明:
- 每个端点一个
AsyncClient。在 httpx 中,代理绑定到客户端,而不是单个请求,因此每个端点一个客户端是最简单的正确轮换模型。 - 无论你提交多少 URL,信号量都会限制进行中的请求数。
407和429与5xx一起被视为可重试状态,并且它们会计为端点失败,因此坏端点会在一段时间内停止接收流量。- 失败惩罚线性增长,并在成功时重置,这可以防止不稳定端点被紧密循环重试。
第 5 步 — 冒烟测试并调整并发
# run_smoke_test.py
import asyncio
import httpx
from fetch_async import RotatingFetcher
from pool import POOL
async def main() -> None:
fetcher = RotatingFetcher(pool=POOL, concurrency=8)
urls = [f"https://example.com/?page={index}" for index in range(1, 21)]
try:
results = await asyncio.gather(*(fetcher.fetch(url) for url in urls), return_exceptions=True)
ok = sum(1 for item in results if isinstance(item, httpx.Response))
print(f"success={ok} failed={len(results) - ok}")
print(f"endpoint health: {fetcher.report()}")
finally:
await fetcher.aclose()
if __name__ == "__main__":
asyncio.run(main())
然后调整:
- 运行一次。如果
failed为 0,将concurrency提高 4 并再次运行。 - 当延迟或超时开始上升时停止增加——这就是你所购买套餐的实际上限。
- 如果失败集中在单个端点上,将该端点视为不健康,并联系支持人员,而不是更用力地重试。
套餐说明:按量付费计费时,你可以为突发流量提高并发;每月固定费用套餐通常更适合稳定、定时的任务。
故障排除
| 症状 | 可能原因 | 修复 |
|---|---|---|
提到 socksio 的导入错误 |
未安装 httpx[socks] |
pip install "httpx[socks]" |
| 每个请求都失败并返回 407 | 凭据错误,或包含特殊字符的密码未编码 | 重新检查环境变量,并保留 quote() 调用 |
| 随着并发上升,超时增加 | 套餐或目标能承受的进行中请求过多 | 降低 concurrency,提高 timeout |
| 请求不断进入冷却 | 重复失败后惩罚叠加 | 先用 curl 验证端点,并考虑限制 cooldown_seconds |
socks5 端点失败而 http 正常 |
本地 DNS 解析问题 | 如果你的库支持,尝试 socks5h 方案 |
| 目标持续返回 403 | 该端点 IP 被此目标过滤 | 在更多端点之间轮换,或改用住宅或移动代理 |
总结
- ProxyScrape 支持 HTTP 和 SOCKS5;将套餐中的每个端点表示为
ProxyEndpoint,并在它们之间轮换。 - 将代理绑定到
httpx.AsyncClient,用信号量限制进行中的工作,并使用带抖动的退避重试。 - 将 407、429 和 5xx 响应视为端点失败,而不是用户错误,并对端点进行冷却。
- 在扩大规模前用 curl 验证端点,然后逐步提高并发,直到延迟——而不是运气——成为你的限制。