Appearance
章节 4:Scrapy 框架与分布式
学习目标
- 理解 Scrapy 五大核心组件的架构与数据流
- 掌握 Spider 编写、start_requests 与 parse 回调机制
- 实现 Item Pipeline 进行数据清洗与多存储
- 编写 Downloader Middleware 拦截/修改请求
- 使用 Scrapy-Redis 实现分布式爬虫(去重、任务分发、增量爬取)
4.1 Scrapy 框架概述
4.1.1 什么是 Scrapy
定义:Scrapy 是一个基于 Twisted 异步网络引擎的高性能爬虫框架,提供了从请求调度到数据存储的完整解决方案。
4.1.2 五大核心组件
┌──────────────────────────┐
│ Scrapy Engine │ ← 引擎(核心调度)
│ (数据流控制中枢) │
└───┬─────┬──────┬──────┬──┘
┌─────────── │ │ │ │ ────────────┐
▼ ▼ ▼ ▼ ▼ ▼
┌──────────┐ ┌────────┐ ┌──────┐ ┌──────────┐ ┌──────────┐
│ Scheduler│ │Downloader│ │Spider│ │Item │ │ │
│ 调度器 │ │ 下载器 │ │爬虫 │ │Pipeline │ │Extensions│
└──────────┘ └────────┘ └──────┘ └──────────┘ └──────────┘
│ │ │ │
▼ ▼ ▼ ▼
请求队列 下载页面 解析响应 数据存储| 组件 | 职责 |
|---|---|
| Engine(引擎) | 控制数据流,触发事件,协调各组件工作 |
| Scheduler(调度器) | 管理请求队列,去重,决定下一个要下载的请求 |
| Downloader(下载器) | 执行 HTTP 请求,返回 Response 给 Spider |
| Spider(爬虫) | 解析 Response,提取数据(Item)或生成新请求 |
| Item Pipeline(管道) | 清洗、验证、存储 Item(支持多个串联) |
数据流(五大步骤):
1. Engine 从 Spider 获取初始 Request
2. Engine 将 Request 交给 Scheduler 排队
3. Scheduler 出队 → Engine → Downloader 执行下载
4. Downloader 返回 Response → Engine → Spider.parse()
5. Spider 解析出 Item → Engine → Item Pipeline
或 解析出新 Request → Engine → Scheduler(循环 1-5)4.1.3 项目结构
bash
# 创建项目
scrapy startproject myproject
# 项目结构
myproject/
├── scrapy.cfg # 项目配置文件
└── myproject/
├── __init__.py
├── items.py # Item 定义(数据结构)
├── middlewares.py # Middleware 实现
├── pipelines.py # Pipeline 实现
├── settings.py # 全局设置
└── spiders/ # Spider 目录
├── __init__.py
└── example_spider.py4.2 Spider 编写
4.2.1 基础 Spider
python
# myproject/spiders/quotes_spider.py
import scrapy
class QuotesSpider(scrapy.Spider):
name = "quotes" # Spider 名称,用于启动
allowed_domains = ["quotes.toscrape.com"] # 允许的域名
start_urls = ["https://quotes.toscrape.com/"] # 起始 URL
def parse(self, response):
"""默认回调函数,处理 start_urls 的响应"""
# 提取每条名言
for quote_block in response.css("div.quote"):
yield {
"text": quote_block.css("span.text::text").get(),
"author": quote_block.css("small.author::text").get(),
"tags": quote_block.css("div.tags a.tag::text").getall(),
}
# 翻页:找到"下一页"链接继续爬取
next_page = response.css("li.next a::attr(href)").get()
if next_page:
yield response.follow(next_page, callback=self.parse)运行爬虫:
bash
# 运行爬虫,输出到 JSON
scrapy crawl quotes -o quotes.json
# 输出到 CSV
scrapy crawl quotes -o quotes.csv
# 输出到 JSON Lines
scrapy crawl quotes -o quotes.jl4.2.2 start_requests 方法
定义:start_requests() 方法替代 start_urls 列表,更灵活地生成初始请求。
python
import scrapy
class CustomSpider(scrapy.Spider):
name = "custom"
def start_requests(self):
"""自定义初始请求——支持 Headers、Cookie、优先级等"""
urls = [
"https://example.com/page/1",
"https://example.com/page/2",
]
for url in urls:
yield scrapy.Request(
url=url,
headers={"Custom-Header": "value"},
cookies={"session": "abc123"},
callback=self.parse_page,
priority=10, # 优先级(越大越先处理)
meta={"page": 1}, # 传递额外数据
dont_filter=True, # 不过滤重复 URL
)
def parse_page(self, response):
page = response.meta.get("page")
yield {"page": page, "url": response.url}4.2.3 parse 回调与数据传递
python
import scrapy
class DetailSpider(scrapy.Spider):
name = "details"
start_urls = ["https://books.toscrape.com/"]
def parse(self, response):
"""列表页:提取每个商品链接,发起详情页请求"""
for book in response.css("article.product_pod"):
detail_url = book.css("h3 a::attr(href)").get()
# 传递 meta 数据到下一个回调
yield response.follow(
detail_url,
callback=self.parse_detail,
meta={
"title": book.css("h3 a::attr(title)").get(),
"price": book.css("p.price_color::text").get(),
},
)
# 翻页
next_page = response.css("li.next a::attr(href)").get()
if next_page:
yield response.follow(next_page, callback=self.parse)
def parse_detail(self, response):
"""详情页:提取完整信息"""
title = response.meta["title"]
price = response.meta["price"]
yield {
"title": title,
"price": price,
"description": response.css(
"#product_description ~ p::text"
).get(),
"upc": response.css("table tr:nth-child(1) td::text").get(),
}4.2.4 Spider 参数传递
bash
# 运行时传递参数
scrapy crawl quotes -a category=tech -a page_count=5python
import scrapy
class ParamSpider(scrapy.Spider):
name = "param_spider"
def __init__(self, category=None, page_count="3", *args, **kwargs):
super().__init__(*args, **kwargs)
self.category = category
self.page_count = int(page_count)
def start_requests(self):
for page in range(1, self.page_count + 1):
url = f"https://example.com/{self.category}?page={page}"
yield scrapy.Request(url, callback=self.parse)4.3 Item Pipeline —— 数据存储与清洗
4.3.1 定义 Item
定义:Item 是 Scrapy 中定义结构化数据字段的容器,类似于字典但有类型声明。
python
# myproject/items.py
import scrapy
class ProductItem(scrapy.Item):
# 定义字段
name = scrapy.Field()
price = scrapy.Field()
description = scrapy.Field()
url = scrapy.Field()
created_at = scrapy.Field()
category = scrapy.Field()4.3.2 编写 Pipeline
定义:Item Pipeline 是处理 Spider 产出 Item 的组件,多个 Pipeline 按优先级顺序执行。
python
# myproject/pipelines.py
import json
from datetime import datetime
from itemadapter import ItemAdapter
class PriceCleanPipeline:
"""Pipeline 1:数据清洗——价格格式化"""
def process_item(self, item, spider):
adapter = ItemAdapter(item)
# 清洗价格字段:去除货币符号,转为浮点数
if adapter.get("price"):
price_str = adapter["price"]
# 去除 ¥、$、, 等符号
price_clean = price_str.replace("¥", "") \
.replace("$", "") \
.replace(",", "") \
.strip()
try:
adapter["price"] = float(price_clean)
except ValueError:
adapter["price"] = 0.0
# 添加采集时间
adapter["created_at"] = datetime.now().isoformat()
return item
class DuplicatesPipeline:
"""Pipeline 2:去重"""
def __init__(self):
self.urls_seen = set()
def process_item(self, item, spider):
adapter = ItemAdapter(item)
if adapter["url"] in self.urls_seen:
raise DropItem(f"重复数据: {adapter['url']}")
else:
self.urls_seen.add(adapter["url"])
return item
class JsonWriterPipeline:
"""Pipeline 3:写入 JSON 文件"""
def open_spider(self, spider):
"""爬虫启动时调用"""
self.file = open("products.jsonl", "w", encoding="utf-8")
def close_spider(self, spider):
"""爬虫关闭时调用"""
self.file.close()
def process_item(self, item, spider):
line = json.dumps(ItemAdapter(item).asdict(), ensure_ascii=False) + "\n"
self.file.write(line)
return item
class MySQLPipeline:
"""Pipeline 4:写入 MySQL(示例)"""
def __init__(self, db_config):
self.db_config = db_config
@classmethod
def from_crawler(cls, crawler):
"""从 settings.py 读取配置"""
return cls(db_config=crawler.settings.getdict("MYSQL_CONFIG"))
def open_spider(self, spider):
"""连接数据库"""
import pymysql
self.conn = pymysql.connect(**self.db_config)
self.cursor = self.conn.cursor()
def close_spider(self, spider):
"""关闭连接"""
self.cursor.close()
self.conn.close()
def process_item(self, item, spider):
adapter = ItemAdapter(item)
sql = """INSERT INTO products (name, price, url, created_at)
VALUES (%s, %s, %s, %s)"""
self.cursor.execute(sql, (
adapter["name"],
adapter["price"],
adapter["url"],
adapter["created_at"],
))
self.conn.commit()
return item4.3.3 注册 Pipeline
python
# myproject/settings.py
# 按优先级注册 Pipeline(数字越小越先执行)
ITEM_PIPELINES = {
"myproject.pipelines.PriceCleanPipeline": 100, # 先清洗
"myproject.pipelines.DuplicatesPipeline": 200, # 再去重
"myproject.pipelines.JsonWriterPipeline": 300, # 再保存
# "myproject.pipelines.MySQLPipeline": 400,
}
# Pipeline 的自定义配置
MYSQL_CONFIG = {
"host": "localhost",
"port": 3306,
"user": "root",
"password": "password",
"database": "scrapy_db",
"charset": "utf8mb4",
}4.4 Downloader Middleware —— 请求拦截
定义:Downloader Middleware 位于 Engine 和 Downloader 之间,可以对请求进行预处理(如添加代理、修改 Headers)或对响应进行后处理。
4.4.1 代理中间件
python
# myproject/middlewares.py
import random
class RandomProxyMiddleware:
"""随机代理中间件"""
def __init__(self, proxies):
self.proxies = proxies
@classmethod
def from_crawler(cls, crawler):
proxies = crawler.settings.getlist("PROXY_LIST")
return cls(proxies)
def process_request(self, request, spider):
"""每个请求前调用——设置代理"""
proxy = random.choice(self.proxies)
request.meta["proxy"] = proxy
spider.logger.info(f"使用代理: {proxy}")
def process_exception(self, request, exception, spider):
"""请求异常时——切换到下一个代理"""
spider.logger.warning(f"代理失败: {request.meta.get('proxy')}")
# 移除失败的代理,重新尝试
request.meta["proxy"] = random.choice(self.proxies)
return request4.4.2 User-Agent 轮换中间件
python
import random
class UserAgentMiddleware:
"""UA 轮换中间件"""
USER_AGENTS = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
"Chrome/120.0.0.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) "
"AppleWebKit/537.36 Chrome/120.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:109.0) "
"Gecko/20100101 Firefox/121.0",
"Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) "
"AppleWebKit/605.1.15 Mobile/15E148",
]
def process_request(self, request, spider):
ua = random.choice(self.USER_AGENTS)
request.headers["User-Agent"] = ua
def process_response(self, request, response, spider):
"""响应处理——检测反爬"""
if response.status in [403, 429]:
spider.logger.warning(
f"被限制: {response.status} - {request.url}"
)
# 可以返回 Request 对象让下载器重试
new_request = request.copy()
new_request.headers["User-Agent"] = random.choice(self.USER_AGENTS)
return new_request
return response4.4.3 注册 Middleware
python
# myproject/settings.py
DOWNLOADER_MIDDLEWARES = {
"myproject.middlewares.UserAgentMiddleware": 400,
"myproject.middlewares.RandomProxyMiddleware": 500,
# Scrapy 内置的 Retry Middleware
"scrapy.downloadermiddlewares.retry.RetryMiddleware": 550,
}
# 代理列表
PROXY_LIST = [
"http://proxy1.example.com:8080",
"http://proxy2.example.com:8080",
"http://proxy3.example.com:8080",
]4.5 Scrapy-Redis —— 分布式爬虫
4.5.1 分布式原理
定义:Scrapy-Redis 将 Scrapy 的 Scheduler(调度器)和 Dupefilter(去重过滤器)从内存迁移到 Redis,使得多个 Scrapy 进程可以共享请求队列和去重集合,从而实现分布式爬虫。
┌──────────────┐
│ Redis │
│ │
│ ┌──────────┐ │
│ │请求队列 │ │ ← 多个爬虫共享
│ ├──────────┤ │
│ │去重集合 │ │ ← 全局去重
│ ├──────────┤ │
│ │Item 集合 │ │ ← 可选统计
│ └──────────┘ │
└──────────────┘
▲ ▲
│ │
┌─────────────┘ └─────────────┐
▼ ▼
┌──────────────┐ ┌──────────────┐
│ Scrapy 进程1 │ │ Scrapy 进程2 │
│ (爬虫节点) │ │ (爬虫节点) │
└──────────────┘ └──────────────┘4.5.2 安装与配置
bash
pip install scrapy-redispython
# settings.py —— 分布式配置
# 使用 Scrapy-Redis 的调度器
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 使用 Scrapy-Redis 的去重过滤器
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# Redis 连接配置
REDIS_URL = "redis://localhost:6379/0"
# 持久化——爬虫结束后不清空请求队列
SCHEDULER_PERSIST = True
# 请求队列模式
# FifoQueue: 先进先出(广度优先)
# LIFOQueue: 后进先出(深度优先)
# PriorityQueue: 优先级队列(默认)
SCHEDULER_QUEUE_CLASS = "scrapy_redis.queue.PriorityQueue"
# 可选的:将采集到的 Items 也存到 Redis
ITEM_PIPELINES = {
"scrapy_redis.pipelines.RedisPipeline": 300,
}4.5.3 分布式 Spider
python
# myproject/spiders/redis_spider.py
from scrapy_redis.spiders import RedisSpider
class MyRedisSpider(RedisSpider):
"""继承 RedisSpider,从 Redis 获取起始 URL"""
name = "my_redis_spider"
# Redis key,爬虫会从这个 key 中 pop 起始 URL
redis_key = "my_redis_spider:start_urls"
def parse(self, response):
# 解析逻辑与普通 Spider 相同
title = response.css("title::text").get()
yield {"url": response.url, "title": title}
# 提取新链接
for link in response.css("a::attr(href)").getall():
yield response.follow(link, callback=self.parse)启动方式:
bash
# 1. 在 Redis 中注入起始 URL
redis-cli lpush my_redis_spider:start_urls "https://example.com"
# 2. 在所有爬虫节点上启动(数量任意)
scrapy crawl my_redis_spider
# 3. 随时添加更多 URL 到队列
redis-cli lpush my_redis_spider:start_urls "https://example.com/page2"4.5.4 Redis 去重原理
定义:Scrapy-Redis 使用 Redis 的 Set 数据结构存储已发送请求的指纹(fingerprint),实现全局去重。
去重流程:
1. 收到新 Request
2. 计算请求指纹:SHA1(method + url + body + headers)
3. 检查指纹是否在 Redis Set 中
→ 已存在:丢弃(重复)
→ 不存在:加入 Set,加入请求队列手动管理去重:
python
# 查看去重集合大小
redis-cli SCARD my_redis_spider:dupefilter
# 清除去重集合(重新爬取)
redis-cli DEL my_redis_spider:dupefilter
# 查看等待队列大小
redis-cli LLEN my_redis_spider:requests
# 清空队列
redis-cli DEL my_redis_spider:requests4.5.5 增量爬取
定义:增量爬取只采集新数据或已更新数据,避免重复下载。结合 Scrapy-Redis 的持久化特性,可以实现高效的增量采集。
策略一:基于时间戳:
python
import scrapy
from datetime import datetime
class IncrementalSpider(scrapy.Spider):
name = "incremental"
start_urls = ["https://example.com/news"]
def parse(self, response):
# 只采集 24 小时内发布的文章
for article in response.css("article"):
pub_time_str = article.css("time::attr(datetime)").get()
if pub_time_str:
pub_time = datetime.fromisoformat(pub_time_str)
if (datetime.now() - pub_time).days > 1:
continue # 跳过旧文章
yield {
"title": article.css("h2::text").get(),
"url": article.css("a::attr(href)").get(),
"pub_time": pub_time_str,
}策略二:基于 Redis 去重(分布式增量):
python
# 分布式增量爬虫——利用 Redis 去重自动过滤已爬 URL
# 只需要配置 SCHEDULER_PERSIST = True
# 爬虫停止后去重集合保留,下次启动不会重复爬取
# 如果要去除旧的去重记录重新爬取:
# redis-cli DEL my_redis_spider:dupefilter4.6 settings.py 完整配置参考
python
# myproject/settings.py
# === 基础设置 ===
BOT_NAME = "myproject"
SPIDER_MODULES = ["myproject.spiders"]
NEWSPIDER_MODULE = "myproject.spiders"
ROBOTSTXT_OBEY = True
# === 并发与延迟 ===
CONCURRENT_REQUESTS = 16 # 最大并发数
CONCURRENT_REQUESTS_PER_DOMAIN = 8 # 单域名并发数
DOWNLOAD_DELAY = 0.5 # 请求间隔(秒)
RANDOMIZE_DOWNLOAD_DELAY = True # 随机化延迟 (±50%)
# === 超时与重试 ===
DOWNLOAD_TIMEOUT = 15 # 下载超时(秒)
RETRY_ENABLED = True
RETRY_TIMES = 2 # 重试次数
RETRY_HTTP_CODES = [500, 502, 503, 504, 408, 429]
# === 缓存与日志 ===
HTTPCACHE_ENABLED = True # HTTP 缓存(调试用)
HTTPCACHE_EXPIRATION_SECS = 3600
HTTPCACHE_DIR = "httpcache"
LOG_LEVEL = "INFO" # DEBUG, INFO, WARNING, ERROR
LOG_FILE = "crawl.log" # 日志文件(不设则输出到控制台)
# === AutoThrottle(自动限速) ===
AUTOTHROTTLE_ENABLED = True
AUTOTHROTTLE_START_DELAY = 1.0
AUTOTHROTTLE_MAX_DELAY = 10.0
AUTOTHROTTLE_TARGET_CONCURRENCY = 4.0
# === Pipeline ===
ITEM_PIPELINES = {
"myproject.pipelines.PriceCleanPipeline": 100,
"myproject.pipelines.JsonWriterPipeline": 300,
}
# === Middleware ===
DOWNLOADER_MIDDLEWARES = {
"myproject.middlewares.UserAgentMiddleware": 400,
}
# === 分布式 ===
# SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# REDIS_URL = "redis://localhost:6379/0"
# SCHEDULER_PERSIST = True小结
- Scrapy 架构:五大组件(Engine/Scheduler/Downloader/Spider/Pipeline)通过 Twisted 异步引擎协作,数据流闭环循环
- Spider 编写:
start_requests()生成初始请求,parse()是默认回调,response.follow()实现翻页 - Item Pipeline:按优先级串联处理,
process_item()清洗/去重/存储,open_spider()和close_spider()管理连接 - Downloader Middleware:在请求发送前拦截修改(代理、UA、Cookie),在响应返回后处理(重试、检测反爬)
- Scrapy-Redis 分布式:将调度器和去重器迁移到 Redis,多个爬虫节点共享请求队列,实现分布式+增量爬取
练习
- 创建项目:使用
scrapy startproject创建一个新的 Scrapy 项目,理解项目目录结构。 - 基础 Spider:编写一个 Spider 爬取
https://quotes.toscrape.com/的所有名言,包括翻页,输出到 JSON 文件。 - Item Pipeline:编写两个 Pipeline——一个将价格字符串转为浮点数,一个将结果写入 SQLite 数据库。在 settings.py 中注册并测试。
- Downloader Middleware:编写一个随机 User-Agent 中间件和一个随机代理中间件,验证请求头中的 User-Agent 是否每次不同。
- 分布式实战:配置 Scrapy-Redis,启动两个爬虫节点(在本地用两个终端窗口运行),通过 Redis 注入起始 URL,观察两个节点如何协同消费队列。