首页
/ Scrapling Spider 爬虫系统架构解析:调度、多会话与断点续爬的数据流全览

Scrapling Spider 爬虫系统架构解析:调度、多会话与断点续爬的数据流全览

2026-09-06 11:34:07作者:伍霜盼Ellen

本文深入解析 Scrapling 的 Spider 爬虫系统架构:它是一套受 Scrapy 启发的异步爬取框架,专为并发、多会话爬取设计,并内置了调度、并发控制与断点续爬(pause/resume)能力。读完后,你将理解一次爬取从 start_urls 到结果导出的完整数据流、各核心组件(Spider、Crawler Engine、Scheduler、Session Manager、Checkpoint、Response Cache)的分工,并能基于 Spider 基类 编写支持断点恢复的多会话爬虫。

Scrapling Spider 架构图:Spider 产生请求,Scheduler 优先级队列与去重,Crawler Engine 驱动并发下载并回调解析,Session Manager 路由到不同会话类型

一、系统定位:把解析引擎与 Fetcher 统一进一个爬取 API

Scrapling 的 Spider 系统不是一个孤立的爬取器,而是把 Scrapling 已有的解析引擎(Selector/Response)与各类型 fetcher 整合进一个统一的异步爬取 API,并在其之上叠加了三样东西:

  • 调度(Scheduling):基于优先级队列的请求队列与内置 URL 去重;
  • 并发控制(Concurrency Control):全局并发、按域名并发、下载延迟、自动节流;
  • 断点检查点(Checkpointing):可配置间隔定期保存爬取状态,支持暂停与恢复。

从源码结构看,整个系统实现于 scrapling/spiders 包,模块划分非常清晰:spider.py(用户基类)、engine.py(引擎)、scheduler.py(调度器)、session.py(会话管理)、checkpoint.py(检查点)、cache.py(开发缓存)、result.py(输出与统计)、request.py(请求与指纹)、throttle.py(自动节流)。公共 API 通过 scrapling/spiders/init.py 导出,包括 SpiderRequestResponseCrawlResult 以及 CrawlSpiderSitemapSpiderShopifySpider 等模板。

二、数据流:一次爬取运行的七个步骤

下面按引擎实际执行顺序拆解一次爬取的数据流(对应 CrawlerEngine.crawl() 的主循环)。

  1. Spider 产生第一批 Request 对象。默认为 start_urls 中的每个 URL 创建一个请求,也可以覆写 start_requests() 实现自定义逻辑。从 spider.py 可以看到默认实现:

    for url in self.start_urls:
        yield Request(url, sid=self._session_manager.default_session_id)
    

    每个请求都携带默认会话的 sid(会话 ID),这是后续路由的依据。若未设置 start_urls 又未覆写 start_requests(),引擎会抛出 RuntimeError

  2. Scheduler 接收请求并入优先级队列,同时生成指纹。指纹由 URL、HTTP 方法、请求体与会话 ID 计算而来,指纹相同的请求会被去重丢弃(除非 dont_filter=True)。优先级高的请求先出队,实现见 scheduler.py

    fingerprint = request.update_fingerprint(self._include_kwargs, self._include_headers, self._keep_fragments)
    if not request.dont_filter and fingerprint in self._seen:
        return False  # Dropped duplicate request
    item = (-request.priority, counter, request)  # 负号使高优先级先出队
    
  3. Crawler Engine 出队下一个请求,遵守全局与按域名的并发上限以及下载延迟。若启用了 robots_txt_obey,引擎会在处理前检查目标域名的 robots.txt 规则——被禁止的请求被静默丢弃,并计入 robots_disallowed_count 统计。拿到请求后,引擎把它交给 Session Manager,由其根据请求的 sid 路由到正确的会话。

    一个容易忽略的细节:当 robots_txt_obey 开启时,引擎还会把 robots.txt 中的 Crawl-delay / Request-ratedownload_delay 取最大值作为该域名的实际延迟下限,见 engine.py 的 _get_domain_delay()

  4. 会话抓取页面并返回 Response 对象。引擎记录统计信息(请求数、响应字节数、状态码分布)并检查响应是否被反爬拦截。若被判定为 blocked,引擎最多重试 max_blocked_retries 次(默认 3 次)。拦截判定与重试逻辑均可自定义:

    • is_blocked() 的默认实现是检查状态码是否属于 BLOCKED_CODES = {401, 403, 407, 429, 444, 500, 502, 503, 504}spider.py#L16);
    • 重试时在 engine.py 中会把新请求的 priority 减 1(避免立即重试)、置 dont_filter=True、并剥离旧的 proxy/proxies 参数,再交给用户钩子 retry_blocked_request() 做最后加工。
  5. 引擎把 Response 交给请求的回调。回调是一个 async 生成器,每次 yield 一个字典(视为抓取到的条目)或一个后续 Request(送回 Scheduler 排队);yield None 表示无产出。条目在入结果前会先经过 on_scraped_item() 钩子,钩子返回 None 时该条目被静默丢弃并计入 items_dropped

  6. 循环重复第 2 步,直到调度器为空且没有活跃任务,或者 Spider 被暂停。

  7. 断点保存。若构造 Spider 时设置了 crawldir,引擎会周期性把检查点(pending 请求 + 已见 URL 集合)写到磁盘;优雅停机(Ctrl+C)时会再存一次最终检查点。下次以相同 crawldir 启动时,Spider 从上次的地方恢复——跳过 start_requests(),直接还原调度器状态engine.py 中明确日志 Resuming from checkpoint, skipping start_requests())。

三、核心组件深度剖析

3.1 Spider:用户交互的中心类

Spider 是抽象基类,用户通过子类化它来定义爬取逻辑。最小可用示例:

from scrapling.spiders import Spider, Response, Request

class MySpider(Spider):
    name = "my_spider"
    start_urls = ["https://example.com"]

    async def parse(self, response: Response):
        for link in response.css("a::attr(href)").getall():
            yield response.follow(link, callback=self.parse_page)

    async def parse_page(self, response: Response):
        yield {"title": response.css("h1::text").get(""))}

spider.py 的源码可以整理出全部类属性及其默认值,这是配置爬虫行为的主要手段:

类属性 默认值 说明
name None Spider 名称,必填,否则初始化报错
start_urls [] 起始 URL 列表
allowed_domains set() 允许的域名白名单,空则不限制
robots_txt_obey False 是否遵守 robots.txt
development_mode False 开发模式:启用响应缓存回放
development_cache_dir None 缓存目录,默认 .scrapling_cache/{name}
concurrent_requests 4 全局并发上限
concurrent_requests_per_domain 0(不限制) 每域名并发上限
download_delay 0.0 每请求基础下载延迟(秒)
max_blocked_retries 3 被拦截请求的最大重试次数
autothrottle_enabled False 是否启用自动节流
autothrottle_start_delay 5.0 节流起始延迟
autothrottle_max_delay 60.0 节流延迟上限
autothrottle_target_concurrency None 目标并发,缺省取 concurrent_requests_per_domain 或 1
autothrottle_block_backoff True 被拦截时按 Retry-After 或翻倍延迟退避
fp_include_kwargs False 指纹是否包含请求 kwargs
fp_keep_fragments False 指纹是否保留 URL 片段
fp_include_headers False 指纹是否包含请求头
logging_level / logging_format / log_file 日志级别(默认 DEBUG)、格式与文件输出

构造函数接受两个参数:crawldir(启用断点续爬的目录)和 interval(周期性检查点间隔,默认 300.0 秒即 5 分钟,见 spider.py#L106-L111)。

生命周期钩子(均可覆写):

  • on_start(resuming: bool):爬取开始前调用,resuming 标识是否从检查点恢复,适合做初始化或恢复分支逻辑;
  • on_close():爬取结束后调用,用于清理;
  • on_error(request, error):请求或回调出错时调用;
  • on_scraped_item(item):每个条目入库前的处理钩子,返回 None 可静默丢弃该条目;
  • is_blocked(response)retry_blocked_request(request, response):自定义拦截判定与重试前的请求加工;
  • configure_sessions(manager):注册会话,默认实现添加一个名为 "default"FetcherSession。第一个注册的会话成为默认会话(供 start_requests() 使用)。

3.2 Crawler Engine:你不直接碰的编排者

CrawlerEngine 负责整个爬取的编排:主循环、并发限制、经 Session Manager 分发请求、处理回调结果。用户不直接实例化它——Spider.start()Spider.stream() 会替你完成(start() 内部通过 anyio 的 asyncio 后端驱动,可选 use_uvloop=True 使用更快的 uvloop/winloop 事件循环)。

主循环的关键机制(engine.py):

  • 基于 anyio task group 并发执行,只有当活跃任务数低于 concurrent_requests 时才从调度器出队并 spawn 新任务,避免一次创建数千个等待中的任务
  • 并发限制通过 CapacityLimiter 实现:全局一个,按域名各一个(当 concurrent_requests_per_domain > 0 时),取域名 limiter 优先;
  • 队列空且无活跃任务时判定爬取完成;
  • Ctrl+C 触发 request_pause():第一次请求优雅暂停(等待在途请求完成),第二次强制停止(立即取消 task group)。强制停止时若启用了检查点,会先保存检查点再取消,避免状态丢失;
  • 爬取正常完成(非暂停退出)时自动清理检查点文件(_checkpoint_manager.cleanup())。

robots.txt 的处理还有"预取"优化:爬取循环开始前,引擎从 start_urls 提取唯一域名并批量预热 robots 缓存(_prefetch_robots_txt()),中途新发现的域名才会在首次访问时拉取。

自动节流(AutoThrottle,实现于 throttle.py)根据观测到的响应延迟动态调整每域名的延迟:服务器快就提速,慢或敌意就退避;当 block_backoff 开启时,被拦截会按 Retry-After 头等待或将该域名延迟翻倍(BLOCK_BACKOFF_FACTOR = 2.0)。

3.3 Scheduler:优先级队列 + 内置去重

Scheduler 是一个基于 asyncio.PriorityQueue 的优先级队列,同时承担 URL 去重职责。其内部结构值得注意:

  • _queue:真正的异步优先队列,元素为 (-priority, counter, request)——负号保证优先级数值大的先出队,counter 作为同优先级下的入队序号打破平局;
  • _seen:已见指纹集合,用于去重;
  • _pending:镜像字典,无需排空队列即可生成快照;
  • _inflight:跟踪出队后尚未完成的任务,供检查点统计使用。

请求指纹的生成逻辑在 Request.update_fingerprint():对 sid、请求体(data/json 参数序列化后的 hex)、HTTP 方法、规范化后的 URL(canonicalize_url,默认丢弃 fragment)做 orjson 稳定序列化后取 SHA1。可选地把全部 kwargs 或请求头也纳入指纹(fp_include_kwargs / fp_include_headers),适用于"同一 URL 不同请求头应视为不同请求"的场景。结果缓存在 request._fp 上,避免重复计算。

检查点支持依赖两个方法:

  • snapshot():返回按优先级排序的 pending 请求列表与 _seen 集合的副本;
  • restore(data):从 CheckpointData 还原 seen 集合,并把请求重新按序放入队列。

3.4 Session Manager:多会话路由

SessionManager 管理一组具名会话实例。源码中的类型标注明确了每个会话只能是三种之一(session.py#L9):

这正是 Spider 系统相对 Scrapy 的差异点之一:Scrapy 的 Downloader + Middlewares 是单一下载通道,而 Scrapling 允许同一个爬虫内混用不同类型的会话——例如常规页面走 FetcherSession,需要 JS 渲染的详情页走 AsyncStealthySession,只需在 configure_sessions() 中注册多个会话、在请求上指定不同的 sid 即可。

add(session_id, session, *, default=False, lazy=False) 支持两种启动方式:

  • 随 Spider 启动而启动(默认):SessionManager.start() 在爬取开始时 __aenter__ 所有非 lazy 会话;
  • 懒加载(lazy):会话直到第一次被 sid 引用时才启动,适合浏览器型会话——爬虫可能从头到尾都用不到它,没必要付出启动浏览器的代价。

路由逻辑见 fetch():按 request.sid(缺省回退到默认会话 ID)取出会话;对 FetcherSession 走底层 _ASyncSessionLogic._make_request,其他会话调用 session.fetch()。返回的 Response 会被回填 request 引用,并把请求的 meta 合并进响应 meta(响应 meta 优先)。

3.5 Checkpoint System:原子写入的断点续爬

CheckpointManager 把爬虫状态(pending 请求 + 已见 URL 指纹集合,封装在 CheckpointData dataclass 中)序列化到 crawldir/checkpoint.pkl

  • 原子写:先写 checkpoint.tmpreplace 重命名为正式文件,防止写一半损坏(checkpoint.py#L42-L61);
  • 周期性保存:间隔由构造函数 interval 参数控制(默认 5 分钟),引擎主循环每轮检查是否到期(_is_checkpoint_time());interval=0 时禁用周期保存,仅在暂停时保存;
  • 保存时机:周期到期、优雅暂停(Ctrl+C 第一次)、强制停止(Ctrl+C 第二次);
  • 自动清理:爬取正常完成(非暂停)时删除检查点文件。

一个精妙的工程细节:Request 在 pickle 时不保存回调函数对象本身,而是保存回调的函数名(request.py 的 __getstate__/__setstate__);恢复时引擎遍历所有请求调用 request._restore_callback(self.spider),按名字从 Spider 实例上重新取回方法。这解决了"闭包/方法不可 pickle"的经典难题,也是断点恢复能跳过 start_requests() 而回调仍然有效的原因。恢复逻辑见 engine.py 的 _restore_from_checkpoint()

3.6 Response Cache:开发模式的响应回放

development_mode = True 时,引擎会启用 ResponseCacheManager:每次抓取到的响应都落盘,之后的运行直接回放,用于反复调试 parse() 逻辑而不重复请求目标服务器(官方定位是开发用途,不是生产用途)。

从源码看其存储格式:

  • 每个响应存为 <fingerprint.hex()>.json 文件,按请求指纹命名;
  • JSON 字段包括 urlcontentbase64 编码,保证二进制内容存活)、statusreasonencodingcookiesheadersrequest_headersmethod
  • cookies 保留原始形状:浏览器引擎响应是完整 cookie 字典的 tuple,静态引擎是扁平 dict,读取时按需还原;
  • 同样采用 tmp 文件 + replace 的原子写;
  • 缓存目录默认 .scrapling_cache/{spider.name}engine.py#L57-L58),可由 development_cache_dir 覆盖。

引擎侧的命中/未命中会计入 cache_hits / cache_misses 统计;命中缓存时直接复用响应跑回调,不再发起真实请求。

3.7 Output:ItemList 与 CrawlStats

抓取到的条目收集在 ItemListlist 的子类)中,内置四种导出方法:

  • to_json(path, *, indent=False):orjson 序列化,支持 numpy 类型,可选 2 空格缩进;
  • to_jsonl(path):JSON Lines,每行一个对象;
  • to_csv(path, *, fields=None, delimiter=","):列默认取全部条目中出现过的所有 key(按出现顺序),缺失单元格留空,非标量值(嵌套 dict/list)写为 JSON;
  • to_xml(path, *, root_tag="items", item_tag="item", indent=True):非法 XML 标签名的 key 会被改写并在 name 属性保留原名,控制字符被剔除。

统计由 CrawlStats dataclass 承载,字段包括:requests_countfailed_requests_countoffsite_requests_countrobots_disallowed_countblocked_requests_countcache_hits/cache_missesitems_scraped/items_droppedresponse_bytesresponse_status_count(状态码分布)、domains_response_bytessessions_requests_count(按会话统计)、proxies(用到的代理)、log_levels_counter(各级别日志条数)等,并提供 elapsed_secondsrequests_per_second 计算属性。爬取结束时引擎会把 stats.to_dict() 整体打印到日志。

最终 Spider.start() 返回 CrawlResultstats + items + paused 标志,completed 属性即 not paused)。Spider.stream() 则返回 async 生成器,基于容量为 100 的内存对象流逐条产出条目;迭代期间可通过 spider.stats 属性访问实时统计——注意源码注释明确说明:stream 模式下没有 SIGINT 暂停/恢复处理

四、与 Scrapy 的概念对照表

如果你熟悉 Scrapy,下面的映射表可以直接迁移心智模型(response.follow() 在两者中的签名一致):

概念 Scrapy Scrapling
Spider 定义 scrapy.Spider 子类 scrapling.spiders.Spider 子类
初始请求 start_requests() async start_requests()
回调 def parse(self, response) async def parse(self, response)
跟进链接 response.follow(url) response.follow(url)
条目输出 yield dictyield Item yield dict
请求调度 Scheduler + Dupefilter 内置去重的 Scheduler
下载 Downloader + Middlewares 支持多会话的 Session Manager
条目处理 Item Pipelines on_scraped_item() 钩子
拦截检测 自定义中间件 内置 is_blocked() + retry_blocked_request() 钩子
并发 CONCURRENT_REQUESTS 设置项 concurrent_requests 类属性
域名过滤 allowed_domains allowed_domains
Robots.txt ROBOTSTXT_OBEY 设置项 robots_txt_obey 类属性
暂停/恢复 JOBDIR 设置项 crawldir 构造函数参数
导出 Feed exports result.items.to_json() / to_jsonl() / to_csv() / to_xml() 或经钩子自定义
运行方式 scrapy crawl spider_name MySpider().start()
流式输出 async for item in spider.stream()
多会话 同一爬虫内多类型会话并存

五、小结

Scrapling 的 Spider 系统用"优先级队列调度器 + 会话路由 + 原子检查点"三个机制,把 Scrapling 已有的静态/动态/隐身 fetcher 与解析引擎统一进一个异步爬取 API:

  • spider.py 的类属性与钩子出发,用最小代码获得并发、robots、拦截重试等全套能力;
  • engine.py 的主循环展示了并发限制、延迟解析、暂停/恢复与检查点时机的完整实现;
  • scheduler.pyrequest.py 解释了去重指纹与优先级排序的具体算法;
  • checkpoint.pycache.pyresult.py 分别对应断点续爬、开发回放与结果导出。

如果你只需要单次请求,直接使用 Fetchers 即可;当任务规模上升到"大规模、可暂停、多会话"的完整爬取时,这套 Spider 系统就是对应的解决方案。相关实现可结合 tests/spiders/ 下的测试用例(如 test_engine.pytest_checkpoint.pytest_scheduler.py)进一步验证各组件行为。

登录后查看全文
热门项目推荐
相关项目推荐