我在三家公司当过"那个写爬虫的人"。每次别人说"帮忙搭个爬虫",真实需求其实都不一样。A 组要电商比价,B 组要新闻聚合,C 组要在登录态里抓订单。如果你按需求一个一个写脚本,最后会攒下十几个几乎一样的文件,网站一改版同时崩三个。
我后来想明白一件事:问题不在爬虫写得够不够好,在于系统把"采集逻辑"和"业务行为"焊死在了一起。解决办法是反过来,内核永远不动,行为全靠插进去的零件决定。这篇文章聊的就是这套"中间件"思路,以及它怎么让系统真的能扩展。
一、先说清楚,"通用"难在哪
很多人以为通用就是"多写几个 if"。错了。if 堆多了以后,改一个分支要读懂整个文件,这比维护十个脚本还累。
通用的真正含义是三件事:
第一,加新能力不碰核心代码。今天要加 UA 随机,明天要加代理,后天要加指纹伪装,内核一行都不用改。
第二,零件能单独换。代理从自建池换成隧道代理,只改一个类,其余全部不动。
第三,能横向铺机器。中间件本身无状态,十台机器跑同一套代码,行为一致。
这三点里,第一点最关键,也最常被忽略。
二、中间件到底是什么
中间件说白了就是一个函数,卡在请求和响应的必经之路上,可以对它们动手脚。比如请求要发出去之前,塞个代理地址进去;响应回来之后,遇到 403 就标记重试。
按作用位置分三类:
- 下载中间件:围着 HTTP 请求转,代理、UA、限速都在这层。
- 解析中间件:围着解析出来的数据转,清洗、字段映射在这层。
- 调度中间件:围着任务队列转,优先级、去重在这层。
我这次重点讲下载中间件,因为它最影响成功率,也最容易写出扩展性问题。
三、整体架构:一条带插口的管道
把系统画成一条从左到右的管道,每个接缝处都能插中间件:
┌─────────┐ ┌──────────────┐ ┌──────────┐ ┌──────────┐
│ 调度层 │───→│ 下载中间件链 │───→│ 下载器 │───→│ 响应中间件链 │───→ 数据
│ 队列/去重 │ │ 代理/UA/限速 │ │ aiohttp │ │ 重试/落库 │
└─────────┘ └──────────────┘ └──────────┘ └──────────┘
代理放在下载中间件链的最前面,因为它必须在请求真正发出之前就把代理地址和 IP 控制头挂上去。顺序错了,代理不生效,后面全白搭。
四、可插拔的三条规矩
契约先行。 所有中间件继承同一个基类,只实现 process_request 和 process_response 两个方法。靠这个统一契约,管理器才能不关心中间件具体干嘛,只管按顺序调用。
顺序可控。 每个中间件带一个 order 数字,越小越先跑(请求方向)。代理设 10,UA 设 20,重试设 200。顺序靠数字说话,不靠文件里的排列,改顺序不用动代码位置。
故障隔离。 单个中间件抛异常,不能把整条链带崩。管理器用 try/except 包住每次调用,中间件出错就记一条日志跳过,请求照常往后走。我见过有人把限速中间件写崩了,结果整批任务卡死,就是因为没做这层隔离。
五、代码:一个最小可用的中间件框架
下面这套代码能直接跑(依赖 aiohttp)。核心是 Middleware 基类、MiddlewareManager 调度器,和一个基于16YUN隧道代理的 YiniuProxyMiddleware。把代理地址换成你自己的账号密码就能用。
"""
通用爬虫中间件框架(最小可用版)
核心思想:把"代理、UA、重试"等能力拆成独立中间件,
通过统一契约插进请求/响应生命周期,新增能力不碰核心代码。
代理层基于16YUN隧道代理,云端自动调度出口IP。
"""
import asyncio
import abc
import logging
import random
from dataclasses import dataclass, field
from typing import Optional, Dict, Any, List
import aiohttp
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger(__name__)
@dataclass
class Request:
"""一次抓取请求,中间件在它身上挂代理、加头、记元数据"""
url: str
headers: Dict[str, str] = field(default_factory=dict)
proxy: Optional[str] = None # 由代理中间件填充
meta: Dict[str, Any] = field(default_factory=dict)
retry_count: int = 0
@dataclass
class Response:
url: str
status: int
text: str
request: Request
class Middleware(abc.ABC):
"""所有中间件的抽象基类。
设计约定:
- order 越小越先执行(请求方向),响应方向反向执行
- process_request 返回 Request 表示放行,返回 None 表示拦截丢弃
- process_response 返回 Response 表示继续往后走
"""
order: int = 100
async def process_request(self, request: Request) -> Optional[Request]:
return request
async def process_response(self, response: Response) -> Optional[Response]:
return response
class YiniuProxyMiddleware(Middleware):
"""隧道代理中间件。
隧道代理的特点:客户端只配置一个固定入口,云端在数十万IP池里
自动调度出口IP,不用自己维护IP池、不用写轮换和失效剔除逻辑。
本中间件只负责把代理地址和IP控制头挂到请求上。
"""
order = 10 # 必须在真正发起请求之前执行
def __init__(self, username: str, password: str,
host: str = "t.16yun.cn", port: int = 31111,
ip_mode: str = "force_switch", tunnel_id: str = "crawler_001"):
# 固定的隧道入口,账号密码拼进URL即可,无需管理IP列表
self.proxy_url = f"http://{username}:{password}@{host}:{port}"
self.ip_mode = ip_mode # 三种模式:force_switch / keep_alive / tunnel_fixed
self.tunnel_id = tunnel_id
def _build_headers(self) -> Dict[str, str]:
"""根据IP控制模式构造请求头(对应隧道代理的几种模式)"""
headers = {
}
if self.ip_mode == "force_switch":
# 模式A:每次断开TCP,云端分配新出口IP,适合匿名批量采集
headers["Connection"] = "close"
elif self.ip_mode == "keep_alive":
# 模式B:复用TCP,同一会话出口IP不变,适合登录态采集
headers["Connection"] = "keep-alive"
elif self.ip_mode == "tunnel_fixed":
# 模式C/D:相同tunnel_id路由到同一出口IP,实现多业务IP隔离
# 注意:HTTPS 目标站点需改用 aiohttp 的 proxy_headers 传该头
headers["Proxy-Tunnel"] = self.tunnel_id
return headers
async def process_request(self, request: Request) -> Optional[Request]:
# 挂上代理地址和IP控制头,后续下载器直接用
request.proxy = self.proxy_url
request.headers.update(self._build_headers())
return request
class UserAgentMiddleware(Middleware):
"""UA随机中间件,order比代理靠后,两者不冲突"""
order = 20
UA_POOL = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/129.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Safari/605.1.15",
"Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/128.0 Safari/537.36",
]
async def process_request(self, request: Request) -> Optional[Request]:
request.headers.setdefault("User-Agent", random.choice(self.UA_POOL))
return request
class RetryMiddleware(Middleware):
"""重试中间件,放在响应链,遇到可重试状态码就标记重抓"""
order = 200
def __init__(self, max_retry: int = 3):
self.max_retry = max_retry
async def process_response(self, response: Response) -> Optional[Response]:
# 403/429/5xx 视为可重试:累加计数并打标记,交给调度层重入队
if response.status in (403, 429, 500, 502, 503):
response.request.retry_count += 1
if response.request.retry_count <= self.max_retry:
response.request.meta["need_retry"] = True
logger.warning(f"[{response.status}] 标记重试: {response.url} "
f"(第{response.request.retry_count}次)")
return response
class MiddlewareManager:
"""中间件调度器:按order排序,依次跑请求链和响应链"""
def __init__(self):
self._middlewares: List[Middleware] = []
def add(self, mw: Middleware) -> None:
"""注册一个中间件。新增能力 = 写一个类 + 调一次add,核心代码不动"""
self._middlewares.append(mw)
def _sorted(self, reverse: bool = False) -> List[Middleware]:
return sorted(self._middlewares, key=lambda m: m.order, reverse=reverse)
async def process_request(self, request: Request) -> Optional[Request]:
for mw in self._sorted():
try:
result = await mw.process_request(request)
if result is None:
return None # 被中间件拦截,不再往下走
request = result
except Exception as e:
# 故障隔离:单个中间件出错不影响整条链
logger.warning(f"中间件 {mw.__class__.__name__} 处理请求异常: {e}")
return request
async def process_response(self, response: Response) -> Optional[Response]:
for mw in self._sorted(reverse=True): # 响应方向反向执行
try:
result = await mw.process_response(response)
if result is None:
return None
response = result
except Exception as e:
logger.warning(f"中间件 {mw.__class__.__name__} 处理响应异常: {e}")
return response
class Downloader:
"""下载器:持有中间件管理器,负责真正发请求"""
def __init__(self, manager: MiddlewareManager,
concurrency: int = 20, force_close: bool = True):
self.manager = manager
self.concurrency = concurrency
self.semaphore = asyncio.Semaphore(concurrency)
# force_close 由代理中间件的IP模式决定(模式A需要True来触发换IP)
self.force_close = force_close
async def _fetch_one(self, session: aiohttp.ClientSession, url: str):
request = Request(url=url)
async with self.semaphore:
# 1) 先过请求中间件链:挂代理、加UA
request = await self.manager.process_request(request)
if request is None:
return None
try:
# 2) 真正发请求,走隧道代理
async with session.get(
request.url,
headers=request.headers,
proxy=request.proxy,
timeout=aiohttp.ClientTimeout(total=15),
) as resp:
text = await resp.text()
response = Response(request.url, resp.status, text, request)
except Exception as e:
logger.warning(f"下载失败 {url}: {e}")
return None
# 3) 再过响应中间件链:重试判定、落库前清洗
return await self.manager.process_response(response)
async def run(self, urls: List[str]) -> List:
# 连接池的force_close与IP模式联动:模式A每次断连换IP,其余复用
connector = aiohttp.TCPConnector(
limit=self.concurrency,
force_close=self.force_close,
enable_cleanup_closed=True,
)
async with aiohttp.ClientSession(connector=connector) as session:
tasks = [self._fetch_one(session, u) for u in urls]
return await asyncio.gather(*tasks, return_exceptions=True)
async def main():
# 组装系统:注册三个中间件,顺序由各自order决定
manager = MiddlewareManager()
manager.add(YiniuProxyMiddleware(
username="your-username", password="your-password",
ip_mode="force_switch", # 匿名批量采集用模式A
))
manager.add(UserAgentMiddleware())
manager.add(RetryMiddleware(max_retry=3))
downloader = Downloader(manager, concurrency=20, force_close=True)
urls = [f"https://httpbin.org/ip" for _ in range(10)]
results = await downloader.run(urls)
ok = sum(1 for r in results if isinstance(r, Response) and r.status == 200)
logger.info(f"完成: 成功 {ok}/{len(urls)}")
if __name__ == "__main__":
asyncio.run(main())
这段代码不长,但把"通用"和"可扩展"都落到了实处:加一个新行为,你只需要写一个 Middleware 子类,然后 manager.add(...) 一行注册,核心的 Downloader 和 MiddlewareManager 一行都不用改。
六、隧道代理为什么适合做中间件
我试过自建代理池,也试过几家服务商,最后稳定用隧道代理,原因跟"中间件"这个主题其实是一回事:它把 IP 管理的复杂度从我的代码里挪走了。
隧道代理是一个固定入口(t.16yun.cn:31111),云端在几十万 IP 池里自动调度出口,失效节点毫秒级剔除。落到我的框架里,YiniuProxyMiddleware 这个类只干一件事:把代理地址和 Connection / Proxy-Tunnel 头挂上去。IP 轮换、健康检查、失效剔除,全在云端完成,我没写一行相关代码。
这和中间件的"故障隔离、独立替换"思路完全对得上。哪天要换 API 代理或者独享代理,我只改 YiniuProxyMiddleware 一个类,甚至换成另一个中间件类,下载器和重试逻辑完全不动。
它家几种 IP 控制模式也能直接映射到中间件配置:
- 匿名批量采集(商品列表、新闻聚合)用模式 A,每请求换 IP;
- 登录态采集用模式 B,同一会话保 IP;
- 多站点业务隔离用模式 C/D,不同
tunnel_id各自的出口互不干扰。
我一个新闻分析项目同时跑 5 个新闻站点,最早 5 个站点混在同一个隧道里,被其中一个站的反爬牵连,整体成功率掉到 72%。改成每个站点一个 tunnel_id 之后,IP 行为彻底隔离,成功率回到 94%。这件事如果用传统"一个代理池"的思路做,得维护 5 个池子,换成中间件就是一个字符串的事。
七、高扩展性怎么真正落地
框架写完只是开始。要做到真能扩展,还有几件工程上的事:
无状态。 中间件不存业务状态,十台机器跑同一套代码行为一致。需要登录态的那种中间件是例外,得把它单独拆出来,别让它拖累整批任务的横向扩展。
插件化加载。 把中间件类放到一个目录,启动时按配置文件动态 import。加功能 = 丢一个文件 + 配置里加一行,不用改核心、不用重新发版。我们生产环境现在是这么干的,新同事加一个反爬中间件从提交流程到上线不到半小时。
指标外抛。 每个中间件在 process_request / process_response 里埋点:耗时、命中次数、拦截次数。没有这些数字,你不知道是代理慢还是解析慢,调优全靠猜。我吃过这个亏,有次整体成功率掉了一半,盯了两天才发现是某个新加的中间件把超时设错了。
非阻塞。 中间件里别写同步重活(比如大循环、慢 IO)。它夹在请求链里,一个中间件卡住,后面的请求全排队。需要重计算的,丢到后台任务或者独立服务。
八、几个我踩过的坑
中间件顺序写反是最常见也最隐蔽的。代理设成 order 999,结果请求先发出去了才挂代理,等于没挂。我现在约定:下载中间件全部用两位数 order,重试类用三位数,基本不会再乱。
第二个坑是状态中间件破坏扩展。早期我把登录态塞进一个全局中间件,结果这台机器登录了,那台机器没有,横向一铺就串号。后来把登录单独做成"账号会话池",按账号维度分配,才解决。
第三个坑和代理有关:模式 A 下忘了开 force_close,aiohttp 默认复用 TCP 连接,代理端看到同一个连接就不再分配新 IP,30 个并发实际上只走了 3 个出口,风控秒封。这个在代码里我用 force_close 跟 IP 模式联动了,避免再犯。
九、能带走的几点
- 通用系统的关键不是多写 if,是把"行为"和"引擎"拆开。中间件就是那道缝。
- 契约、顺序、故障隔离,这三件事定下来,扩展性就有了地基。
- 代理这类外部依赖,尽量做成可替换的中间件。隧道代理的价值不只是 IP 稳,更在于它把 IP 管理挪出了我的代码,让我能专心写业务。
- 框架搭完不等于能扩展。无状态、插件化加载、指标埋点、非阻塞,这几条不做,系统还是长不大。
如果今天只做一件事,我建议先把你现有爬虫里"代理、UA、重试"这三段抽成一个个中间件类,哪怕还没上分布式。这一步做完,你后面每一个新需求都会轻松很多。