深入理解 Scrapy 架构:Engine、Scheduler、Downloader 与 Spiders 的完整数据流
【免费下载链接】scrapyScrapy, a fast high-level web crawling & scraping framework for Python.项目地址: https://gitcode.com/GitHub_Trending/sc/scrapy
本文以 Scrapy 官方架构文档 architecture.rst 为主体,完整讲解 Scrapy 的组件划分与数据流(从 Spider 发起初始请求,到响应解析、Item 落库的 9 步循环),并结合仓库中 scrapy/core/engine.py、scrapy/core/scheduler.py、scrapy/core/downloader/init.py 等核心源码,说明每个组件的真实实现位置、关键方法与默认配置,帮助你在编写 Spider、中间件或自定义调度器时,能准确理解每个数据包的流转路径与事件触发时机。
架构总览
Scrapy 采用"执行引擎(Engine)为中心"的组件化架构:Engine 负责控制系统内所有组件间的数据流,并在特定事件发生时触发信号。系统内的组件职责如下(摘自官方文档并补充源码位置):
| 组件 | 职责 | 源码位置 |
|---|---|---|
| Engine | 控制各组件之间的数据流,触发事件(信号) | scrapy/core/engine.py |
| Scheduler | 接收 Engine 发来的 Request,入队后在 Engine 索取时再吐出 | scrapy/core/scheduler.py |
| Downloader | 抓取网页,并把 Response 送回 Engine | scrapy/core/downloader/init.py |
| Spiders | 用户编写的自定义类,解析响应、提取 Items 与新的 Requests | scrapy/spiders/ |
| Item Pipeline | 处理 Spider 提取出的 Items:清洗、校验、持久化(如入库) | scrapy/pipelines/ |
| Downloader Middlewares | 位于 Engine 与 Downloader 之间的钩子,处理下行 Request 与上行 Response | scrapy/core/downloader/middleware.py |
| Spider Middlewares | 位于 Engine 与 Spiders 之间的钩子,处理 Spider 输入(Response)与输出(Items、Requests) | scrapy/core/spidermw.py |
| Extensions | 不参与数据流,用于在特定事件发生时执行逻辑 | scrapy/extensions/ |
图中红色箭头勾勒出的数据流,即下面逐条展开的 9 步循环。
数据流(Data Flow)
官方文档将 Scrapy 的数据流归纳为以下 9 步,由执行引擎驱动:
- Engine 从 Spider 获取初始 Requests:引擎打开 Spider 后,消费其
start()(即 start_urls 等起始请求)产出的 Request。 - Engine 将 Requests 交给 Scheduler 排队,并向 Scheduler 索取下一个要爬的 Request。
- Scheduler 返回下一个 Request给 Engine。
- Engine 将 Request 发送给 Downloader,途中依次经过Downloader Middlewares(对应
process_request钩子)。 - 页面下载完成后,Downloader 生成携带该页面的 Response送回 Engine,途中经过Downloader Middlewares(对应
process_response钩子)。 - Engine 接收 Response 并转交 Spider 处理,途中经过Spider Middlewares(对应
process_spider_input钩子)。 - Spider 解析 Response,产出被抓取的数据项(Items)与待跟进的新 Requests 返回 Engine,途中经过Spider Middlewares(对应
process_spider_output钩子)。 - Engine 将 Items 送交 Item Pipelines,将新的 Requests 送回Scheduler,并继续向 Scheduler 索取下一个 Request。
- 上述过程从第 3 步开始循环重复,直到 Scheduler 中不再有未处理的请求。
源码印证:Engine 如何驱动这个循环
从 scrapy/core/engine.py 中的ExecutionEngine类可以看到上述流程的具体落点:
- 入队(对应第 1、2 步):
ExecutionEngine.crawl()是所有新请求进入管线的入口,它调用_schedule_request()——先发送request_scheduled信号,再调用scheduler.enqueue_request(request);若 Scheduler 拒绝(返回False,例如被去重过滤器拦截),则发送request_dropped信号。对应源码见 engine.py 的 crawl/_schedule_request。 - 出队并下载(对应第 3、4 步):
_start_scheduled_request()不断从self._slot.scheduler.next_request()取出请求,取不到时发送scheduler_empty信号;取到后调用_download()交给 Downloader,并在结果处理中判断:若 Downloader 中间件返回的是新的Request(典型场景是重定向),会再次crawl()注入管线。对应源码见 engine.py 的 _start_scheduled_request。 - 背压控制(needs_backout):
needs_backout()会在"引擎停止运行、槽位正在关闭、Downloader 并发占满、Scraper 积压超限"任一情况发生时返回True,Engine 随即暂停向 Downloader 投递请求。这是保证高吞吐下内存稳定的关键机制,见 engine.py 的 needs_backout。 - 响应转交 Spider(对应第 5、6 步):
_handle_downloader_output()在拿到 Response 后调用scraper.enqueue_scrape(result, request),由 Scraper 负责经过 Spider Middlewares 喂给 Spider。见 engine.py 的 _handle_downloader_output。 - 空闲判定与结束(对应第 8、9 步):
spider_is_idle()只有在"Scraper 槽位空闲、Downloader 无在途请求、Scheduler 无待处理请求"三者同时成立时才返回True;此时触发spider_idle信号,若没有DontCloseSpider处理器阻止,则关闭 Spider,本轮爬取结束。见 engine.py 的 spider_is_idle 与 _spider_idle。
一个容易被忽略的细节:_Slot(engine.py L65-L98)中维护着一个inprogress请求集合和heartbeat心跳调用(间隔_SLOT_HEARTBEAT_INTERVAL = 5.0秒)。即使 Scheduler 报告仍有待处理请求却暂时取不出请求,Engine 也会通过心跳定期重试,防止爬虫卡死。
组件深潜
Scrapy Engine:数据流的总调度者
Engine 的职责在文档中的原话是"控制系统内所有组件之间的数据流,并在特定动作发生时触发事件"。在实现上,ExecutionEngine在__init__中装配三大协作对象(engine.py L104-L154):
Downloader:由设置项DOWNLOADER(默认"scrapy.core.downloader.Downloader",见 default_settings.py L315)加载;Scheduler:由设置项SCHEDULER加载,且必须通过BaseScheduler接口检查,否则抛出TypeError;Scraper:负责"响应 → Spider"这一段,内部封装了 Spider Middlewares 与 Item Pipelines。
值得注意的是,当前仓库版本的 Engine 以async def协程为主体(如start_async()、stop_async()、open_spider_async()),旧的 Deferred 风格方法(start()、stop())已标记为弃用,仅做兼容转发——这说明 Scrapy 的事件驱动底座正在从纯 Twisted 向 Twisted + asyncio 双轨演进(详见后文"事件驱动网络"一节)。
Scheduler:请求的队列与顺序
scrapy/core/scheduler.py 定义了最小调度器接口BaseScheduler(L52-L124),任何自定义 Scheduler 必须实现三个抽象方法:
| 方法 | 语义 |
|---|---|
has_pending_requests() | 是否还有已入队的请求 |
enqueue_request(request) | 接收 Engine 的请求;返回False时 Engine 会发出request_dropped信号且不再重试(默认实现中,被去重过滤器拒绝的请求即返回False) |
next_request() | 返回下一个待处理请求;返回None表示当前 reactor 周期内无可发送请求,Engine 会继续调用直到has_pending_requests()为False |
元类BaseSchedulerMeta还通过__subclasscheck__在运行时检查这三个方法是否存在且可调用——这也解释了 Engine 初始化时对 Scheduler 的接口校验逻辑。
默认实现Scheduler(scheduler.py L127-L498)将请求存入按Request.priority排序的优先级队列(SCHEDULER_PRIORITY_QUEUE)。几个直接影响爬取行为的实现事实:
- 去重:
enqueue_request()先调用DUPEFILTER_CLASS(默认scrapy.dupefilters.RFPDupeFilter)的request_seen(),已被过滤且未设置dont_filter的请求直接返回False; - 内存/磁盘双队列:默认全部请求走内存队列;启用
JOBDIR(断点续爬)时同时创建磁盘队列,且不可序列化的请求自动回落到内存队列(见enqueue_request与_dqpush的ValueError分支,scheduler.py L364-L436);同一优先级下,内存队列优先于磁盘队列; - 抓取顺序:默认内存队列是 LIFO 栈,因此爬取呈"深度优先"(DFO)顺序;文档明确给出改为 BFO(广度优先)的三个设置:
DEPTH_PRIORITY = 1、SCHEDULER_DISK_QUEUE = "scrapy.squeues.PickleFifoDiskQueue"、SCHEDULER_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue"(见 scheduler.py 文档字符串 L181-L194); - 统计:
scheduler/enqueued、scheduler/enqueued/disk、scheduler/enqueued/memory、scheduler/dequeued等计数在入队/出队时逐次累加,可用于观察调度行为; - JOBDIR 文件布局:
requests.queue/目录保存未下载请求,active.json保存优先级队列状态并在作业停止时写出、恢复时读入(_write_dqs_state/_read_dqs_state)。源码文档字符串特别提示:这些文件属于实现细节,不要依赖其结构。
Downloader:并发槽位与下载处理器
Downloader(scrapy/core/downloader/init.py L83 起)初始化时读取的关键设置:
CONCURRENT_REQUESTS(总并发)、CONCURRENT_REQUESTS_PER_DOMAIN(按域名并发)、CONCURRENT_REQUESTS_PER_IP(按 IP 并发);DOWNLOAD_DELAY与RANDOMIZE_DOWNLOAD_DELAY:每个域名槽位(Slot数据类,L44-L80)按delay控制请求节奏,开启随机化后实际延迟在0.5×delay ~ 1.5×delay之间均匀取随机值;DOWNLOAD_SLOTS:可为特定槽名(如download_slot、<domain>)单独配置并发与延迟;- 实际的请求发送委托给
DownloadHandlers(scrapy/core/downloader/handlers/),按 URL scheme 分发到 http11、http2、httpx、ftp、file、data 等具体处理器。
Downloader.fetch()的调用链是:请求先注册进active集合,再交给DownloaderMiddlewareManager.download_async()(走完process_request→ 下载 →process_response的完整中间件链),最后按槽位节流发出。Engine 正是通过downloader.needs_backout()(由这些槽位容量决定)与downloader.active(在途请求集合)实现背压与空闲判断。
Spiders 与 Scraper
Spider 是用户代码:解析响应、产出 Items 与新 Requests。Scrapy 侧真正承载"喂给 Spider"这一职责的是Scraper(scrapy/core/scraper.py),它在 L106-L127 中装配了:
SpiderMiddlewareManager(SPIDER_MIDDLEWARES 中间件链);ItemPipelineManager(由ITEM_PROCESSOR设置加载,默认即 Item Pipelines 管理器);CONCURRENT_ITEMS:并发处理 Item 的并发度。
Scraper.Slot(scraper.py L62-L103)为每个运行中的 Spider 维护一个响应队列(queue)、活跃请求集合(active)与活跃数据量(active_size):当active_size超过max_active_size(默认 5,000,000 字节)时needs_backout()返回True,Engine 暂停投递新响应,防止大页面拖垮内存。is_idle()则在队列、活跃请求、Item 处理三者均为空时返回True——这正是 Engine 判定爬虫空闲的输入之一。
Item Pipeline
Item Pipeline 在 Items 被提取后对其进行处理,典型任务是清洗、校验与持久化(如写入数据库)。框架内建实现位于 scrapy/pipelines/(如FilesPipeline、ImagesPipeline、MediaPipeline),用户管线通过设置项ITEM_PIPELINES注册,由Scraper中的ItemPipelineManager统一驱动。
Downloader Middlewares:位于 Engine 与 Downloader 之间
文档说明:Downloader 中间件是位于 Engine 与 Downloader 之间的特定钩子,处理从 Engine 流向 Downloader 的请求,以及从 Downloader 流回 Engine 的响应。管理器DownloaderMiddlewareManager(scrapy/core/downloader/middleware.py)把三类方法挂成两条方向相反的链:
process_request按优先级正序执行(append),任一中间件返回 Response/Request 则短路,直接跳回响应链;process_response与process_exception按优先级逆序执行(appendleft),形成"洋葱"结构(middleware.py L43-L52, L96-L156)。
仓库内置的默认中间件及优先级见 default_settings.py 的 DOWNLOADER_MIDDLEWARES_BASE:
DOWNLOADER_MIDDLEWARES_BASE = { # Engine side "scrapy.downloadermiddlewares.offsite.OffsiteMiddleware": 50, "scrapy.downloadermiddlewares.robotstxt.RobotsTxtMiddleware": 100, "scrapy.downloadermiddlewares.httpauth.HttpAuthMiddleware": 300, "scrapy.downloadermiddlewares.downloadtimeout.DownloadTimeoutMiddleware": 350, "scrapy.downloadermiddlewares.defaultheaders.DefaultHeadersMiddleware": 400, "scrapy.downloadermiddlewares.useragent.UserAgentMiddleware": 500, "scrapy.downloadermiddlewares.retry.RetryMiddleware": 550, "scrapy.downloadermiddlewares.redirect.MetaRefreshMiddleware": 580, "scrapy.downloadermiddlewares.httpcompression.HttpCompressionMiddleware": 590, "scrapy.downloadermiddlewares.redirect.RedirectMiddleware": 600, "scrapy.downloadermiddlewares.cookies.CookiesMiddleware": 700, "scrapy.downloadermiddlewares.httpproxy.HttpProxyMiddleware": 750, "scrapy.downloadermiddlewares.stats.DownloaderStats": 850, "scrapy.downloadermiddlewares.httpcache.HttpCacheMiddleware": 900, # Downloader side }注释中的 "Engine side" / "Downloader side" 直观标出了中间件链两端;用户中间件通过DOWNLOADER_MIDDLEWARES设置与上述默认项合并,优先级数值越小越靠近 Engine。每个中间件类的具体实现位于 scrapy/downloadermiddlewares/。
Spider Middlewares:位于 Engine 与 Spiders 之间
Spider 中间件处理 Spider 的输入(Response)与输出(Items 和 Requests)。管理器SpiderMiddlewareManager(scrapy/core/spidermw.py)维护四条方法链:process_start、process_spider_input、process_spider_output、process_spider_exception。其中异常路径值得一提:若 Spider 回调抛出异常,_process_spider_exception()会依次询问各中间件;一旦某个中间件返回可迭代对象,控制权立即交还process_spider_output链继续处理(spidermw.py L116-L150),这让中间件有机会"吞掉"错误请求并产出替代结果。
默认 Spider 中间件见 default_settings.py 的 SPIDER_MIDDLEWARES_BASE:
SPIDER_MIDDLEWARES_BASE = { # Engine side "scrapy.spidermiddlewares.start.StartSpiderMiddleware": 25, "scrapy.spidermiddlewares.httperror.HttpErrorMiddleware": 50, "scrapy.spidermiddlewares.referer.RefererMiddleware": 700, "scrapy.spidermiddlewares.urllength.UrlLengthMiddleware": 800, "scrapy.spidermiddlewares.depth.DepthMiddleware": 900, "scrapy.spidermiddlewares.metacopy.MetaCopyDetectionMiddleware": 1000, # Spider side }实现代码位于 scrapy/spidermiddlewares/。
Extensions
Extensions 与上面组件不同,它们在数据流中没有固定角色,而是监听爬虫生命周期事件(Spider 打开/关闭、错误计数、内存占用等)并执行逻辑。默认启用的扩展见 default_settings.py 的 EXTENSIONS_BASE,包括CoreStats、LogCount、TelnetConsole、MemoryUsage、CloseSpider、FeedExporter、LogStats、SpiderState、AutoThrottle、RemoteControl等,实现位于 scrapy/extensions/。
事件驱动的网络模型
官方文档明确指出:Scrapy 构建于Twisted——Python 流行的事件驱动网络框架——之上,因此采用非阻塞(异步)代码实现并发。这意味着:整个引擎内没有线程等待网络 I/O,Engine、Downloader、Scheduler 的协作全部由事件循环(reactor)调度,一个进程即可支撑大量并发连接。
从当前仓库的源码结构可以进一步看到底座的演进轨迹:ExecutionEngine的核心方法已全部提供async版本(如open_spider_async()、close_spider_async(),见 engine.py L547 起),Deferred 风格旧 API 被标记为ScrapyDeprecationWarning;同时scrapy/utils/asyncio.py提供了create_looping_call、maybe_deferred_to_future等桥接工具,让循环调用、心跳等机制在 Twisted reactor 与 asyncio 事件循环下都能运行。换言之,文档所描述的"事件驱动 + 非阻塞"结论在当前版本中依然成立,只是事件循环的实现细节兼容了 asyncio 生态。学习时建议掌握 Deferred/协程的基本心智模型后再阅读 scrapy/utils/defer.py 中的桥接实现。
小结:把 9 步数据流映射到代码
| 数据流步骤 | 对应源码入口 |
|---|---|
| 1. 获取初始请求 | scrapy/core/spidermw.py 的 process_start → engine.py 的 _start_request_processing |
| 2/3. 入队与出队 | engine.py 的 crawl/_schedule_request 与 _start_scheduled_request、scheduler.py 的 enqueue_request/next_request |
| 4/5. 经 Downloader 中间件下载 | middleware.py 的 download_async + downloader/init.py 的 fetch |
| 6/7. 经 Spider 中间件解析 | spidermw.py 的 scrape_response_async、scraper.py |
| 8. Items 入管线、Requests 回调度 | scraper.py 中 ItemPipelineManager 的驱动 与 engine 的crawl()循环 |
| 9. 循环直至 Scheduler 清空 | engine.py 的 spider_is_idle/_spider_idle 与scheduler_empty信号 |
理解了这套"Engine 居中调度、队列控速、中间件分层拦截、信号驱动事件"的骨架后,你在配置CONCURRENT_REQUESTS、选择抓取顺序(DFO/BFO)、编写自定义 Middleware 或 Scheduler 时,都能准确判断改动会影响数据流中的哪一段。
【免费下载链接】scrapyScrapy, a fast high-level web crawling & scraping framework for Python.项目地址: https://gitcode.com/GitHub_Trending/sc/scrapy
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考