title: Scrapy 爬虫框架深入:从基础到分布式实战 slug: scrapy-framework-deep-dive date: 2026-06-23 tags: [Python, Scrapy, 爬虫, 分布式, Web Scraping] description: 深入解析 Scrapy 框架核心架构、组件机制与分布式实战,涵盖 Middleware、Pipeline、Extensions 及性能优化等进阶主题。
Scrapy 爬虫框架深入:从基础到分布式实战
前言
Scrapy 是 Python 生态中最强大的异步爬虫框架之一,广泛应用于数据采集、内容聚合、监控预警等场景。与简单的 requests + BeautifulSoup 脚本不同,Scrapy 提供了完整的请求调度、并发控制、数据管道、扩展机制等企业级能力。
本文将带你从架构原理出发,逐步深入到分布式实战,帮助你在实际项目中写出健壮、高效、可维护的爬虫。
一、Scrapy 架构深度解析
1.1 核心组件
Scrapy 的架构围绕五个核心组件展开,它们协同工作构成了一个高效的数据流引擎:
+-------------+
| Scheduler |
+------+------+
|
+------------+------------+
| |
+-----v-----+ +-------v-------+
| Downloader | | Engine |
+-----+-----+ +-------+-------+
| |
| +------v------+
| | Spiders |
| +------+------+
| |
+-----v-----+ |
| Pipeline | <----------------+
+-----------+
| 组件 | 职责 |
|---|---|
| Engine(引擎) | 核心调度中枢,协调所有组件之间的数据流动 |
| Scheduler(调度器) | 接收引擎发来的 Request,压入队列,去重后按优先级出队 |
| Downloader(下载器) | 执行 HTTP 请求,返回 Response 给 Engine |
| Spider(爬虫) | 解析 Response,提取 Item 或生成新的 Request |
| Item Pipeline(管道) | 依次处理提取到的 Item,清洗、验证、持久化 |
1.2 数据流生命周期
一次完整的 Scrapy 爬取流程如下:
- Engine 从 Spider 获取初始 Request,交给 Scheduler
- Scheduler 将 Request 入队,Engine 轮询获取下一个 Request
- Engine 将 Request 交给 Downloader Middleware 链,再由 Downloader 执行 HTTP 请求
- Downloader 返回 Response,经 Downloader Middleware 处理后交给 Engine
- Engine 将 Response 交给 Spider Middleware 链,最终进入 Spider 解析
- Spider 解析出 Item 或新的 Request
- Item → Engine → Item Pipeline(清洗、存储)
- Request → Engine → Scheduler(继续循环)
整个流程基于 Twisted 事件循环实现异步非阻塞 I/O,吞吐量远高于同步方案。
二、Request / Response 生命周期与 Middleware 机制
2.1 Middleware 处理链
Middleware 是 Scrapy 扩展能力的核心。每个 Request 和 Response 都会依次经过 Middleware 链,按优先级顺序处理。
Downloader Middleware 处理顺序:
Request → [Middleware 1 → Middleware 2 → ...] → Downloader
Response ← [Middleware 1 ← Middleware 2 ← ...] ← Downloader
2.2 自定义 Middleware 示例
# middlewares.py
import random
from scrapy import signals
from scrapy.downloadermiddlewares.useragent import UserAgentMiddleware
class RandomUserAgentMiddleware(UserAgentMiddleware):
"""随机 User-Agent 中间件"""
USER_AGENTS = [
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 '
'(KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 '
'(KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36',
'Mozilla/5.0 (X11; Linux x86_64; rv:127.0) Gecko/20100101 Firefox/127.0',
]
def process_request(self, request, spider):
request.headers['User-Agent'] = random.choice(self.USER_AGENTS)
return None # 继续传递请求
2.3 process_* 回调详解
每个 Middleware 可以定义以下方法,分别在不同阶段被调用:
| 方法 | 触发时机 | 返回值说明 |
|---|---|---|
process_request(req, spider) |
Request 经过中间件时 | None(继续) / Response(跳过下载) / Request(重定向) / IgnoreRequest(忽略) |
process_response(req, resp, spider) |
Response 返回时 | Response(继续) / Request(重试或重定向) |
process_exception(req, exc, spider) |
下载异常时 | None(继续异常) / Response(降级处理) / Request(重试) |
class RetryMiddleware:
"""带指数退避的重试中间件"""
def process_response(self, request, response, spider):
if response.status in (429, 503) and request.meta.get('retry_times', 0) < 3:
retry_times = request.meta.get('retry_times', 0) + 1
wait = 2 ** retry_times # 指数退避
spider.logger.warning(f'Retry {request.url} after {wait}s (attempt {retry_times})')
new_request = request.copy()
new_request.meta['retry_times'] = retry_times
new_request.meta['download_timeout'] = wait
return new_request
return response
三、Spider 开发模式
Scrapy 内置了多种 Spider 基类,适应不同的数据源结构。
3.1 基础 Spider
最通用的模式,手动定义 start_requests 和 parse:
import scrapy
from myproject.items import ArticleItem
class BlogSpider(scrapy.Spider):
name = 'blog'
start_urls = ['https://example.com/blog']
def parse(self, response):
for article in response.css('article.post'):
yield ArticleItem(
title=article.css('h2 a::text').get(),
url=article.css('h2 a::attr(href)').get(),
summary=article.css('p.summary::text').get(),
)
# 翻页
next_page = response.css('a.next::attr(href)').get()
if next_page:
yield response.follow(next_page, callback=self.parse)
3.2 CrawlSpider —— 自动链接跟踪
适用于站点内容采集,通过 Rule 自动发现和跟踪链接:
from scrapy.spiders import CrawlSpider, Rule
from scrapy.linkextractors import LinkExtractor
class NewsCrawlSpider(CrawlSpider):
name = 'news'
allowed_domains = ['news.example.com']
start_urls = ['https://news.example.com/']
rules = (
# 文章详情页 -> 解析内容
Rule(LinkExtractor(allow=r'/article/\d+'), callback='parse_article'),
# 列表页 -> 继续跟进(不回调)
Rule(LinkExtractor(allow=r'/page/\d+'), follow=True),
)
def parse_article(self, response):
yield {
'title': response.css('h1::text').get(),
'content': response.css('div.content::text').getall(),
'url': response.url,
}
3.3 XMLFeedSpider —— 结构化数据源
适合 RSS / Atom / XML Sitemap 等结构化数据:
from scrapy.spiders import XMLFeedSpider
class RssSpider(XMLFeedSpider):
name = 'rss'
start_urls = ['https://example.com/rss']
iterator = 'iternodes'
itertag = 'item'
def parse_node(self, response, node):
self.logger.info('Parsing RSS item: %s', node)
yield {
'title': node.xpath('title/text()').get(),
'link': node.xpath('link/text()').get(),
'pub_date': node.xpath('pubDate/text()').get(),
}
四、Item Pipeline —— 数据处理与持久化
Pipeline 是链式处理架构,每个 Item 会依次经过所有已启用的 Pipeline 组件。
4.1 启用 Pipeline
在 settings.py 中配置(数值越小优先级越高):
ITEM_PIPELINES = {
'myproject.pipelines.DuplicatesPipeline': 100,
'myproject.pipelines.ValidationPipeline': 200,
'myproject.pipelines.DatabasePipeline': 300,
}
4.2 实战:去重 + 验证 + MySQL 存储
# pipelines.py
import pymysql
from scrapy.exceptions import DropItem
from itemadapter import ItemAdapter
class DuplicatesPipeline:
"""基于集合的 URL 去重"""
def __init__(self):
self.urls_seen = set()
def process_item(self, item, spider):
adapter = ItemAdapter(item)
if adapter.get('url') in self.urls_seen:
raise DropItem(f'Duplicate item found: {adapter["url"]}')
self.urls_seen.add(adapter['url'])
return item
class ValidationPipeline:
"""数据清洗与校验"""
def process_item(self, item, spider):
adapter = ItemAdapter(item)
# 清洗:去除首尾空格
for field in adapter.keys():
if isinstance(adapter[field], str):
adapter[field] = adapter[field].strip()
# 校验:标题不能为空
if not adapter.get('title'):
raise DropItem('Missing title')
return item
class DatabasePipeline:
"""MySQL 持久化"""
def open_spider(self, spider):
self.conn = pymysql.connect(
host='localhost', user='root', password='secret',
database='scrapy_db', charset='utf8mb4',
)
self.cursor = self.conn.cursor()
def close_spider(self, spider):
self.conn.commit()
self.cursor.close()
self.conn.close()
def process_item(self, item, spider):
sql = """INSERT INTO articles (title, url, summary, created_at)
VALUES (%s, %s, %s, NOW())"""
self.cursor.execute(sql, (
item['title'], item['url'], item.get('summary', ''),
))
self.conn.commit()
return item
五、下载中间件实战
5.1 代理轮换中间件
class RotateProxyMiddleware:
"""IP 代理池轮换"""
def __init__(self, proxy_list):
self.proxy_list = proxy_list
self.current = 0
@classmethod
def from_crawler(cls, crawler):
proxies = crawler.settings.getlist('PROXY_LIST')
return cls(proxies)
def process_request(self, request, spider):
if not self.proxy_list:
return None
proxy = self.proxy_list[self.current % len(self.proxy_list)]
request.meta['proxy'] = proxy
self.current += 1
spider.logger.debug(f'Using proxy: {proxy}')
5.2 Cookie 会话管理
class SessionMiddleware:
"""维持 Session 会话"""
def process_request(self, request, spider):
# 首次请求获取 Cookie
if 'cookies' not in request.meta:
request.cookies = {'session_id': spider.session_id}
return None
def process_response(self, request, response, spider):
# 从响应中提取并更新 Cookie
set_cookie = response.headers.getlist('Set-Cookie')
if set_cookie:
spider.logger.debug(f'Updated cookies: {set_cookie}')
return response
5.3 在 settings.py 中启用
DOWNLOADER_MIDDLEWARES = {
'myproject.middlewares.RandomUserAgentMiddleware': 400,
'myproject.middlewares.RotateProxyMiddleware': 500,
'myproject.middlewares.SessionMiddleware': 600,
# 关闭默认的 UserAgentMiddleware
'scrapy.downloadermiddlewares.useragent.UserAgentMiddleware': None,
}
六、Extensions 与信号系统
6.1 信号机制
Scrapy 的信号系统允许你在爬虫生命周期关键节点挂载自定义逻辑:
| 信号 | 触发时机 |
|---|---|
engine_started |
Engine 启动 |
spider_opened |
Spider 打开 |
spider_closed |
Spider 关闭 |
item_scraped |
Item 被成功提取 |
item_dropped |
Item 被 Pipeline 丢弃 |
request_scheduled |
Request 被调度 |
response_received |
Response 被接收 |
6.2 自定义 Extension:统计与监控
# extensions.py
from scrapy import signals
from datetime import datetime
class StatsExtension:
"""爬虫运行状态统计扩展"""
def __init__(self, stats):
self.stats = stats
self.start_time = None
@classmethod
def from_crawler(cls, crawler):
ext = cls(crawler.stats)
crawler.signals.connect(ext.spider_opened, signal=signals.spider_opened)
crawler.signals.connect(ext.spider_closed, signal=signals.spider_closed)
crawler.signals.connect(ext.item_scraped, signal=signals.item_scraped)
crawler.signals.connect(ext.item_dropped, signal=signals.item_dropped)
return ext
def spider_opened(self, spider):
self.start_time = datetime.now()
spider.logger.info(f'Spider started: {spider.name}')
def spider_closed(self, spider, reason):
elapsed = datetime.now() - self.start_time
spider.logger.info(
f'Spider finished: {spider.name} | '
f'Items: {self.stats.get("item_scraped_count", 0)} | '
f'Dropped: {self.stats.get("item_dropped_count", 0)} | '
f'Duration: {elapsed}'
)
def item_scraped(self, item, spider):
self.stats.inc_value('custom_items_scraped')
def item_dropped(self, item, spider, exception):
self.stats.inc_value('custom_items_dropped')
在 settings.py 中启用:
EXTENSIONS = {
'myproject.extensions.StatsExtension': 500,
}
七、Scrapy-Redis 分布式爬虫架构
当单机无法满足数据量需求时,分布式是必然选择。Scrapy-Redis 通过 Redis 共享调度队列和去重集合,实现多机协同。
7.1 架构原理
+------------+ +------------+ +------------+
| Spider 节点 | | Spider 节点 | | Spider 节点 |
+-----+------+ +-----+------+ +-----+------+
| | |
+--------------------+---------------------+
|
+------v------+
| Redis |
| (队列+去重) |
+------+------+
|
+------v------+
| 共享调度队列 |
+-------------+
7.2 配置实现
# settings.py — 分布式配置
# 使用 Redis 调度器和去重组件
SCHEDULER = 'scrapy_redis.scheduler.Scheduler'
DUPEFILTER_CLASS = 'scrapy_redis.dupefilter.RFPDupeFilter'
# Redis 连接配置
REDIS_HOST = 'redis-master.example.com'
REDIS_PORT = 6379
REDIS_PARAMS = {
'password': 'your-redis-password',
'db': 0,
}
# 持久化调度队列(爬虫重启后不丢失)
SCHEDULER_PERSIST = True
# 队列模式:SpiderPriorityQueue(优先级)/ SpiderQueue(FIFO)/ SpiderStack(LIFO)
SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.SpiderPriorityQueue'
# Item 直接推入 Redis 供消费者处理
ITEM_PIPELINES = {
'scrapy_redis.pipelines.RedisPipeline': 300,
}
7.3 Redis Spider 示例
from scrapy_redis.spiders import RedisSpider
class DistributedBlogSpider(RedisSpider):
"""从 Redis 获取起始 URL,支持多机协同"""
name = 'distributed_blog'
redis_key = 'blog:start_urls'
def parse(self, response):
# 与普通 Spider 相同
for article in response.css('article'):
yield {'title': article.css('h2::text').get()}
启动方式:
# 在所有节点上运行
scrapy crawl distributed_blog
# 向 Redis 推送起始 URL(任意一台机器执行一次即可)
redis-cli lpush blog:start_urls 'https://example.com/blog'
7.4 去重机制深入
Redis 去重基于 RFPDupeFilter,默认以请求的 方法 + URL + Body 生成指纹并存入 Redis Set。你可以通过重写 request_fingerprint 实现自定义去重:
from scrapy.utils.request import request_fingerprint
def custom_fingerprint(request):
"""忽略查询参数的指纹(同一页面不同参数视为重复)"""
from w3lib.url import url_query_cleaner
cleaned_url = url_query_cleaner(request.url, parameterlist=['utm_source', 'utm_campaign'])
new_req = request.replace(url=cleaned_url)
return request_fingerprint(new_req)
八、性能优化
8.1 并发控制参数
# settings.py
# 并发请求数(默认 16)
CONCURRENT_REQUESTS = 32
# 每个域名最大并发(防封禁)
CONCURRENT_REQUESTS_PER_DOMAIN = 8
# 每个 IP 最大并发(使用代理时)
CONCURRENT_REQUESTS_PER_IP = 0 # 0 表示不限制
# 下载延迟(秒)
DOWNLOAD_DELAY = 1.0
# 自动限速扩展(根据服务器响应动态调整)
AUTOTHROTTLE_ENABLED = True
AUTOTHROTTLE_START_DELAY = 1.0
AUTOTHROTTLE_MAX_DELAY = 30.0
AUTOTHROTTLE_TARGET_CONCURRENCY = 8.0
8.2 请求去重优化
# 使用布隆过滤器减少内存占用
# 先安装:pip install scrapy-bloomfilter
DUPEFILTER_CLASS = 'scrapy_bloomfilter.dupefilter.BloomFilterDupeFilter'
BLOOMFILTER_CAPACITY = 10_000_000 # 一千万容量
BLOOMFILTER_ERROR_RATE = 0.001 # 千分之一误判率
8.3 连接复用与 Keep-Alive
# 启用 HTTP 连接池复用
DOWNLOADER_CLIENT_TLS_METHOD = 'TLSv1.2'
DOWNLOADER_CLIENT_TLS_VERBOSE_LOGGING = False
# Twisted 连接池配置
CONNECTION_POOL_SIZE = 64
8.4 内存优化
# 限制 Item 在内存中的数量(避免大数据量时 OOM)
# Pipeline 处理完一批后释放
ITEM_PIPELINES = {
'myproject.pipelines.BulkDatabasePipeline': 300,
}
# 配合批量写入
class BulkDatabasePipeline:
BULK_SIZE = 500
def open_spider(self, spider):
self.items = []
def process_item(self, item, spider):
self.items.append(dict(item))
if len(self.items) >= self.BULK_SIZE:
self._flush()
return item
def close_spider(self, spider):
if self.items:
self._flush()
def _flush(self):
# 批量写入数据库
self.cursor.executemany(
'INSERT INTO articles (title, url) VALUES (%(title)s, %(url)s)',
self.items,
)
self.conn.commit()
self.items = []
九、实战最佳实践
9.1 结构化异常处理
class SafeSpider(scrapy.Spider):
name = 'safe_spider'
def parse(self, response):
try:
# 使用 extract_first 提供默认值,避免 None 传播
title = response.css('h1::text').get(default='').strip()
if not title:
self.logger.warning(f'Empty title at {response.url}')
return
yield {'title': title, 'url': response.url}
except Exception as e:
self.logger.error(
f'Parse error on {response.url}: {e}',
exc_info=True,
)
# 可以将异常 URL 记录到文件以便后续重试
with open('failed_urls.txt', 'a') as f:
f.write(response.url + '\n')
9.2 日志系统配置
# settings.py — 日志配置
LOG_ENABLED = True
LOG_LEVEL = 'INFO' # DEBUG / INFO / WARNING / ERROR / CRITICAL
LOG_FILE = 'logs/scrapy.log'
LOG_FORMAT = '%(asctime)s [%(name)s] %(levelname)s: %(message)s'
LOG_DATEFORMAT = '%Y-%m-%d %H:%M:%S'
# 自定义日志轮转
LOG_ENCODING = 'utf-8'
9.3 结合 Prometheus + Grafana 监控
# prometheus_extension.py
from scrapy import signals
from prometheus_client import Counter, Gauge, start_http_server
class PrometheusExporter:
"""导出 Scrapy 指标到 Prometheus"""
def __init__(self, port=9100):
self.items_scraped = Counter('scrapy_items_total', 'Total scraped items')
self.items_dropped = Counter('scrapy_items_dropped', 'Total dropped items')
self.active_requests = Gauge('scrapy_active_requests', 'Active requests')
start_http_server(port)
@classmethod
def from_crawler(cls, crawler):
ext = cls()
crawler.signals.connect(ext.item_scraped, signal=signals.item_scraped)
crawler.signals.connect(ext.item_dropped, signal=signals.item_dropped)
return ext
def item_scraped(self, item, spider):
self.items_scraped.inc()
def item_dropped(self, item, spider, exception):
self.items_dropped.inc()
9.4 快速故障恢复
# 意外中断后断点续爬
# 在 settings.py 中启用:
# 持久化调度队列(爬虫重启后继续未完成的请求)
SCHEDULER_PERSIST = True
# 默认情况下 Scrapy 会清除去重记录,保留以便后续续爬
DUPEFILTER_CLASS = 'scrapy_redis.dupefilter.RFPDupeFilter'
# 或者使用本地 JSON 保存已爬取 URL
import json
class PersistenceExtension:
def spider_opened(self, spider):
try:
with open('crawled_urls.json') as f:
spider.crawled_urls = set(json.load(f))
except FileNotFoundError:
spider.crawled_urls = set()
def spider_closed(self, spider):
with open('crawled_urls.json', 'w') as f:
json.dump(list(spider.crawled_urls), f)
9.5 反爬虫策略总结
| 策略 | 实现方式 |
|---|---|
| IP 封禁 | 代理池 + 自动轮换 |
| User-Agent 检测 | 随机 UA 中间件 |
| Cookie 验证 | Session 维持 + 登录态管理 |
| 访问频率限制 | DOWNLOAD_DELAY + AutoThrottle |
| 页面渲染检测 | Splash / Selenium 中间件 |
| 验证码 | 打码平台 API 集成 |
| 签名参数 | 逆向 JS + process_request 注入 |
十、总结与进阶路线
Scrapy 是一个设计精良的框架,掌握其核心组件和扩展机制后,你可以:
- 快速开发:利用 Spider 基类和内置组件快速搭建爬虫
- 深度定制:通过 Middleware、Pipeline、Extensions 实现任意自定义逻辑
- 横向扩展:引入 Scrapy-Redis 实现多机分布式采集
- 生产可用:配合日志、监控、异常处理构建企业级采集系统
进阶推荐阅读:
- Scrapy 官方文档
- Scrapy-Redis 源码分析
- 《Python 网络数据采集》—— Ryan Mitchell
- Twisted 异步编程模型
采集是手段,不是目的。请尊重目标网站的
robots.txt和使用条款,合理控制采集频率,做负责任的爬虫工程师。
本文由团队内部技术分享整理而成,欢迎在评论区交流讨论。
评论