架构师视角:如何设计一套高扩展性的通用爬虫中间件系统?

简介: 这篇文章讨论了如何设计一套高扩展性的通用爬虫中间件系统。作者从架构师视角出发,分析了爬虫系统的需求和挑战,并提出了相应的设计原则和解决方案。文章强调了模块化、可配置性和容错性的重要性,旨在帮助读者构建一个稳定、高效且易于扩展的爬虫系统。

我在三家公司当过"那个写爬虫的人"。每次别人说"帮忙搭个爬虫",真实需求其实都不一样。A 组要电商比价,B 组要新闻聚合,C 组要在登录态里抓订单。如果你按需求一个一个写脚本,最后会攒下十几个几乎一样的文件,网站一改版同时崩三个。

我后来想明白一件事:问题不在爬虫写得够不够好,在于系统把"采集逻辑"和"业务行为"焊死在了一起。解决办法是反过来,内核永远不动,行为全靠插进去的零件决定。这篇文章聊的就是这套"中间件"思路,以及它怎么让系统真的能扩展。

一、先说清楚,"通用"难在哪

很多人以为通用就是"多写几个 if"。错了。if 堆多了以后,改一个分支要读懂整个文件,这比维护十个脚本还累。

通用的真正含义是三件事:

第一,加新能力不碰核心代码。今天要加 UA 随机,明天要加代理,后天要加指纹伪装,内核一行都不用改。

第二,零件能单独换。代理从自建池换成隧道代理,只改一个类,其余全部不动。

第三,能横向铺机器。中间件本身无状态,十台机器跑同一套代码,行为一致。

这三点里,第一点最关键,也最常被忽略。

二、中间件到底是什么

中间件说白了就是一个函数,卡在请求和响应的必经之路上,可以对它们动手脚。比如请求要发出去之前,塞个代理地址进去;响应回来之后,遇到 403 就标记重试。

按作用位置分三类:

  • 下载中间件:围着 HTTP 请求转,代理、UA、限速都在这层。
  • 解析中间件:围着解析出来的数据转,清洗、字段映射在这层。
  • 调度中间件:围着任务队列转,优先级、去重在这层。

我这次重点讲下载中间件,因为它最影响成功率,也最容易写出扩展性问题。

三、整体架构:一条带插口的管道

把系统画成一条从左到右的管道,每个接缝处都能插中间件:

┌─────────┐    ┌──────────────┐    ┌──────────┐    ┌──────────┐
│ 调度层   │───→│  下载中间件链  │───→│  下载器   │───→│ 响应中间件链 │───→ 数据
│ 队列/去重 │    │ 代理/UA/限速   │    │ aiohttp  │    │ 重试/落库   │
└─────────┘    └──────────────┘    └──────────┘    └──────────┘

代理放在下载中间件链的最前面,因为它必须在请求真正发出之前就把代理地址和 IP 控制头挂上去。顺序错了,代理不生效,后面全白搭。

四、可插拔的三条规矩

契约先行。 所有中间件继承同一个基类,只实现 process_requestprocess_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(...) 一行注册,核心的 DownloaderMiddlewareManager 一行都不用改。

六、隧道代理为什么适合做中间件

我试过自建代理池,也试过几家服务商,最后稳定用隧道代理,原因跟"中间件"这个主题其实是一回事:它把 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 模式联动了,避免再犯。

九、能带走的几点

  1. 通用系统的关键不是多写 if,是把"行为"和"引擎"拆开。中间件就是那道缝。
  2. 契约、顺序、故障隔离,这三件事定下来,扩展性就有了地基。
  3. 代理这类外部依赖,尽量做成可替换的中间件。隧道代理的价值不只是 IP 稳,更在于它把 IP 管理挪出了我的代码,让我能专心写业务。
  4. 框架搭完不等于能扩展。无状态、插件化加载、指标埋点、非阻塞,这几条不做,系统还是长不大。

如果今天只做一件事,我建议先把你现有爬虫里"代理、UA、重试"这三段抽成一个个中间件类,哪怕还没上分布式。这一步做完,你后面每一个新需求都会轻松很多。

相关文章
|
8天前
|
人工智能 JSON 安全
Fastjson远程代码执行漏洞,阿里云AI安全为您保驾护航
阿里云AI安全产品联动防御Fastjson攻击
2187 12
Fastjson远程代码执行漏洞,阿里云AI安全为您保驾护航
|
8天前
|
云安全 人工智能 安全
|
8天前
|
人工智能 自然语言处理 数据挖掘
Qwen3.8-Max-Preview深度全解析:2.4万亿参数旗舰MoE模型+Token Plan限时优惠完整落地指南
2026年7月,全新旗舰级混合专家大模型Qwen3.8-Max-Preview正式开放抢先体验,作为通义千问Qwen3系列规格最高、综合推理能力顶尖的新一代模型,该模型总参数量达到2.4万亿(2.4T),是当前线上可调用的原生多模态旗舰模型,综合推理水准对标海外顶级Fable 5模型,在复杂工程开发、长文档深度分析、多步骤智能体自治、跨境多语言创作、海量数据挖掘五大高难度业务场景实现跨越式性能提升。
981 1
|
10天前
|
人工智能
Qwen3.8抢先体验!正式版即将发布并开源!
千问Qwen3.8即将开源,参数达2.4T,进化速度以“天”计,实力媲美Fable 5。预览版Qwen3.8-Max已上线阿里Token Plan等平台,限时优惠:日间Credits低至1折,夜间更优,个人/团队版月付仅35元起!
986 44
|
8天前
|
人工智能 自然语言处理 数据挖掘
最新版通义千问(Qwen3.8-Max-Preview)功能介绍
2026年,通义千问正式推出全新旗舰级大模型 **Qwen3.8-Max-Preview 预览版**,作为首款突破万亿参数规格的新一代基座模型,该模型总参数量达到**2.4万亿**,采用全新迭代的MoE混合专家架构,综合推理性能、长文本处理、多模态理解、复杂任务规划能力全面超越前代Qwen3.7-Max版本,整体实力跻身全球第一梯队,可对标海外顶级旗舰模型,是当前面向复杂工程开发、多智能体协同、超长文档解析、专业办公自动化场景的最优国产基座模型。
988 0
|
6天前
|
自然语言处理 测试技术 API
通义千问Qwen3.8-Max-Preview全功能解析:2.4万亿参数旗舰模型深度使用指南
在大模型技术持续迭代的当下,通义千问推出的Qwen3.8-Max-Preview作为新一代旗舰预览版模型,凭借2.4万亿参数的超大规模、多模态融合能力与全场景适配特性,成为开发者与企业用户探索AI应用的核心工具。该模型采用稀疏混合专家(MoE)架构,是通义千问首个突破万亿参数的多模态模型,可同时处理文本、图像、视频与文档等多种数据形态,在全栈代码开发、复杂逻辑推理、长文档分析与多智能体协作等场景实现跨越式升级。本文将全面拆解Qwen3.8-Max-Preview的核心功能,详解API调用流程与配置方法,覆盖多场景实战技巧,帮助用户快速掌握这款旗舰模型的使用方法,充分释放其性能潜力。
476 1
|
9天前
|
人工智能 自然语言处理 数据挖掘
Qwen3.8-Max 预览版全解析:2.4 万亿参数旗舰模型,Token Plan 限时优惠指南
Qwen3.8-Max-Preview是通义千问Qwen3系列旗舰MoE大模型,参数达2.4万亿,综合推理能力居行业第一梯队。支持思考/快速双模式,擅长大模型五大高难场景。现于阿里云百炼Token Plan、Qoder及QoderWork上线体验,个人版低至39元/月。在阿里云百炼官网:https://t.aliyun.com/U/fPVHqY 免费领取千万Tokens
687 1
Qwen3.8-Max 预览版全解析:2.4 万亿参数旗舰模型,Token Plan 限时优惠指南