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 爬取流程如下:

  1. Engine 从 Spider 获取初始 Request,交给 Scheduler
  2. Scheduler 将 Request 入队,Engine 轮询获取下一个 Request
  3. Engine 将 Request 交给 Downloader Middleware 链,再由 Downloader 执行 HTTP 请求
  4. Downloader 返回 Response,经 Downloader Middleware 处理后交给 Engine
  5. Engine 将 Response 交给 Spider Middleware 链,最终进入 Spider 解析
  6. 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_requestsparse

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 是一个设计精良的框架,掌握其核心组件和扩展机制后,你可以:

  1. 快速开发:利用 Spider 基类和内置组件快速搭建爬虫
  2. 深度定制:通过 Middleware、Pipeline、Extensions 实现任意自定义逻辑
  3. 横向扩展:引入 Scrapy-Redis 实现多机分布式采集
  4. 生产可用:配合日志、监控、异常处理构建企业级采集系统

进阶推荐阅读:

采集是手段,不是目的。请尊重目标网站的 robots.txt 和使用条款,合理控制采集频率,做负责任的爬虫工程师。


本文由团队内部技术分享整理而成,欢迎在评论区交流讨论。