Skip to content

章节 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.py

4.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.jl

4.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=5
python
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 item

4.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 request

4.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 response

4.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-redis
python
# 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:requests

4.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:dupefilter

4.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

小结

  1. Scrapy 架构:五大组件(Engine/Scheduler/Downloader/Spider/Pipeline)通过 Twisted 异步引擎协作,数据流闭环循环
  2. Spider 编写start_requests() 生成初始请求,parse() 是默认回调,response.follow() 实现翻页
  3. Item Pipeline:按优先级串联处理,process_item() 清洗/去重/存储,open_spider()close_spider() 管理连接
  4. Downloader Middleware:在请求发送前拦截修改(代理、UA、Cookie),在响应返回后处理(重试、检测反爬)
  5. Scrapy-Redis 分布式:将调度器和去重器迁移到 Redis,多个爬虫节点共享请求队列,实现分布式+增量爬取

练习

  1. 创建项目:使用 scrapy startproject 创建一个新的 Scrapy 项目,理解项目目录结构。
  2. 基础 Spider:编写一个 Spider 爬取 https://quotes.toscrape.com/ 的所有名言,包括翻页,输出到 JSON 文件。
  3. Item Pipeline:编写两个 Pipeline——一个将价格字符串转为浮点数,一个将结果写入 SQLite 数据库。在 settings.py 中注册并测试。
  4. Downloader Middleware:编写一个随机 User-Agent 中间件和一个随机代理中间件,验证请求头中的 User-Agent 是否每次不同。
  5. 分布式实战:配置 Scrapy-Redis,启动两个爬虫节点(在本地用两个终端窗口运行),通过 Redis 注入起始 URL,观察两个节点如何协同消费队列。

Python 学习资料