爬虫反反爬与分布式爬虫架构
1 反爬机制深度分析
1.1 请求头检测
网站通过校验 HTTP 请求头来区分正常浏览器访问和爬虫请求。
User-Agent 检测
User-Agent 是最基础的检测手段。正常浏览器的 User-Agent 包含浏览器名称、版本、操作系统等信息,而爬虫的 User-Agent 可能缺失或包含 Python-requests、curl、Scrapy 等特征字符串。
常见检测策略:
- 黑名单匹配:拦截已知爬虫 User-Agent 字符串
- 白名单验证:只允许已知浏览器 User-Agent 通过
- 版本号合理性:校验 User-Agent 中的版本号是否与当前主流浏览器一致
import re
SUSPICIOUS_UA_PATTERNS = [
r'Python-httpx', r'Python-requests', r'curl', r'wget', r'Scrapy',
r'Java/[\d.]+', r'Go-http-client', r'libcurl', r'HttpClient',
r'AHC/[\d.]+', r'okhttp', r'Dalvik', r'Apache-HttpClient'
]
def check_user_agent(ua: str) -> bool:
"""返回 True 表示疑似爬虫"""
if not ua or len(ua) < 20:
return True
for pattern in SUSPICIOUS_UA_PATTERNS:
if re.search(pattern, ua, re.IGNORECASE):
return True
return FalseCookie 检测
服务器通过检测 Cookie 的完整性、时效性和生成逻辑来判断请求是否来自真实浏览器。
- 首次访问无 Cookie:正常浏览器首次访问会携带空 Cookie,但后续请求必须携带服务器下发的 Cookie
- Cookie 生成逻辑:某些 Cookie 值由 JavaScript 计算生成,爬虫若未执行 JS 则无法获取正确值
- Session 一致性:同一会话的请求应使用相同的 Session ID
Referer / Origin 检测
- Referer:校验请求来源页面是否在网站域名内,跨域请求的 Referer 不符合预期则拦截
- Origin:CORS 场景下校验 Origin 头是否在白名单内
- 缺少 Referer:对静态资源请求缺少 Referer 视为异常
Sec-Fetch-* 头检测
现代浏览器会自动添加 Sec-Fetch-* 请求头系列:
Sec-Fetch-Site: same-origin | same-site | cross-site | none
Sec-Fetch-Mode: navigate | same-origin | cors | no-cors | websocket
Sec-Fetch-Dest: document | image | script | style | empty
Sec-Fetch-User: ?1 | ?0爬虫若未正确设置这些请求头,服务器可据此识别。
Accept / Content-Type 检测
- Accept 头:正常浏览器请求 HTML 时携带
text/html,application/xhtml+xml,...,而爬虫可能只携带*/*或application/json - Content-Type:POST 请求的 Content-Type 应与请求体格式匹配(如
application/x-www-form-urlencoded或multipart/form-data)
CSRF Token 检测
网站在表单中嵌入 CSRF Token,并在提交时校验。爬虫需要先抓取 Token 再提交请求,增加了采集难度。
import requests
from bs4 import BeautifulSoup
session = requests.Session()
# 先抓取登录页获取 CSRF Token
login_page = session.get('https://example.com/login')
soup = BeautifulSoup(login_page.text, 'html.parser')
csrf_token = soup.find('input', {'name': 'csrf_token'})['value']
# 提交登录请求
session.post('https://example.com/login', data={
'username': 'user',
'password': 'pass',
'csrf_token': csrf_token
})1.2 IP 检测与限流
单 IP QPS 限制
对单一 IP 的请求频率进行统计,超过阈值则返回 429 状态码或直接封禁。
import time
from collections import defaultdict
from datetime import datetime, timedelta
class SlidingWindowCounter:
"""滑动窗口计数限流"""
def __init__(self, limit: int = 60, window_secs: int = 60):
self.limit = limit
self.window_secs = window_secs
self.counters = defaultdict(list) # ip -> [timestamp, ...]
def allow(self, ip: str) -> bool:
now = time.time()
window_start = now - self.window_secs
# 移除窗口外的记录
self.counters[ip] = [t for t in self.counters[ip] if t > window_start]
if len(self.counters[ip]) >= self.limit:
return False
self.counters[ip].append(now)
return True令牌桶限流
import time
import threading
class TokenBucket:
"""令牌桶限流"""
def __init__(self, rate: float, capacity: int):
self.rate = rate # 每秒生成令牌数
self.capacity = capacity # 桶容量
self.tokens = capacity
self.last_refill = time.time()
self.lock = threading.Lock()
def allow(self) -> bool:
with self.lock:
now = time.time()
elapsed = now - self.last_refill
self.tokens = min(self.capacity, self.tokens + elapsed * self.rate)
self.last_refill = now
if self.tokens >= 1:
self.tokens -= 1
return True
return False频率分析
通过统计单位时间内的请求次数、请求间隔的均值和方差来识别机器行为。正常用户的请求间隔呈现自然分布,而爬虫通常呈现均匀或 Burst 模式。
地理位置与 ASN 检测
- 地理位置:检测 IP 所属国家/城市,对非目标地区 IP 进行限制或验证
- ASN 检测:识别 IP 所属 AS(自治系统),对已知机房 IP(AWS、GCP、阿里云等)进行严格限制
CDN IP 段识别
维护云服务商 IP 段列表,对来自 CDN 或云主机的请求提高验证强度。
CLOUD_IP_RANGES = {
'aws': ['3.0.0.0/9', '35.0.0.0/8'],
'gcp': ['34.0.0.0/8', '35.192.0.0/12'],
'azure': ['13.64.0.0/11', '20.0.0.0/8'],
'aliyun': ['47.0.0.0/8', '59.0.0.0/8'],
}
def is_cloud_ip(ip: str) -> bool:
import ipaddress
addr = ipaddress.ip_address(ip)
for ranges in CLOUD_IP_RANGES.values():
for cidr in ranges:
if addr in ipaddress.ip_network(cidr):
return True
return False动态黑名单
基于实时分析将异常 IP 加入动态黑名单,设置不同的封禁时长(5 分钟、1 小时、24 小时、永久)。
验证码触发阈值
当请求频率或异常行为达到一定阈值时,触发验证码验证:
| 指标 | 低风险阈值 | 中风险阈值 | 高风险阈值 |
|---|---|---|---|
| 单 IP QPS | < 5 | 5-20 | > 20 |
| 请求间隔变异系数 | > 0.8 | 0.3-0.8 | < 0.3 |
| 页面停留时间 | > 10s | 3-10s | < 3s |
| 无 JS 环境请求占比 | < 5% | 5-30% | > 30% |
1.3 行为分析
鼠标轨迹
真实用户的鼠标移动轨迹呈现自然曲线,具有加速度和微小的抖动。爬虫的模拟轨迹通常过于规则或直线移动。
// 前端采集鼠标轨迹
document.addEventListener('mousemove', function(e) {
const data = {
x: e.clientX,
y: e.clientY,
t: Date.now(),
type: 'mousemove'
};
// 通过 Beacon API 异步发送
navigator.sendBeacon('/collect', JSON.stringify(data));
});点击间隔与页面停留时间
- 点击间隔:真实用户的点击间隔符合正太分布,爬虫的点击间隔可能过于均匀
- 页面停留时间:正常用户阅读页面需要数秒到数分钟,爬虫可能在毫秒级完成"阅读"
滚动速度
真实用户的页面滚动速度不均匀,存在停顿和反复。爬虫的滚动通常速度恒定或直接跳转到底部。
浏览器窗口尺寸
采集 window.innerWidth、window.innerHeight、window.screen 等信息,爬虫的窗口尺寸可能不符合常规显示器分辨率。
操作序列
记录用户在页面上的操作序列(点击→滚动→悬停→输入),真实用户的操作序列具有随机性和多样性,爬虫的操作序列高度重复。
浏览器指纹
现代反爬系统综合多种指纹信息进行设备识别:
// WebGL 指纹
const canvas = document.createElement('canvas');
const gl = canvas.getContext('webgl');
const debugInfo = gl.getExtension('WEBGL_debug_renderer_info');
const renderer = gl.getParameter(debugInfo.UNMASKED_RENDERER_WEBGL);
const vendor = gl.getParameter(debugInfo.UNMASKED_VENDOR_WEBGL);
// Canvas 指纹
const c = document.createElement('canvas');
c.width = 200; c.height = 50;
const ctx = c.getContext('2d');
ctx.textBaseline = 'top';
ctx.fillText('CanvasFingerprint', 2, 2);
const fingerprint = c.toDataURL();
// AudioContext 指纹
const audioCtx = new (window.AudioContext || window.webkitAudioContext)();
const analyser = audioCtx.createAnalyser();
const oscillator = audioCtx.createOscillator();
oscillator.connect(analyser);
// 获取音频处理结果作为指纹
// Font 指纹
function detectFonts() {
const testFonts = ['Arial', 'Verdana', 'Times New Roman', 'Courier New'];
const baseWidth = measureText('mmmmmmmmmmlli', 'monospace');
return testFonts.map(font => ({
name: font,
width: measureText('mmmmmmmmmmlli', font)
}));
}
// 时区
const timezone = Intl.DateTimeFormat().resolvedOptions().timeZone;
// Battery 信息
navigator.getBattery().then(function(battery) {
console.log(battery.level, battery.charging);
});
// WebRTC 内网 IP
const pc = new RTCPeerConnection({ iceServers: [] });
pc.createDataChannel('');
pc.createOffer().then(pc.setLocalDescription.bind(pc));
pc.onicecandidate = function(e) {
if (e.candidate && e.candidate.candidate) {
// 解析 candidate 获取内网 IP
}
};1.4 JS 挑战与验证
Cloudflare Challenge (5 秒盾)
Cloudflare 的防爬机制会在首次访问时返回一个 JS 挑战页面,浏览器需要执行 JavaScript 计算并通过验证后才能访问实际内容。挑战过程:
- 服务器返回包含 JS 计算任务的 HTML 页面
- 浏览器执行 JS 代码,进行数学计算或 Proof-of-Work
- 计算完成后设置
cf_clearanceCookie - 携带 Cookie 重新请求获取真实内容
JavaScript 计算结果验证
服务器下发一段 JavaScript 代码,客户端执行后返回计算结果。爬虫若未执行 JS 则无法通过验证。
// 服务器下发的 JS 挑战示例
function challenge(a, b, c) {
var d = a * b + c;
var e = d.toString(16);
var f = e.split('').reverse().join('');
return f;
}
// 要求返回 challenge(127, 89, 43) 的结果Cookie 验证(Cookie Challenge)
服务器通过 JavaScript 设置特定的 Cookie 值,后续请求校验该 Cookie:
__cfduid:Cloudflare 用于识别访客_ga/_gid:Google Analytics 跟踪 Cookie- 自定义加密 Cookie:随时间和 IP 变化的签名 Cookie
浏览器环境检测
// 检测正常浏览器环境
function detectHeadless() {
const checks = {
// Chrome Headless 检测
chromeHeadless: navigator.webdriver === true,
// 缺少 plugins
noPlugins: navigator.plugins.length === 0,
// 缺少 mimeTypes
noMimeTypes: navigator.mimeTypes.length === 0,
// languages 异常
languages: navigator.languages.length === 0,
// WebDriver 属性
webdriver: window.navigator.webdriver,
$chrome_asyncScript: window.$chrome_asyncScript,
};
return checks;
}WebDriver 检测
// 检测是否通过 Selenium/Puppeteer 控制
function detectWebDriver() {
const indicators = [
document.documentElement.getAttribute('webdriver'),
window.navigator.webdriver,
window._phantom,
window.__nightmare,
window.__selenium_evaluate,
window.document.__selenium_unwrapped,
window.callPhantom,
window.chrome && window.chrome.runtime,
];
return indicators.some(i => i !== undefined && i !== null);
}Chrome Headless 检测
// 检测 Headless Chrome 的多种方法
function detectHeadlessChrome() {
const tests = {
// Headless Chrome 中 UserAgent 不包含 Headless 字样
userAgent: /Headless/.test(navigator.userAgent),
// 图片缺少某些属性
imageMime: !HTMLImageElement.prototype.decode,
// 权限查询行为差异
permissionQuery: false,
};
// Headless Chrome 的 permissions.query 行为不同
if (navigator.permissions) {
navigator.permissions.query({ name: 'notifications' })
.then(p => { tests.permissionQuery = p.state === 'denied'; });
}
return tests;
}2 反反爬策略
2.1 IP 代理池
代理来源
| 来源类型 | 优点 | 缺点 | 典型服务 |
|---|---|---|---|
| 付费代理 | 稳定、高速、可用率高 | 成本较高 | 快代理、芝麻代理、讯代理 |
| 免费代理 | 零成本 | 稳定性差、可用率低、匿名性差 | 西刺、ProxyList+ |
| 自建代理 | 完全可控、IP 独享 | 需要服务器资源 | 云主机搭建 Squid/HAProxy |
| 住宅代理 | 匿名性极高、不易被检测 | 价格昂贵 | Luminati/oxylabs/Smartproxy |
| 机房代理 | 速度块、价格低 | 易被识别为数据中心 IP | AWS/GCP 弹性 IP |
代理验证
代理在使用前需要进行验证,确保其可用性和匿名性:
import requests
import asyncio
import aiohttp
async def verify_proxy(proxy: str, test_url: str = 'https://httpbin.org/ip') -> dict:
"""验证代理可用性"""
start = time.time()
try:
async with aiohttp.ClientSession() as session:
async with session.get(test_url, proxy=f'http://{proxy}', timeout=10) as resp:
latency = time.time() - start
if resp.status == 200:
data = await resp.json()
return {
'proxy': proxy,
'ok': True,
'latency': round(latency, 3),
'ip': data.get('origin', ''),
'anonymous': True
}
except Exception as e:
return {'proxy': proxy, 'ok': False, 'error': str(e)}
return {'proxy': proxy, 'ok': False}匿名级别
| 级别 | 说明 | 检测方式 |
|---|---|---|
| 透明代理 | 请求头携带 X-Forwarded-For,服务器知道真实 IP | HTTP_X_FORWARDED_FOR 存在 |
| 匿名代理 | 不携带真实 IP,但声明了代理身份 | HTTP_VIA 或 HTTP_PROXY_CONNECTION 存在 |
| 高匿代理 | 不携带任何代理信息,完全模拟原生请求 | 无任何代理相关头信息 |
def check_anonymity(proxy_response_headers: dict) -> str:
"""检查代理匿名级别"""
headers_lower = {k.lower(): v for k, v in proxy_response_headers.items()}
if 'x-forwarded-for' in headers_lower:
return 'transparent'
if 'via' in headers_lower or 'proxy-connection' in headers_lower:
return 'anonymous'
return 'elite' # 高匿Socks5 vs HTTP 代理
| 特性 | HTTP 代理 | Socks5 代理 |
|---|---|---|
| 协议支持 | HTTP/HTTPS | TCP/UDP 全协议 |
| 认证方式 | Basic Auth | 用户名密码/GSSAPI |
| 性能 | 需解析 HTTP 协议 | 更轻量 |
| 适用场景 | Web 爬虫 | 全场景(含 TCP/UDP) |
| Python 支持 | requests 原生支持 | 需 PySocks 或 aiohttp-socks |
# HTTP 代理
requests.get('https://example.com', proxies={'http': 'http://proxy:8080'})
# Socks5 代理
import socks
import socket
socks.set_default_proxy(socks.SOCKS5, 'proxy_host', 1080)
socket.socket = socks.socksocket
requests.get('https://example.com')代理质量检测指标
- 延迟:请求响应时间,通常要求 < 3s
- 可用率:一段时间内代理可用的比例,要求 > 80%
- 地区:与目标网站所在地区匹配,减少检测风险
- 匿名级别:高匿 > 匿名 > 透明
- 支持协议:HTTPS > HTTP
轮换策略
import random
from collections import deque
import time
class ProxyRotator:
"""代理轮换器"""
ROTATE_EVERY_REQUEST = 'request' # 每个请求换 IP
ROTATE_EVERY_N = 'every_n' # 每 N 个请求换 IP
ROTATE_ON_ERROR = 'on_error' # 错误时切换
def __init__(self, proxies: list, strategy: str = 'request',
rotate_interval: int = 10):
self.proxy_pool = deque(proxies)
self.strategy = strategy
self.rotate_interval = rotate_interval
self.request_count = 0
self.current_proxy = None
self.failed_proxies = set()
def get_proxy(self) -> str:
"""获取当前代理"""
if self.strategy == self.ROTATE_EVERY_REQUEST:
self._rotate()
elif self.strategy == self.ROTATE_EVERY_N:
self.request_count += 1
if self.request_count >= self.rotate_interval:
self._rotate()
self.request_count = 0
else:
if not self.current_proxy:
self._rotate()
return self.current_proxy
def report_error(self, proxy: str):
"""报告代理错误,触发切换"""
self.failed_proxies.add(proxy)
self.proxy_pool = deque(
p for p in self.proxy_pool if p != proxy
)
if self.strategy == self.ROTATE_ON_ERROR:
self._rotate()
def _rotate(self):
if self.proxy_pool:
self.proxy_pool.rotate(-1)
self.current_proxy = self.proxy_pool[0]
def remove_failed(self, max_failures: int = 3):
"""剔除失效代理"""
self.proxy_pool = deque(
p for p in self.proxy_pool
if p not in self.failed_proxies
)2.2 请求伪装
请求头顺序随机化
不同浏览器的 HTTP 头顺序存在差异,随机化请求头顺序有助于规避基于头部顺序的检测。
import random
from collections import OrderedDict
import requests
def randomize_headers(base_headers: dict) -> OrderedDict:
"""随机化请求头顺序"""
items = list(base_headers.items())
random.shuffle(items)
return OrderedDict(items)
# 使用随机化请求头
session = requests.Session()
session.headers = randomize_headers({
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
'Accept-Language': 'zh-CN,zh;q=0.9,en;q=0.8',
'Accept-Encoding': 'gzip, deflate, br',
'Connection': 'keep-alive',
'Upgrade-Insecure-Requests': '1',
'Sec-Fetch-Dest': 'document',
'Sec-Fetch-Mode': 'navigate',
'Sec-Fetch-Site': 'none',
'Sec-Fetch-User': '?1',
})浏览器指纹随机化
使用 Playwright/Puppeteer 时,通过参数注入随机化浏览器指纹:
import random
from playwright.sync_api import sync_playwright
def create_stealth_browser():
"""创建带随机指纹的浏览器实例"""
playwright = sync_playwright().start()
# 随机化窗口尺寸
width = random.choice([1366, 1440, 1536, 1680, 1920])
height = random.choice([768, 900, 864, 1050, 1080])
browser = playwright.chromium.launch(
headless=True,
args=[
f'--window-size={width},{height}',
'--disable-blink-features=AutomationControlled'
]
)
context = browser.new_context(
viewport={'width': width, 'height': height},
user_agent=random.choice(USER_AGENTS),
locale='zh-CN',
timezone_id='Asia/Shanghai',
# 随机化硬件并发数
hardware_concurrency=random.choice([4, 6, 8, 12, 16]),
)
return browser, contextNavigator 属性覆盖
// 通过 Playwright 注入脚本覆盖 navigator 属性
const overrideScript = `
Object.defineProperty(navigator, 'webdriver', { get: () => false });
Object.defineProperty(navigator, 'plugins', {
get: () => [
{ name: 'Chrome PDF Plugin', filename: 'internal-pdf-viewer' },
{ name: 'Chrome PDF Viewer', filename: 'mhjfbmdgcfjbbpaeojofohoefgiehjai' },
{ name: 'Native Client', filename: 'internal-nacl-plugin' }
]
});
Object.defineProperty(navigator, 'mimeTypes', {
get: () => [
{ type: 'application/pdf', suffixes: 'pdf', description: 'Portable Document Format' }
]
});
Object.defineProperty(navigator, 'languages', { get: () => ['zh-CN', 'zh', 'en'] });
Object.defineProperty(navigator, 'hardwareConcurrency', { get: () => 8 });
Object.defineProperty(navigator, 'deviceMemory', { get: () => 8 });
Object.defineProperty(navigator, 'platform', { get: () => 'Win32' });
`;WebRTC 禁用(防止内网 IP 泄露)
# Playwright 中禁用 WebRTC
context = browser.new_context(
permissions=[],
# 通过 args 禁用 WebRTC
)
# 或者在浏览器 args 中
browser = playwright.chromium.launch(
args=[
'--disable-webrtc',
'--enforce-webrtc-ip-permission-check',
'--force-webrtc-ip-handling-policy=disable_non_proxied_udp'
]
)Canvas / WebGL / Font 指纹随机化
// 注入脚本随机化 Canvas 指纹
const randomizeCanvas = `
const originalToDataURL = HTMLCanvasElement.prototype.toDataURL;
HTMLCanvasElement.prototype.toDataURL = function() {
const imageData = this.getContext('2d').getImageData(0, 0, this.width, this.height);
// 对像素数据做微小的随机扰动
for (let i = 0; i < imageData.data.length; i += 4) {
imageData.data[i] += Math.floor(Math.random() * 2); // R
imageData.data[i+1] += Math.floor(Math.random() * 2); // G
imageData.data[i+2] += Math.floor(Math.random() * 2); // B
}
this.getContext('2d').putImageData(imageData, 0, 0);
return originalToDataURL.apply(this, arguments);
};
`;时序随机化与点击移动模拟
import asyncio
import random
async def human_like_mouse_move(page, target_x: int, target_y: int):
"""模拟人类鼠标移动"""
current_x = random.randint(100, 500)
current_y = random.randint(100, 500)
steps = random.randint(8, 15)
for i in range(steps):
progress = (i + 1) / steps
# 贝塞尔曲线插值,模拟自然曲线
x = current_x + (target_x - current_x) * progress + random.randint(-5, 5)
y = current_y + (target_y - current_y) * progress + random.randint(-5, 5)
await page.mouse.move(x, y)
await asyncio.sleep(random.uniform(0.01, 0.05)) # 随机延迟
async def human_like_click(page, selector: str):
"""模拟人类点击"""
element = await page.wait_for_selector(selector)
box = await element.bounding_box()
# 点击位置在元素中心附近随机偏移
target_x = box['x'] + box['width'] * random.uniform(0.3, 0.7)
target_y = box['y'] + box['height'] * random.uniform(0.3, 0.7)
await human_like_mouse_move(page, target_x, target_y)
await asyncio.sleep(random.uniform(0.1, 0.3))
await page.mouse.click(target_x, target_y)2.3 验证码对抗方案
OCR 识别
| 方案 | 适用场景 | 准确率 | 速度 | 成本 |
|---|---|---|---|---|
| Tesseract OCR | 简单数字/字母验证码 | 50-70% | 快 | 免费 |
| ddddocr | 常见图形验证码 | 70-90% | 快 | 免费 |
| 深度学习 CNN | 固定类型验证码 | 90-95% | 中 | 需训练 |
| 打码平台 | 复杂验证码 | 85-98% | 中 | 按次计费 |
Tesseract OCR 示例
import pytesseract
from PIL import Image
def ocr_captcha(image_path: str) -> str:
"""使用 Tesseract 识别验证码"""
img = Image.open(image_path)
# 预处理:灰度化 + 二值化 + 降噪
img = img.convert('L')
threshold = 128
img = img.point(lambda x: 255 if x > threshold else 0)
# 识别
code = pytesseract.image_to_string(img, config='--psm 8 --oem 3')
return code.strip()ddddocr 示例
import ddddocr
ocr = ddddocr.DdddOcr()
with open('captcha.png', 'rb') as f:
image_bytes = f.read()
result = ocr.classification(image_bytes)
print(f'验证码识别结果: {result}')打码平台接入
import requests
import time
class CaptchaSolver:
"""打码平台对接(以 2captcha 为例)"""
def __init__(self, api_key: str):
self.api_key = api_key
self.base_url = 'https://2captcha.com'
def solve_normal(self, image_path: str) -> str:
"""识别普通图片验证码"""
with open(image_path, 'rb') as f:
files = {'file': f}
resp = requests.post(
f'{self.base_url}/in.php',
files=files,
data={'key': self.api_key, 'method': 'post'}
)
if resp.text.startswith('OK|'):
captcha_id = resp.text.split('|')[1]
return self._poll_result(captcha_id)
return ''
def solve_recaptcha_v2(self, site_key: str, page_url: str) -> str:
"""识别 reCAPTCHA v2"""
resp = requests.post(f'{self.base_url}/in.php', data={
'key': self.api_key,
'method': 'userrecaptcha',
'googlekey': site_key,
'pageurl': page_url,
})
if resp.text.startswith('OK|'):
captcha_id = resp.text.split('|')[1]
return self._poll_result(captcha_id)
return ''
def _poll_result(self, captcha_id: str, timeout: int = 120) -> str:
"""轮询获取识别结果"""
start = time.time()
while time.time() - start < timeout:
resp = requests.get(f'{self.base_url}/res.php', params={
'key': self.api_key,
'action': 'get',
'id': captcha_id
})
if resp.text == 'CAPCHA_NOT_READY':
time.sleep(5)
continue
if resp.text.startswith('OK|'):
return resp.text.split('|')[1]
break
return ''验证码方案对比
| 验证码类型 | ddddocr | 打码平台 | 深度学习定制 | 建议方案 |
|---|---|---|---|---|
| 数字字母验证码 | 80-90% | 90-95% | 95-98% | ddddocr |
| 中文验证码 | 60-70% | 85-90% | 90-95% | 打码平台 |
| 滑块验证码 | 不支持 | 85-95% | 需专门模型 | 打码平台 |
| 点选验证码 | 不支持 | 80-90% | 需专门模型 | 打码平台 |
| reCAPTCHA v2 | 不支持 | 85-95% | 不支持 | 打码平台 |
| reCAPTCHA v3 | 不支持 | 部分支持 | 不支持 | Playwright 模拟 |
2.4 JS 挑战处理
Puppeteer / Playwright 自动执行
from playwright.sync_api import sync_playwright
import time
def solve_cloudflare_challenge(url: str) -> str:
"""处理 Cloudflare JS 挑战"""
with sync_playwright() as p:
browser = p.chromium.launch(
headless=False, # 首次可设为 True,失败时改用 False
args=['--disable-blink-features=AutomationControlled']
)
context = browser.new_context(
user_agent=(
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) '
'AppleWebKit/537.36 (KHTML, like Gecko) '
'Chrome/120.0.0.0 Safari/537.36'
),
locale='zh-CN',
timezone_id='Asia/Shanghai'
)
page = context.new_page()
page.goto(url, wait_until='networkidle')
# 等待 Cloudflare 挑战完成
try:
page.wait_for_selector('#challenge-form', timeout=5000)
# 正在处理 JS 挑战,等待完成
page.wait_for_selector('#challenge-form', state='hidden', timeout=30000)
except:
pass
# 获取最终页面内容和 Cookie
content = page.content()
cookies = context.cookies()
browser.close()
return content, cookiesCloudflare 通过率优化
class CloudflareSolver:
"""Cloudflare 挑战处理器"""
def __init__(self):
self.context = None
self.browser = None
def create_persistent_context(self):
"""创建持久化浏览器上下文"""
from playwright.sync_api import sync_playwright
self.playwright = sync_playwright().start()
self.browser = self.playwright.chromium.launch(
headless=True,
args=[
'--disable-blink-features=AutomationControlled',
'--no-sandbox',
'--disable-setuid-sandbox',
]
)
# 使用持久化的 Cookie 和本地存储
self.context = self.browser.new_context(
storage_state='cloudflare_state.json' if os.path.exists('cloudflare_state.json') else None,
user_agent=random.choice(REAL_USER_AGENTS),
viewport={'width': 1920, 'height': 1080},
)
def save_state(self):
"""保存浏览器状态以供复用"""
if self.context:
self.context.storage_state(path='cloudflare_state.json')浏览器上下文持久化
通过持久化浏览器上下文(Cookie、LocalStorage、Session Storage),避免每次请求都重新执行 JS 挑战。
| 策略 | 说明 | Cookie 有效期 |
|---|---|---|
| 内存存储 | 运行期间保持 Cookie | 进程生命周期 |
| 文件存储 | 序列化到 JSON 文件 | 多次运行间共享 |
| Redis 存储 | 分布式共享 Cookie | 跨节点共享 |
| Session Pool | 维护多个会话池 | 多个 IP 维度 |
import json
import redis
class SessionPersister:
"""会话持久化管理器"""
def __init__(self, redis_url: str = 'redis://localhost:6379/0'):
self.redis = redis.from_url(redis_url)
def save_session(self, key: str, cookies: list):
"""保存会话 Cookie"""
self.redis.setex(
f'session:{key}',
3600, # 1 小时过期
json.dumps(cookies)
)
def load_session(self, key: str) -> list:
"""加载会话 Cookie"""
data = self.redis.get(f'session:{key}')
return json.loads(data) if data else []
def refresh_session(self, key: str, cookies: list):
"""刷新会话"""
ttl = self.redis.ttl(f'session:{key}')
if ttl > 0:
self.redis.setex(f'session:{key}', ttl, json.dumps(cookies))3 分布式爬虫架构
3.1 任务队列
Redis 队列
基于 Redis 实现高效的爬虫任务队列:
import redis
import json
import time
import hashlib
class RedisTaskQueue:
"""基于 Redis 的爬虫任务队列"""
def __init__(self, host='localhost', port=6379, db=0):
self.client = redis.Redis(host=host, port=port, db=db,
decode_responses=True)
# ---------- 普通 FIFO 队列 ----------
def push(self, queue: str, task: dict):
"""使用 LPUSH 将任务推入队列"""
self.client.lpush(queue, json.dumps(task))
def pop(self, queue: str, timeout: int = 0) -> dict:
"""使用 BRPOP 阻塞弹出任务"""
result = self.client.brpop(queue, timeout=timeout)
if result:
return json.loads(result[1])
return None
# ---------- 优先级队列 ----------
def push_priority(self, queue: str, task: dict, priority: int = 0):
"""使用 ZSET 实现优先级队列,priority 越大越优先"""
score = time.time() - priority * 10000 # 优先级权重
self.client.zadd(f'{queue}:priority', {json.dumps(task): -score})
def pop_priority(self, queue: str) -> dict:
"""弹出优先级最高的任务"""
results = self.client.zpopmin(f'{queue}:priority', count=1)
if results:
return json.loads(results[0][0])
return None
# ---------- 延迟队列 ----------
def push_delayed(self, queue: str, task: dict, delay_secs: int):
"""使用 ZSET 实现延迟队列"""
execute_at = time.time() + delay_secs
self.client.zadd(f'{queue}:delayed', {json.dumps(task): execute_at})
def poll_delayed(self, queue: str) -> list:
"""轮询到期的延迟任务"""
now = time.time()
tasks = self.client.zrangebyscore(
f'{queue}:delayed', 0, now
)
if tasks:
self.client.zremrangebyscore(f'{queue}:delayed', 0, now)
return [json.loads(t) for t in tasks]
return []
# ---------- 去重集合 ----------
def seen(self, key: str, value: str, expire: int = 86400) -> bool:
"""检查并记录已处理的任务"""
dedup_key = f'seen:{key}'
if self.client.sadd(dedup_key, value):
self.client.expire(dedup_key, expire)
return False # 新任务
return True # 已处理过消息队列选型对比
| 特性 | Redis | RabbitMQ | Kafka |
|---|---|---|---|
| 吞吐量 | 10万+/s | 1-10万/s | 百万+/s |
| 持久化 | 支持(RDB/AOF) | 支持 | 支持 |
| 延迟 | 微秒级 | 毫秒级 | 毫秒级 |
| 消息可靠性 | 一般 | 高(Confirm+持久化) | 高(ISR 副本) |
| 消费者组 | 不支持原生 | 支持 | 支持 |
| 死信队列 | 手动实现 | 原生支持 | 需要配置 |
| 适用场景 | 小规模分布式 | 中等规模、需要可靠投递 | 大规模、高吞吐 |
| 运维复杂度 | 低 | 中 | 高 |
# RabbitMQ 作为任务队列
import pika
import json
class RabbitMQTaskQueue:
def __init__(self, host='localhost', queue='crawler_tasks'):
self.connection = pika.BlockingConnection(
pika.ConnectionParameters(host=host)
)
self.channel = self.connection.channel()
self.queue = queue
self.channel.queue_declare(queue=queue, durable=True)
def push(self, task: dict):
self.channel.basic_publish(
exchange='',
routing_key=self.queue,
body=json.dumps(task),
properties=pika.BasicProperties(delivery_mode=2) # 持久化
)
def consume(self, callback):
self.channel.basic_qos(prefetch_count=1)
self.channel.basic_consume(
queue=self.queue,
on_message_callback=lambda ch, method, properties, body: (
callback(json.loads(body)),
ch.basic_ack(delivery_tag=method.delivery_tag)
)
)
self.channel.start_consuming()3.2 爬虫节点设计
Master 调度器
import redis
import json
import time
import uuid
from typing import Optional, Dict, List
class MasterScheduler:
"""Master 调度器:负责任务分发和节点管理"""
def __init__(self, redis_client: redis.Redis):
self.redis = redis_client
self.worker_ttl = 60 # Worker 心跳超时
def register_worker(self, worker_id: str, capacity: int = 10) -> bool:
"""注册 Worker 节点"""
self.redis.hset('workers', worker_id, json.dumps({
'id': worker_id,
'capacity': capacity,
'current_load': 0,
'status': 'idle',
'last_heartbeat': time.time(),
'registered_at': time.time()
}))
return True
def heartbeat(self, worker_id: str, load: int) -> bool:
"""接收 Worker 心跳"""
key = f'worker:{worker_id}:hb'
self.redis.setex(key, self.worker_ttl, json.dumps({
'load': load,
'timestamp': time.time()
}))
return True
def dispatch_task(self, task: dict) -> Optional[str]:
"""分发任务到合适的 Worker"""
workers = self.get_alive_workers()
if not workers:
return None
# 选择负载最低的 Worker
best_worker = min(workers, key=lambda w: w['load'])
# 将任务推入该 Worker 的任务队列
self.redis.lpush(f'worker:{best_worker["id"]}:tasks', json.dumps(task))
return best_worker['id']
def get_alive_workers(self) -> List[Dict]:
"""获取所有活跃的 Worker"""
workers = []
for key in self.redis.scan_iter('worker:*:hb'):
data = self.redis.get(key)
if data:
worker_info = json.loads(data)
workers.append({
'id': key.split(':')[1],
'load': worker_info['load'],
'last_seen': worker_info['timestamp']
})
return workers
def get_task_stats(self) -> Dict:
"""获取任务统计"""
return {
'pending': self.redis.llen('tasks:pending'),
'processing': self.redis.scard('tasks:processing'),
'completed': self.redis.get('tasks:completed') or 0,
'failed': self.redis.get('tasks:failed') or 0,
}Worker 执行器
import redis
import json
import time
import threading
import requests
from typing import Optional
class WorkerNode:
"""Worker 节点:执行爬虫任务"""
def __init__(self, worker_id: str, master_redis: redis.Redis,
max_concurrent: int = 5):
self.worker_id = worker_id
self.redis = master_redis
self.max_concurrent = max_concurrent
self.current_load = 0
self.running = True
def start(self):
"""启动 Worker"""
# 启动心跳线程
heartbeat_thread = threading.Thread(target=self._heartbeat_loop)
heartbeat_thread.daemon = True
heartbeat_thread.start()
# 启动任务处理
self._task_loop()
def _heartbeat_loop(self):
"""心跳线程"""
while self.running:
self.redis.setex(
f'worker:{self.worker_id}:hb',
60,
json.dumps({'load': self.current_load, 'timestamp': time.time()})
)
time.sleep(15)
def _task_loop(self):
"""任务处理循环"""
task_queue = f'worker:{self.worker_id}:tasks'
while self.running:
if self.current_load < self.max_concurrent:
result = self.redis.brpop(task_queue, timeout=5)
if result:
task = json.loads(result[1])
self.current_load += 1
# 使用线程池执行任务
thread = threading.Thread(
target=self._execute_task,
args=(task,)
)
thread.start()
def _execute_task(self, task: dict):
"""执行单个爬虫任务"""
task_id = task.get('id', str(time.time()))
try:
# 标记任务开始
self.redis.hset('tasks:status', task_id, 'processing')
# 执行爬取逻辑
result = self._crawl(task)
# 上报结果
self._report_result(task_id, result)
self.redis.incr('tasks:completed')
except Exception as e:
self._report_failure(task_id, str(e))
self.redis.incr('tasks:failed')
# 失败重试逻辑
retries = task.get('retries', 0)
if retries < task.get('max_retries', 3):
task['retries'] = retries + 1
task['retry_delay'] = task.get('retry_delay', 5) * 2
time.sleep(task['retry_delay'])
self.redis.lpush(f'worker:{self.worker_id}:tasks',
json.dumps(task))
finally:
self.current_load -= 1
def _crawl(self, task: dict) -> dict:
"""实际爬取逻辑"""
url = task['url']
headers = task.get('headers', {})
proxies = task.get('proxies', {})
resp = requests.get(
url,
headers=headers,
proxies=proxies,
timeout=30
)
return {
'url': url,
'status': resp.status_code,
'content_length': len(resp.text),
'content': resp.text
}
def _report_result(self, task_id: str, result: dict):
"""上报结果到结果队列"""
self.redis.lpush('results', json.dumps({
'task_id': task_id,
'worker': self.worker_id,
'result': result,
'timestamp': time.time()
}))
self.redis.hset('tasks:status', task_id, 'completed')
def _report_failure(self, task_id: str, error: str):
"""上报失败信息"""
self.redis.hset('tasks:status', task_id, 'failed')
self.redis.lpush('tasks:failed:log', json.dumps({
'task_id': task_id,
'worker': self.worker_id,
'error': error,
'timestamp': time.time()
}))3.3 去重策略
Redis Set 去重
class RedisSetDeduplicator:
"""基于 Redis Set 的去重"""
def __init__(self, redis_client: redis.Redis, key: str = 'dedup:urls'):
self.redis = redis_client
self.key = key
def is_duplicate(self, value: str) -> bool:
"""检查是否重复"""
return self.redis.sismember(self.key, value)
def mark_processed(self, value: str) -> bool:
"""标记已处理"""
return self.redis.sadd(self.key, value) == 1 # 返回 True 表示新添加
def count(self) -> int:
return self.redis.scard(self.key)
def clear(self):
self.redis.delete(self.key)Bloom Filter (布隆过滤器)
布隆过滤器是一种空间效率极高的概率型数据结构,用于判断一个元素是否在集合中。
特点:
- 判断"不在集合中"是绝对准确的
- 判断"在集合中"有一定误判率(False Positive)
- 不支持删除元素
import math
import hashlib
import redis
class BloomFilter:
"""基于 Redis 的布隆过滤器"""
def __init__(self, redis_client: redis.Redis, key: str,
capacity: int = 1000000, error_rate: float = 0.01):
self.redis = redis_client
self.key = key
self.capacity = capacity
self.error_rate = error_rate
# 计算最优的位数组大小和哈希函数数量
self.bit_size = self._optimal_bit_size()
self.hash_count = self._optimal_hash_count()
def _optimal_bit_size(self) -> int:
"""计算最优位数组大小(位)"""
return int(-self.capacity * math.log(self.error_rate) / (math.log(2) ** 2))
def _optimal_hash_count(self) -> int:
"""计算最优哈希函数数量"""
return int(self.bit_size / self.capacity * math.log(2))
def _get_offsets(self, item: str) -> list:
"""计算 item 的哈希偏移位置"""
offsets = []
for i in range(self.hash_count):
# 使用不同的 seed 计算多个哈希
hash_val = hashlib.md5(f'{item}:{i}'.encode()).hexdigest()
offset = int(hash_val, 16) % self.bit_size
offsets.append(offset)
return offsets
def add(self, item: str):
"""添加元素"""
offsets = self._get_offsets(item)
for offset in offsets:
self.redis.setbit(self.key, offset, 1)
def contains(self, item: str) -> bool:
"""检查元素是否可能存在"""
offsets = self._get_offsets(item)
for offset in offsets:
if not self.redis.getbit(self.key, offset):
return False
return True # 可能存在(有误判可能)
def estimated_size(self) -> int:
"""估算已存储的元素数量"""
# 基于已设置的位数估算
bits_set = self.redis.bitcount(self.key)
if bits_set == 0:
return 0
ratio = bits_set / self.bit_size
return int(-self.bit_size / self.hash_count * math.log(1 - ratio))布隆过滤器误判率与大小计算
误判率公式: P = (1 - e^(-kn/m))^k
其中:
- n = 已插入元素数量
- m = 位数组大小(位)
- k = 哈希函数数量
给定 n 和期望误判率 p,最优参数:
m = -n * ln(p) / (ln2)^2
k = m/n * ln2
示例:
n = 1000万, p = 1% => m ≈ 120MB, k = 7
n = 1亿, p = 0.1% => m ≈ 1.8GB, k = 10| 元素数量 | 期望误判率 | 位数组大小 | 哈希函数数 | 内存占用 |
|---|---|---|---|---|
| 100万 | 1% | 9.6 Mb | 7 | 1.2 MB |
| 1000万 | 1% | 96 Mb | 7 | 12 MB |
| 1亿 | 1% | 958 Mb | 7 | 120 MB |
| 1亿 | 0.1% | 1.4 Gb | 10 | 180 MB |
RoaringBitmap
RoaringBitmap 是一种高效的压缩位图数据结构,相比传统 BitSet 大幅节省内存:
from roaringbitmap import RoaringBitmap
class RoaringBitmapDeduplicator:
"""使用 RoaringBitmap 进行去重(将 URL 映射为整数 ID)"""
def __init__(self):
self.bitmap = RoaringBitmap()
def add(self, doc_id: int) -> bool:
"""添加文档 ID,返回 True 表示新添加"""
if doc_id in self.bitmap:
return False
self.bitmap.add(doc_id)
return True
def contains(self, doc_id: int) -> bool:
return doc_id in self.bitmap
def size(self) -> int:
return len(self.bitmap)
def serialize(self) -> bytes:
return self.bitmap.serialize()
def deserialize(self, data: bytes):
self.bitmap = RoaringBitmap.deserialize(data)去重字段设计
| 去重依据 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| URL | 简单直接 | 同一内容可能多个 URL 访问 | 通用爬虫 |
| URL 规范化 | 避免同义 URL | 需要 URL 归一化逻辑 | 新闻/文章采集 |
| 内容 Hash (MD5/SHA1) | 精确去重 | 需要下载后才判断 | 防重复内容 |
| 指纹 (SimHash) | 检测相似内容 | 计算复杂 | 去重相似页面 |
| 自定义 Key | 灵活组合 | 需设计组合规则 | 业务定制 |
import hashlib
import re
from urllib.parse import urlparse, urlunparse
def normalize_url(url: str) -> str:
"""URL 归一化处理"""
parsed = urlparse(url)
# 移除 fragment
parsed = parsed._replace(fragment='')
# 移除尾随斜杠(路径部分)
path = parsed.path.rstrip('/') or '/'
parsed = parsed._replace(path=path)
# 排序查询参数
if parsed.query:
params = sorted(parsed.query.split('&'))
parsed = parsed._replace(query='&'.join(params))
# 移除默认端口
if parsed.port == 80 and parsed.scheme == 'http':
parsed = parsed._replace(netloc=parsed.hostname)
elif parsed.port == 443 and parsed.scheme == 'https':
parsed = parsed._replace(netloc=parsed.hostname)
return urlunparse(parsed)
def content_fingerprint(html: str) -> str:
"""计算内容指纹(去除 HTML 标签和空白后取 MD5)"""
clean = re.sub(r'<[^>]+>', '', html)
clean = re.sub(r'\s+', '', clean)
# 只取前 N 个字符计算指纹,提高性能
return hashlib.md5(clean[:5000].encode()).hexdigest()3.4 数据采集架构
+-------------+ +-----------+ +-----------+ +-----------+
| 爬虫节点 | --> | 消息队列 | --> | 清洗模块 | --> | 存储 |
| (Worker池) | | (Kafka/ | | (ETL) | | (DB/OSS/ |
| | | RabbitMQ) | | | | ES) |
+-------------+ +-----------+ +-----------+ +-----------+
| |
| 监控 | 数据质量
v v
+-------------+ +-----------+
| 监控系统 | | 数据校验 |
| (Prometheus | | (Schema |
| + Grafana) | | 校验器) |
+-------------+ +-----------+数据 Pipeline 设计
import json
import hashlib
from datetime import datetime
from typing import Generator, Dict, Any
from kafka import KafkaProducer, KafkaConsumer
class CrawlerPipeline:
"""完整的数据采集 Pipeline"""
def __init__(self, kafka_bootstrap: str = 'localhost:9092'):
self.producer = KafkaProducer(
bootstrap_servers=kafka_bootstrap,
value_serializer=lambda v: json.dumps(v).encode(),
acks='all',
retries=3
)
self.topic_raw = 'crawler:raw'
self.topic_clean = 'crawler:clean'
def push_raw_data(self, data: dict):
"""爬虫节点产出原始数据"""
self.producer.send(self.topic_raw, value=data)
def clean_data(self, data: dict) -> dict:
"""数据清洗"""
# 字段标准化
if 'timestamp' not in data:
data['timestamp'] = datetime.now().isoformat()
# 去除 HTML 标签
import re
if 'content' in data:
data['text'] = re.sub(r'<[^>]+>', '', data['content'])
data['text'] = data['text'].strip()
# 计算唯一标识
id_source = f"{data.get('url', '')}:{data.get('title', '')}"
data['_id'] = hashlib.md5(id_source.encode()).hexdigest()
return data
def push_clean_data(self, data: dict):
"""清洗后数据入队列"""
clean = self.clean_data(data)
self.producer.send(self.topic_clean, value=clean)
def consume_and_store(self, es_host: str = 'localhost:9200'):
"""消费清洗后的数据并存储到 Elasticsearch"""
from elasticsearch import Elasticsearch
es = Elasticsearch([es_host])
consumer = KafkaConsumer(
self.topic_clean,
bootstrap_servers='localhost:9092',
value_deserializer=lambda v: json.loads(v.decode()),
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='crawler-storage-group'
)
for msg in consumer:
data = msg.value
doc_id = data.get('_id', hashlib.md5(
data.get('url', '').encode()
).hexdigest())
# 检查是否已存在
if not es.exists(index='crawled_data', id=doc_id):
es.index(index='crawled_data', id=doc_id, body=data)
print(f'Stored: {doc_id} - {data.get("url", "")}')
# ---------- 监控统计 ----------
def get_pipeline_stats(self):
"""获取 Pipeline 统计信息"""
return {
'total_push': 0, # 通过计数器实现
'total_clean': 0,
'total_stored': 0,
'error_rate': 0.0,
'processing_latency_ms': 0
}4 分布式爬虫框架
4.1 Scrapy-Redis 架构
Scrapy-Redis 是 Scrapy 的分布式扩展,通过 Redis 共享请求队列和去重集合,实现多台机器的协同爬取。
核心组件
+-----------+
| Redis |
| (中央队列) |
+-----+-----+
|
+-----------------+-----------------+
| | |
+-----+-----+ +-----+-----+ +-----+-----+
| Spider 1 | | Spider 2 | | Spider N |
| (Worker 1) | | (Worker 2) | | (Worker N) |
+-----------+ +-----------+ +-----------+| 组件 | 类名 | 说明 |
|---|---|---|
| Spider | RedisSpider | 从 Redis 队列读取起始 URL |
| CrawlSpider | RedisCrawlSpider | 支持自动发现链接 |
| 调度器 | RedisScheduler | 通过 Redis 管理请求队列 |
| 去重 | RFPDupeFilter | 基于 Redis 的请求指纹去重 |
| 请求队列 | SpiderQueue | FIFO/LIFO/Priority 队列 |
| 数据管道 | 自定义 | 可写回 Redis 或其他存储 |
配置示例
# settings.py - Scrapy-Redis 分布式配置
# 调度器
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 去重过滤器
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# 请求队列(可选:SpiderQueue / SpiderStack / SpiderPriorityQueue)
SCHEDULER_QUEUE_CLASS = "scrapy_redis.queue.SpiderQueue"
# Redis 连接
REDIS_HOST = '127.0.0.1'
REDIS_PORT = 6379
REDIS_PARAMS = {
'db': 0,
'password': None,
'decode_responses': True, # 自动解码
}
# 持久化调度器状态(爬虫停止后不清空队列)
SCHEDULER_PERSIST = True
# 可选:不清理去重集合(方便断点续爬)
DUPEFILTER_PERSIST = True
# 并发设置
CONCURRENT_REQUESTS = 16
CONCURRENT_REQUESTS_PER_DOMAIN = 8
# 下载延迟
DOWNLOAD_DELAY = 0.5
RANDOMIZE_DOWNLOAD_DELAY = True
# 禁用 Cookies(根据目标网站需求)
COOKIES_ENABLED = False
# 下载中间件
DOWNLOADER_MIDDLEWARES = {
'scrapy.downloadermiddlewares.useragent.UserAgentMiddleware': None,
'scrapy.downloadermiddlewares.retry.RetryMiddleware': 90,
'scrapy.downloadermiddlewares.httpproxy.HttpProxyMiddleware': 110,
}
# 启用 Redis 管道(将数据写入 Redis)
ITEM_PIPELINES = {
'scrapy_redis.pipelines.RedisPipeline': 300,
}Spider 示例
# spiders/example_spider.py
from scrapy_redis.spiders import RedisSpider
class ExampleSpider(RedisSpider):
"""继承 RedisSpider,从 Redis 获取起始 URL"""
name = 'example'
redis_key = 'example:start_urls' # Redis 中的起始 URL 列表键
def __init__(self, *args, **kwargs):
# 动态设置允许域名
domain = kwargs.pop('domain', '')
self.allowed_domains = [domain] if domain else []
super().__init__(*args, **kwargs)
def parse(self, response):
# 提取数据
yield {
'url': response.url,
'title': response.css('title::text').get(),
'content_length': len(response.text),
}
# 提取并跟进链接
for href in response.css('a::attr(href)'):
yield response.follow(href, self.parse)
# 启动后,向 Redis 注入起始 URL:
# $ redis-cli lpush example:start_urls "https://example.com"Spider 共享机制
多台机器上的 Spider 实例共享同一个 Redis 队列和去重集合,天然实现负载均衡:
# 配置文件:不同机器使用相同的配置即可组成集群
# 机器 A
# $ scrapy crawl example
# 机器 B
# $ scrapy crawl example
# 机器 C
# $ scrapy crawl example
# 向 Redis 注入起始 URL,三台机器会协同爬取
# $ for url in $(cat urls.txt); do
# redis-cli lpush example:start_urls "$url"
# done4.2 Scrapy-Redis 配置详解
Scheduler 配置
| 配置项 | 说明 | 可选值 |
|---|---|---|
SCHEDULER | 调度器类 | scrapy_redis.scheduler.Scheduler |
SCHEDULER_PERSIST | 是否持久化调度器状态 | True / False |
SCHEDULER_FLUSH_ON_START | 启动时是否清空队列 | True / False |
SCHEDULER_IDLE_BEFORE_CLOSE | 空闲等待秒数后关闭 | 数值(默认 0) |
SCHEDULER_QUEUE_CLASS | 请求队列类型 | 见下方 |
请求队列类型
| 队列类 | 策略 | 适用场景 |
|---|---|---|
SpiderQueue | FIFO(先进先出) | 广度优先爬取 |
SpiderStack | LIFO(后进先出) | 深度优先爬取 |
SpiderPriorityQueue | 优先级队列 | 按优先级调度 |
# 设置优先级队列
SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.SpiderPriorityQueue'
# 在 Spider 中设置请求优先级
def parse(self, response):
# 高优先级
yield scrapy.Request(
'https://example.com/important',
callback=self.parse_important,
priority=10 # 数值越大优先级越高
)DupeFilter 配置
# 去重过滤器
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# 持久化去重集合(断点续爬)
DUPEFILTER_PERSIST = True
# 可选:自定义去重键函数
# 默认使用请求指纹(method + url + body + headers)
# 可通过继承 RFPDupeFilter 自定义4.3 PySpider
架构组件
+------------------+
| WebUI | <-- 浏览器访问,管理爬虫
| (Flask + Admin) |
+--------+---------+
|
+--------+---------+
| Scheduler | <-- 任务调度,URL 管理
| (任务队列 + |
| 去重 + Rate) |
+--------+---------+
|
+--------+---------+
| Fetcher | <-- 异步抓取,支持 PhantomJS
| (Tornado + |
| PyQuery) |
+--------+---------+
|
+--------+---------+
| Processor | <-- 执行用户的脚本逻辑
| (Python 脚本) |
+------------------+各组件说明
| 组件 | 技术栈 | 功能 |
|---|---|---|
| WebUI | Flask | 项目管理、代码编辑、任务监控、结果查看 |
| Scheduler | 自实现 | URL 去重、任务优先级、定时调度、Rate 控制 |
| Fetcher | Tornado | 异步 HTTP 抓取、支持 PhantomJS/Chrome |
| Processor | Python | 执行用户定义的爬虫脚本(on_start / @every / @config) |
爬虫脚本示例
#!/usr/bin/env python
# -*- coding: utf-8 -*-
from pyspider.libs.base_handler import *
class Handler(BaseHandler):
"""PySpider 爬虫脚本"""
@every(minutes=30) # 每 30 分钟执行一次
def on_start(self):
"""爬虫入口"""
self.crawl(
'https://example.com/list/1',
callback=self.parse_list,
validate_cert=False
)
@config(age=10 * 24 * 60 * 60) # 10 天内不重复抓取
def parse_list(self, response):
"""列表页解析"""
for item in response.doc('.item a').items():
self.crawl(
item.attr.href,
callback=self.parse_detail,
priority=10 # 优先级
)
# 翻页
next_page = response.doc('.next').attr.href
if next_page:
self.crawl(
next_page,
callback=self.parse_list
)
def parse_detail(self, response):
"""详情页解析"""
return {
'url': response.url,
'title': response.doc('h1').text(),
'content': response.doc('.content').text(),
'publish_time': response.doc('.time').text(),
}
@config(priority=5)
def parse_error_page(self, response):
"""错误页面处理"""
if response.status_code != 200:
# 记录失败
self.send_message(self.project_name, {
'url': response.url,
'status': response.status_code
})框架对比
| 特性 | Scrapy-Redis | PySpider | 自研框架 |
|---|---|---|---|
| 部署复杂度 | 中 | 低 | 高 |
| WebUI | 无 | 有 | 可自研 |
| 分布式 | 天然支持 | 需单独配置 | 完全可控 |
| 去重 | Redis/BloomFilter | 自带 | 自定义 |
| 调度策略 | 队列维度 | 任务维度 | 灵活 |
| JS 渲染 | 需集成 Splash | 内置 PhantomJS | 自集成 |
| 扩展性 | 通过 Middleware | 通过 Handler | 完全自由 |
| 社区生态 | 丰富 | 一般 | 取决于开发 |
| 适用场景 | 大规模分布式 | 中规模/可视化 | 高度定制 |
适用场景建议
- Scrapy-Redis:适合数据量大、需要多机分布式采集的场景,如全网新闻采集、电商商品数据采集
- PySpider:适合需要可视化管理和快速开发的中小规模项目,脚本热加载方便调试
- 自研框架:适合有特殊需求(定制反反爬策略、复杂调度逻辑)的大型项目
5 爬虫监控与管理
5.1 日志采集
结构化日志
结构化日志以 JSON 格式记录,便于后续采集和分析:
import json
import logging
from datetime import datetime
class StructuredLogger:
"""结构化日志记录器"""
def __init__(self, name: str):
self.logger = logging.getLogger(name)
self.logger.setLevel(logging.INFO)
def log_request(self, url: str, status: int, latency_ms: float,
proxy: str = '', error: str = ''):
"""记录请求日志"""
record = {
'timestamp': datetime.now().isoformat(),
'type': 'request',
'url': url,
'status': status,
'latency_ms': latency_ms,
'proxy': proxy,
'error': error
}
self.logger.info(json.dumps(record))
def log_error(self, task_id: str, error_msg: str, stack_trace: str = ''):
"""记录错误日志"""
record = {
'timestamp': datetime.now().isoformat(),
'type': 'error',
'task_id': task_id,
'error': error_msg,
'stack_trace': stack_trace
}
self.logger.error(json.dumps(record))
def log_stats(self, stats: dict):
"""记录统计日志"""
record = {
'timestamp': datetime.now().isoformat(),
'type': 'stats',
**stats
}
self.logger.info(json.dumps(record))日志分类与采集
| 日志类型 | 内容 | 级别 | 采集方式 |
|---|---|---|---|
| 请求日志 | URL、状态码、延迟、代理 IP | INFO | Filebeat -> Logstash -> ES |
| 错误日志 | 异常堆栈、失败原因 | ERROR | Filebeat -> Logstash -> ES |
| 数据量统计 | 采集数量、存储量 | INFO | 定时上报 -> Prometheus |
| 成功率统计 | 成功/失败比例 | INFO | 定时上报 -> Prometheus |
Elasticsearch 日志存储
from elasticsearch import Elasticsearch
from elasticsearch.helpers import bulk
class LogStorer:
"""日志存储到 Elasticsearch"""
def __init__(self, hosts: list = ['localhost:9200']):
self.es = Elasticsearch(hosts)
def index_log(self, index: str, doc: dict):
"""单条写入"""
self.es.index(
index=index,
body=doc,
pipeline='crawler_logs'
)
def bulk_index(self, index: str, docs: list):
"""批量写入"""
actions = [
{
'_index': index,
'_source': doc
}
for doc in docs
]
success, _ = bulk(self.es, actions)
return success
def search_logs(self, index: str, query: dict, size: int = 50):
"""搜索日志"""
return self.es.search(
index=index,
body=query,
size=size
)Kibana 可视化配置
推荐创建以下可视化面板:
- 请求成功率趋势图:折线图,展示单位时间内的请求成功率
- 延迟分布直方图:展示请求延迟的分布情况
- 代理 IP 地理分布:地图,展示代理 IP 的来源地区
- 错误类型饼图:展示各类错误的占比
- TOP 10 错误 URL:展示出错最多的 URL
- 爬取数据量趋势:展示单位时间内的数据产出量
5.2 告警系统
告警规则
import json
import smtplib
import requests
from datetime import datetime, timedelta
class AlertManager:
"""告警管理器"""
def __init__(self, redis_client):
self.redis = redis_client
self.webhook_url = '' # 企业微信/钉钉/Slack Webhook
self.thresholds = {
'success_rate': 0.95, # 成功率低于 95% 触发告警
'proxy_pool_min': 50, # 代理池少于 50 个触发告警
'memory_threshold': 85, # 内存使用超过 85% 触发告警
'captcha_frequency': 0.1, # 验证码出现频率超过 10% 触发告警
}
def check_alerts(self, stats: dict) -> list:
"""检查所有告警规则"""
alerts = []
# 成功率下降告警
if stats.get('success_rate', 1.0) < self.thresholds['success_rate']:
alerts.append({
'level': 'critical',
'type': 'success_rate_drop',
'message': f"请求成功率降至 {stats['success_rate']:.1%}",
'current': stats['success_rate'],
'threshold': self.thresholds['success_rate']
})
# 代理池枯竭告警
if stats.get('proxy_pool_size', 0) < self.thresholds['proxy_pool_min']:
alerts.append({
'level': 'warning',
'type': 'proxy_pool_exhaustion',
'message': f"代理池仅剩 {stats['proxy_pool_size']} 个可用代理",
'current': stats['proxy_pool_size'],
'threshold': self.thresholds['proxy_pool_min']
})
# 资源耗尽告警
if stats.get('memory_usage', 0) > self.thresholds['memory_threshold']:
alerts.append({
'level': 'critical',
'type': 'resource_exhaustion',
'message': f"内存使用率达到 {stats['memory_usage']}%",
'current': stats['memory_usage'],
'threshold': self.thresholds['memory_threshold']
})
# 验证码频率激增告警
if stats.get('captcha_rate', 0) > self.thresholds['captcha_frequency']:
alerts.append({
'level': 'warning',
'type': 'captcha_spike',
'message': f"验证码出现频率 {stats['captcha_rate']:.1%}",
'current': stats['captcha_rate'],
'threshold': self.thresholds['captcha_frequency']
})
return alerts
def send_alert(self, alert: dict):
"""发送告警通知"""
message = f"[{alert['level'].upper()}] {alert['message']}"
# 发送到 Webhook(企业微信/钉钉/Slack)
if self.webhook_url:
payload = {
'msgtype': 'text',
'text': {'content': message}
}
try:
requests.post(self.webhook_url, json=payload, timeout=5)
except Exception as e:
print(f'Webhook 发送失败: {e}')
# 记录告警到 Redis
self.redis.lpush('alerts:history', json.dumps({
**alert,
'timestamp': datetime.now().isoformat()
}))
self.redis.ltrim('alerts:history', 0, 999) # 保留最近 1000 条
def evaluate_and_alert(self, stats: dict):
"""评估指标并发送告警"""
alerts = self.check_alerts(stats)
for alert in alerts:
self.send_alert(alert)
return alerts告警通知渠道对比
| 渠道 | 实时性 | 支持格式 | 适用场景 |
|---|---|---|---|
| 企业微信机器人 | 高 | Text/Markdown | 团队内部通知 |
| 钉钉机器人 | 高 | Text/Markdown | 团队内部通知 |
| Slack Webhook | 高 | Text/Attachment | 国际化团队 |
| 邮件 | 低 | HTML | 非紧急报告 |
| 短信 | 高 | 纯文本 | 紧急告警 |
| 电话 | 最高 | 语音 | P0 级故障 |
5.3 爬虫管理平台
平台功能对比
| 功能 | Crawlab | SpiderFlow | 自研管理后台 |
|---|---|---|---|
| 任务部署 | 支持 | 支持 | 自定义 |
| 定时调度 | 支持(Cron) | 支持(Cron) | 自定义 |
| 节点管理 | 支持 | 不支持 | 自定义 |
| 运行日志 | 实时查看 | 实时查看 | 自定义 |
| 爬虫管理 | 通用 | 通用 | 业务定制 |
| 可视化编排 | 不支持 | 支持 | 可选 |
| 权限管理 | 基础 | 基础 | 可深度定制 |
| 数据导出 | 支持 | 支持 | 自定义 |
| 安装复杂度 | 中(Docker) | 低(Java 单体) | 高 |
| 扩展性 | 中 | 低 | 高 |
Crawlab 部署配置
# docker-compose.yml
version: '3.3'
services:
master:
image: tikazyq/crawlab:latest
container_name: crawlab_master
environment:
CRAWLAB_SERVER_MASTER: "Y"
CRAWLAB_MONGO_HOST: "mongo"
CRAWLAB_REDIS_ADDRESS: "redis"
ports:
- "8080:8080"
depends_on:
- mongo
- redis
worker:
image: tikazyq/crawlab:latest
container_name: crawlab_worker
environment:
CRAWLAB_SERVER_MASTER: "N"
CRAWLAB_MONGO_HOST: "mongo"
CRAWLAB_REDIS_ADDRESS: "redis"
CRAWLAB_SERVER_REGISTER_IP: "<worker_ip>"
depends_on:
- mongo
- redis
mongo:
image: mongo:5.0
restart: always
redis:
image: redis:7.0
restart: always自研管理后台接口设计
# api/scheduler_api.py - 调度管理 API
from flask import Flask, request, jsonify
import redis
import json
from datetime import datetime
app = Flask(__name__)
redis_client = redis.Redis(host='localhost', port=6379, decode_responses=True)
# ---------- 任务管理 ----------
@app.route('/api/tasks', methods=['POST'])
def create_task():
"""创建爬虫任务"""
data = request.json
task = {
'id': str(uuid.uuid4()),
'name': data['name'],
'spider': data['spider'],
'config': data.get('config', {}),
'schedule': data.get('schedule', ''), # Cron 表达式
'status': 'pending',
'created_at': datetime.now().isoformat(),
}
redis_client.hset('tasks', task['id'], json.dumps(task))
return jsonify(task), 201
@app.route('/api/tasks/<task_id>/start', methods=['POST'])
def start_task(task_id: str):
"""启动任务"""
task_data = redis_client.hget('tasks', task_id)
if not task_data:
return jsonify({'error': 'task not found'}), 404
task = json.loads(task_data)
task['status'] = 'running'
task['started_at'] = datetime.now().isoformat()
redis_client.hset('tasks', task_id, json.dumps(task))
# 将任务注入爬虫队列
redis_client.lpush('crawler:start_urls', json.dumps({
'task_id': task_id,
'spider': task['spider'],
'config': task['config']
}))
return jsonify(task)
@app.route('/api/tasks/<task_id>/stop', methods=['POST'])
def stop_task(task_id: str):
"""停止任务"""
task_data = redis_client.hget('tasks', task_id)
if not task_data:
return jsonify({'error': 'task not found'}), 404
task = json.loads(task_data)
task['status'] = 'stopped'
task['stopped_at'] = datetime.now().isoformat()
redis_client.hset('tasks', task_id, json.dumps(task))
return jsonify(task)
# ---------- 节点管理 ----------
@app.route('/api/nodes', methods=['GET'])
def list_nodes():
"""列出所有爬虫节点"""
workers = []
for key in redis_client.scan_iter('worker:*:hb'):
worker_id = key.split(':')[1]
data = redis_client.get(key)
if data:
worker_info = json.loads(data)
workers.append({
'id': worker_id,
'load': worker_info['load'],
'last_seen': worker_info['timestamp']
})
return jsonify(workers)
# ---------- 统计仪表盘 ----------
@app.route('/api/stats', methods=['GET'])
def get_stats():
"""获取爬虫统计信息"""
return jsonify({
'total_tasks': redis_client.hlen('tasks'),
'completed_tasks': redis_client.get('tasks:completed') or 0,
'failed_tasks': redis_client.get('tasks:failed') or 0,
'active_nodes': len(list(redis_client.scan_iter('worker:*:hb'))),
'proxy_pool_size': redis_client.scard('proxy:pool'),
'queue_size': redis_client.llen('crawler:start_urls'),
})
if __name__ == '__main__':
app.run(host='0.0.0.0', port=8000)Cron 定时调度
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger
import redis
import json
class TaskScheduler:
"""基于 APScheduler 的定时任务调度器"""
def __init__(self, redis_client: redis.Redis):
self.redis = redis_client
self.scheduler = BackgroundScheduler()
def load_tasks(self):
"""加载所有定时任务"""
tasks = self.redis.hgetall('tasks')
for task_id, task_data in tasks.items():
task = json.loads(task_data)
if task.get('schedule'):
self.scheduler.add_job(
func=self.execute_task,
trigger=CronTrigger.from_crontab(task['schedule']),
args=[task_id],
id=f'task_{task_id}',
replace_existing=True
)
def execute_task(self, task_id: str):
"""执行定时任务"""
task_data = self.redis.hget('tasks', task_id)
if not task_data:
return
task = json.loads(task_data)
task['status'] = 'running'
task['started_at'] = datetime.now().isoformat()
self.redis.hset('tasks', task_id, json.dumps(task))
# 注入爬虫队列
self.redis.lpush('crawler:start_urls', json.dumps({
'task_id': task_id,
'spider': task['spider'],
'config': task['config']
}))
def start(self):
"""启动调度器"""
self.load_tasks()
self.scheduler.start()