这是《Scrapy框架实战教程》的下篇。如果你还没读过上篇,建议先看上篇,那里有Scrapy的基础概念和第一个实战案例。

        上篇我们完成了从0到1的入门,搭建了一个可以跑通的基本爬虫。但这只是开始。在实际项目中,你会遇到各种挑战:

        • 反爬机制越来越强,需要随机UA、代理IP、动态渲染

        • 单机性能有限,需要分布式架构

        • 数据规模扩大,需要高效的存储和监控

        • 代码需要维护,需要模块化和扩展性

这篇下篇,我会带你深入Scrapy的高级功能,解决这些真实场景的问题。

高级功能一:中间件

        中间件是Scrapy最强大的扩展机制。简单说,它可以在请求发送前、响应返回后插入自定义逻辑,像一个个"过滤器"处理所有进出爬虫的请求和响应。

中间件的工作原理

中间件分为两种:

        • 下载中间件(Downloader Middleware): 处理请求和响应

        • 爬虫中间件(Spider Middleware): 处理爬虫输出的Item和Request

我们主要讲下载中间件,因为它更常用。

请求流程:
引擎 → 中间件(入站) → 下载器 → 中间件(出站) → 爬虫
响应流程:
爬虫 → 中间件(入站) → 引擎 → 中间件(出站) → 下载器

每个中间件必须实现三个方法:

        • process_request(request, spider): 请求发送前调用

        • process_response(request, response, spider): 响应返回后调用

        • process_exception(request, exception, spider): 请求异常时调用

实战1:随机UA中间件

最基础的反爬是User-Agent检测。我们来写一个随机UA中间件。

# middlewares.py
import random
from scrapy import signals
class RandomUserAgentMiddleware:
    """随机User-Agent中间件"""
    def __init__(self, user_agent_list):
        self.user_agent_list = user_agent_list
    @classmethod
    def from_crawler(cls, crawler):
        """从settings读取配置"""
        return cls(
            user_agent_list=crawler.settings.get('USER_AGENT_LIST')
        )
    def process_request(self, request, spider):
        """为每个请求随机分配UA"""
        request.headers['User-Agent'] = random.choice(self.user_agent_list)
        spider.logger.debug(f"User-Agent: {request.headers['User-Agent']}")

在settings.py中配置:

# 启用中间件
DOWNLOADER_MIDDLEWARES = {
    'book_scraper.middlewares.RandomUserAgentMiddleware': 400,
}
# UA列表
USER_AGENT_LIST = [
    'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
    'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
    'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:121.0) Gecko/20100101 Firefox/121.0',
    'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.1 Safari/605.1.15',
    'Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Mobile/15E148 Safari/604.1',
]

实战2:代理IP池中间件

当IP被封时,需要自动切换代理。我们来实现一个带重试的代理池中间件。

# middlewares.py
import random
from scrapy import signals
from scrapy.exceptions import NotConfigured
class ProxyPoolMiddleware:
    """代理IP池中间件"""
    def __init__(self, proxy_list):
        self.proxy_list = proxy_list
        if not self.proxy_list:
            raise NotConfigured("Proxy list is empty")
    @classmethod
    def from_crawler(cls, crawler):
        return cls(
            proxy_list=crawler.settings.get('PROXY_LIST', [])
        )
    def process_request(self, request, spider):
        """为每个请求随机分配代理"""
        proxy = random.choice(self.proxy_list)
        request.meta['proxy'] = proxy
        spider.logger.debug(f"Using proxy: {proxy}")
    def process_exception(self, request, exception, spider):
        """请求异常时记录"""
        spider.logger.warning(f"Proxy failed: {request.meta.get('proxy')}, Exception: {exception}")

配置settings.py:

DOWNLOADER_MIDDLEWARES = {
    'book_scraper.middlewares.ProxyPoolMiddleware': 410,
}
# 代理列表(实际项目中建议从API动态获取)
PROXY_LIST = [
    'http://proxy1.example.com:8080',
    'http://proxy2.example.com:8080',
    'http://proxy3.example.com:8080',
]

实战3:重试中间件

当请求失败时自动重试,提高成功率。

# middlewares.py
from scrapy.downloadermiddlewares.retry import RetryMiddleware
from scrapy.utils.response import response_status_message
class CustomRetryMiddleware(RetryMiddleware):
    """自定义重试中间件"""
    def __init__(self, settings):
        super().__init__(settings)
        self.max_retry_times = settings.getint('RETRY_TIMES', 3)
        self.retry_http_codes = set(int(x) for x in settings.getlist('RETRY_HTTP_CODES'))
    def process_response(self, request, response, spider):
        """检查响应状态码,需要重试时重新请求"""
        if response.status in self.retry_http_codes:
            reason = response_status_message(response.status)
            return self._retry(request, reason, spider) or response
        return response

配置:

# 重试设置
RETRY_ENABLED = True
RETRY_TIMES = 3
RETRY_HTTP_CODES = [500, 502, 503, 504, 408, 429]
# 启用重试中间件(默认已启用,优先级越高越先执行)
DOWNLOADER_MIDDLEWARES = {
    'scrapy.downloadermiddlewares.retry.RetryMiddleware': 90,
    'book_scraper.middlewares.ProxyPoolMiddleware': 410,
}

实战4:动态渲染中间件

很多网站使用JavaScript动态加载内容,普通的HTTP请求无法获取完整数据。这时需要用Selenium或Playwright渲染页面。

我们用Playwright(比Selenium更快、更稳定):

# middlewares.py
from scrapy.http import HtmlResponse
from playwright.sync_api import sync_playwright
class PlaywrightMiddleware:
    """Playwright动态渲染中间件"""
    def __init__(self):
        self.playwright = None
        self.browser = None
    @classmethod
    def from_crawler(cls, crawler):
        middleware = cls()
        crawler.signals.connect(middleware.spider_opened, signals.spider_opened)
        crawler.signals.connect(middleware.spider_closed, signals.spider_closed)
        return middleware
    def spider_opened(self, spider):
        """爬虫启动时初始化Playwright"""
        self.playwright = sync_playwright().start()
        self.browser = self.playwright.chromium.launch(headless=True)
    def spider_closed(self, spider):
        """爬虫关闭时清理资源"""
        if self.browser:
            self.browser.close()
        if self.playwright:
            self.playwright.stop()
    def process_request(self, request, spider):
        """使用Playwright渲染页面"""
        # 只对需要渲染的请求启用
        if not request.meta.get('render'):
            return None
        page = self.browser.new_page()
        try:
            page.goto(request.url, wait_until='networkidle')
            # 等待特定元素出现
            if request.meta.get('wait_for'):
                page.wait_for_selector(request.meta['wait_for'])
            # 获取渲染后的HTML
            html = page.content()
            # 返回HtmlResponse,后续解析流程不变
            return HtmlResponse(
                url=request.url,
                body=html.encode('utf-8'),
                encoding='utf-8',
                request=request
            )
        finally:
            page.close()

在爬虫中使用:

def parse(self, response):
    # 渲染后的解析逻辑
    pass
def start_requests(self):
    for url in self.start_urls:
        yield scrapy.Request(
            url,
            callback=self.parse,
            meta={
                'render': True,           # 启用渲染
                'wait_for': '.product'    # 等待元素出现
            }
        )

安装Playwright:

pip install playwright
playwright install chromium

中间件执行顺序

中间件的优先级(数字越小越先执行):

下载中间件执行顺序:
0-400: Scrapy内置中间件
400-700: 自定义中间件
700+: Scrapy内置中间件
请求流程(从大到小):
引擎 → 中间件(400) → 中间件(410) → 下载器
响应流程(从小到大):
下载器 → 中间件(410) → 中间件(400) → 引擎

高级功能二:管道高级用法

上篇我们用了单个Pipeline做数据清洗和保存。实际项目中,你可能需要:

        • 同时保存到多种存储

        • 批量处理提高性能

        • 数据验证和补全

        • 实时去重和报警

实战1:多存储Pipeline

同时保存到JSON、CSV、MySQL三种存储。

# pipelines.py
import json
import csv
import pymysql
from scrapy.exceptions import DropItem
class MultipleStoragePipeline:
    """多存储管道"""
    def __init__(self):
        self.json_file = None
        self.csv_writer = None
        self.csv_file = None
        self.conn = None
        self.cursor = None
        self.items = []
    @classmethod
    def from_crawler(cls, crawler):
        pipeline = cls()
        crawler.signals.connect(pipeline.spider_opened, signals.spider_opened)
        crawler.signals.connect(pipeline.spider_closed, signals.spider_closed)
        return pipeline
    def spider_opened(self, spider):
        """爬虫启动时初始化连接"""
        # JSON
        self.json_file = open('books_data.json', 'w', encoding='utf-8')
        self.json_file.write('[')
        # CSV
        self.csv_file = open('books_data.csv', 'w', encoding='utf-8', newline='')
        self.csv_writer = csv.writer(self.csv_file)
        self.csv_writer.writerow(['title', 'price', 'rating', 'stock', 'url', 'crawl_time'])
        # MySQL
        self.conn = pymysql.connect(
            host='localhost',
            user='root',
            password='password',
            database='scrapy_db'
        )
        self.cursor = self.conn.cursor()
        # 创建表
        self.cursor.execute('''
            CREATE TABLE IF NOT EXISTS books (
                id INT AUTO_INCREMENT PRIMARY KEY,
                title VARCHAR(255),
                price DECIMAL(10,2),
                rating VARCHAR(50),
                stock VARCHAR(100),
                url TEXT,
                crawl_time DATETIME
            )
        ''')
    def process_item(self, item, spider):
        """处理每条数据"""
        # 数据验证
        if not item.get('title'):
            raise DropItem("Missing title")
        # 添加到列表
        self.items.append(item)
        # JSON: 逐条写入(逗号分隔)
        json_str = json.dumps(dict(item), ensure_ascii=False)
        self.json_file.write('\n  ' + json_str + ',')
        # CSV: 逐行写入
        self.csv_writer.writerow([
            item['title'],
            item['price'],
            item['rating'],
            item['stock'],
            item['url'],
            item['crawl_time']
        ])
        # MySQL: 插入数据
        self.cursor.execute('''
            INSERT INTO books (title, price, rating, stock, url, crawl_time)
            VALUES (%s, %s, %s, %s, %s, %s)
        ''', (
            item['title'],
            item['price'],
            item['rating'],
            item['stock'],
            item['url'],
            item['crawl_time']
        ))
        self.conn.commit()
        spider.logger.info(f"Saved: {item['title']}")
        return item
    def spider_closed(self, spider):
        """爬虫关闭时清理"""
        # JSON: 修正最后的逗号并闭合数组
        self.json_file.write('\n]')
        self.json_file.close()
        # CSV: 关闭文件
        self.csv_file.close()
        # MySQL: 关闭连接
        self.cursor.close()
        self.conn.close()
        spider.logger.info(f"Total saved: {len(self.items)} items")

配置:

ITEM_PIPELINES = {
    'book_scraper.pipelines.MultipleStoragePipeline': 300,
}

实战2:批量处理Pipeline

逐条保存到数据库效率太低,改为批量插入。

# pipelines.py
import pymysql
from scrapy.exceptions import DropItem
class BatchMySQLPipeline:
    """批量插入MySQL管道"""
    def __init__(self, batch_size=100):
        self.batch_size = batch_size
        self.conn = None
        self.cursor = None
        self.batch = []
    @classmethod
    def from_crawler(cls, crawler):
        return cls(
            batch_size=crawler.settings.getint('BATCH_SIZE', 100)
        )
    def open_spider(self, spider):
        """爬虫启动时初始化连接"""
        self.conn = pymysql.connect(
            host='localhost',
            user='root',
            password='password',
            database='scrapy_db'
        )
        self.cursor = self.conn.cursor()
    def process_item(self, item, spider):
        """暂存数据,达到批量大小后插入"""
        if not item.get('title'):
            raise DropItem("Missing title")
        self.batch.append(item)
        # 达到批量大小时插入
        if len(self.batch) >= self.batch_size:
            self._insert_batch()
        return item
    def close_spider(self, spider):
        """爬虫关闭时插入剩余数据"""
        if self.batch:
            self._insert_batch()
        self.cursor.close()
        self.conn.close()
    def _insert_batch(self):
        """批量插入数据"""
        if not self.batch:
            return
        sql = '''
            INSERT INTO books (title, price, rating, stock, url, crawl_time)
            VALUES (%s, %s, %s, %s, %s, %s)
        '''
        values = [
            (item['title'], item['price'], item['rating'],
             item['stock'], item['url'], item['crawl_time'])
            for item in self.batch
        ]
        self.cursor.executemany(sql, values)
        self.conn.commit()
        self.batch = []  # 清空批次

配置:

ITEM_PIPELINES = {
    'book_scraper.pipelines.BatchMySQLPipeline': 300,
}
BATCH_SIZE = 100  # 每批100条

实战3:数据验证Pipeline

确保数据完整性和质量。

# pipelines.py
from scrapy.exceptions import DropItem
class DataValidationPipeline:
    """数据验证管道"""
    def process_item(self, item, spider):
        """验证数据质量"""
        errors = []
        # 必填字段检查
        required_fields = ['title', 'price', 'rating']
        for field in required_fields:
            if not item.get(field):
                errors.append(f"Missing field: {field}")
        # 价格范围检查
        try:
            price = float(item.get('price', 0))
            if price < 0 or price > 10000:
                errors.append(f"Invalid price: {price}")
        except ValueError:
            errors.append(f"Invalid price format: {item.get('price')}")
        # URL格式检查
        url = item.get('url', '')
        if not url.startswith('http'):
            errors.append(f"Invalid URL: {url}")
        # 有错误则丢弃数据
        if errors:
            spider.logger.warning(f"Validation failed for {item.get('title', 'unknown')}: {errors}")
            raise DropItem(f"Validation failed: {errors}")
        return item

Pipeline执行顺序

Pipeline按优先级(数字越小越先执行):

ITEM_PIPELINES = {
    'book_scraper.pipelines.DataValidationPipeline': 100,  # 先验证
    'book_scraper.pipelines.BookPipeline': 200,            # 再清洗
    'book_scraper.pipelines.BatchMySQLPipeline': 300,      # 最后存储
}

高级功能三:分布式爬虫(Scrapy-Redis)

单机爬虫受限于硬件性能,当数据量达到百万级时,需要分布式架构。Scrapy-Redis是基于Redis的分布式爬虫框架,可以轻松扩展到多台机器。

Scrapy-Redis的工作原理

传统Scrapy架构:
调度器 → 爬虫1, 爬虫2, 爬虫3 (每个爬虫独立调度)
Scrapy-Redis架构:
Redis调度器 → 爬虫1, 爬虫2, 爬虫3 (共享调度器)

核心改进:

        • 共享调度器:所有爬虫实例共享一个Redis请求队列

        • 共享去重:使用Redis集合去重,避免重复请求

        • 增量爬取:支持断点续传,爬取过程中可随时停止恢复

安装Scrapy-Redis

pip install scrapy-redis
# 启动Redis
redis-server

创建分布式爬虫

修改原有的爬虫,继承RedisSpider:

# spiders/books_distributed.py
import scrapy
from scrapy_redis.spiders import RedisSpider
from book_scraper.items import BookItem
class BooksDistributedSpider(RedisSpider):
    """分布式商品爬虫"""
    name = 'books_distributed'
    redis_key = 'books:start_urls'  # Redis中存储起始URL的key
    # 原有的allowed_domains在分布式爬虫中不需要指定
    # allowed_domains = ['books.toscrape.com']
    def parse(self, response):
        """解析列表页"""
        books = response.css('article.product_pod')
        for book in books:
            item = BookItem()
            item['title'] = book.css('h3 a::attr(title)').get()
            item['price'] = book.css('p.price_color::text').get().replace('£', '')
            item['rating'] = book.css('p.star-rating::attr(class)').get().split()[-1]
            item['stock'] = book.css('p.instock availability::text').re_first(r'\w+.*').strip()
            item['url'] = response.urljoin(book.css('h3 a::attr(href)').get())
            import datetime
            item['crawl_time'] = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
            yield item
        # 处理下一页
        next_page = response.css('li.next a::attr(href)').get()
        if next_page:
            yield scrapy.Request(response.urljoin(next_page), callback=self.parse)

修改settings.py

# 启用Scrapy-Redis调度器
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
# 启用Redis去重
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
# 不清除Redis队列,支持断点续传
SCHEDULER_PERSIST = True
# Redis配置
REDIS_HOST = 'localhost'
REDIS_PORT = 6379
# 其他配置保持不变
ITEM_PIPELINES = {
    'book_scraper.pipelines.BookPipeline': 300,
}

启动分布式爬虫

方式1:手动推送起始URL

# 在Redis中添加起始URL
redis-cli
> LPUSH books:start_urls "http://books.toscrape.com/"
# 启动爬虫(可以启动多个实例)
scrapy crawl books_distributed

方式2:使用RedisSpider的自动模式

爬虫启动后会自动从books:start_urls列表中获取起始URL。

启动多个实例

在多台机器上同时运行:
# 机器1
scrapy crawl books_distributed
# 机器2
scrapy crawl books_distributed
# 机器3
scrapy crawl books_distributed

所有实例会共享Redis中的请求队列,协同工作。

监控Redis状态

# 查看请求队列长度
redis-cli
> LLEN books:requests
# 查看已处理的URL数量
> SCARD books:dupefilter
# 查看当前正在爬取的URL
> LINDEX books:requests 0

Scrapy-Redis vs 传统Scrapy

特性

传统Scrapy

Scrapy-Redis

调度器

内存队列

Redis队列

去重

内存集合

Redis集合

扩展性

单机

多机分布式

断点续传

不支持

支持

性能

受限于单机

可横向扩展

部署复杂度

简单

需要Redis

性能优化

当爬取规模达到十万、百万级时,性能优化变得至关重要。我总结了几个最有效的优化策略。

1. 并发调优

Scrapy默认的并发配置比较保守,需要根据实际情况调整。
# settings.py
# 全局并发数(所有域名)
CONCURRENT_REQUESTS = 32  # 默认16,可提高到64-128
# 单域名并发数
CONCURRENT_REQUESTS_PER_DOMAIN = 16  # 默认8,可提高到32-64
# 下载延迟(秒)
DOWNLOAD_DELAY = 0.5  # 根据目标网站承受能力调整,建议0.1-1.0
# 随机化延迟(避免规律性)
RANDOMIZE_DOWNLOAD_DELAY = True
# 超时设置
DOWNLOAD_TIMEOUT = 15  # 默认180秒,太快可能导致超时

并发数设置原则:

        • CPU密集型爬虫:并发数 = CPU核心数 × 2

        • IO密集型爬虫:并发数可以更高,64-128都没问题

        • 具体数值需要通过压力测试确定

2. 内存优化

Scrapy默认所有请求都在内存中,爬取大量URL时容易内存溢出。

# settings.py
# 限制内存中请求的数量
CONCURRENT_REQUESTS = 32
# 禁用调试模式(减少内存占用)
DEBUG = False
# 使用更高效的数据结构(需要自定义中间件)

3. 请求队列优化

Scrapy-Redis的请求队列可以优化:

# 使用更高效的队列算法
SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.PriorityQueue'  # 优先级队列
# SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.FifoQueue'   # FIFO队列
# SCHEDULER_QUEUE_CLASS = 'scrapy_redis.queue.LifoQueue'   # LIFO队列

4. 数据库优化

批量插入比逐条插入快10-20倍:

# pipelines.py
# 详见上面的BatchMySQLPipeline示例

5. 禁用不必要的功能

# settings.py
# 禁用Telnet控制台
TELNETCONSOLE_ENABLED = False
# 禁用统计信息收集(不推荐,但可以提升性能)
STATS_CLASS = None
# 减少日志输出
LOG_LEVEL = 'INFO'

6. 使用连接池

# settings.py
# 下载器连接池大小
REACTOR_THREADPOOL_MAXSIZE = 20
# 数据库连接池
DBPOOL_SIZE = 10

反爬策略应对

反爬是爬虫开发的永恒话题。这里总结几个实用的反爬应对方案。

1. User-Agent轮换

我们已经实现了RandomUserAgentMiddleware,这里补充一个更高级的版本:

# middlewares.py
import random
from fake_useragent import UserAgent
class AdvancedUserAgentMiddleware:
    """高级UA中间件(使用fake_useragent库)"""
    def __init__(self):
        self.ua = UserAgent()
    def process_request(self, request, spider):
        request.headers['User-Agent'] = self.ua.random
        spider.logger.debug(f"User-Agent: {request.headers['User-Agent']}")

安装fake-useragent:

pip install fake-useragent

2. IP代理池

我们已经实现了基础版本,实际项目中建议使用专业代理服务(如快代理、芝麻代理等):

# middlewares.py
import requests
class DynamicProxyMiddleware:
    """动态代理池中间件(从API获取)"""
    def __init__(self, proxy_api):
        self.proxy_api = proxy_api
        self.proxies = []
        self._fetch_proxies()
    def _fetch_proxies(self):
        """从API获取代理列表"""
        try:
            response = requests.get(self.proxy_api)
            self.proxies = response.json().get('proxies', [])
        except Exception as e:
            self.proxies = []
    def process_request(self, request, spider):
        """使用代理"""
        if self.proxies:
            proxy = random.choice(self.proxies)
            request.meta['proxy'] = proxy['url']
            spider.logger.debug(f"Using proxy: {proxy['url']}")

3. 请求频率控制

通过时间窗口限制请求频率:

# middlewares.py
import time
from collections import defaultdict
class RateLimitMiddleware:
    """频率限制中间件"""
    def __init__(self, max_requests=10, window=60):
        self.max_requests = max_requests  # 时间窗口内最大请求数
        self.window = window              # 时间窗口(秒)
        self.request_times = defaultdict(list)
    def process_request(self, request, spider):
        """检查请求频率"""
        domain = request.url.split('/')[2]
        now = time.time()
        # 移除过期的请求记录
        self.request_times[domain] = [
            t for t in self.request_times[domain] if now - t < self.window
        ]
        # 检查是否超限
        if len(self.request_times[domain]) >= self.max_requests:
            # 等待直到有名额
            wait_time = self.window - (now - self.request_times[domain][0])
            if wait_time > 0:
                time.sleep(wait_time)
        # 记录本次请求
        self.request_times[domain].append(now)

4. Cookie池

有些网站需要登录才能访问数据,需要维护Cookie池。

# middlewares.py
import random
import json
class CookiePoolMiddleware:
    """Cookie池中间件"""
    def __init__(self, cookie_file):
        with open(cookie_file, 'r', encoding='utf-8') as f:
            self.cookies = json.load(f)
    def process_request(self, request, spider):
        """随机分配Cookie"""
        if self.cookies:
            cookie = random.choice(self.cookies)
            request.cookies = cookie
            spider.logger.debug(f"Using cookie")

Cookie文件格式:

[
  {
    "sessionid": "abc123",
    "csrftoken": "xyz456"
  },
  {
    "sessionid": "def789",
    "csrftoken": "uvw012"
  }
]

5. 行为模拟

模拟真实用户的浏览行为:

# middlewares.py
import random
class HumanBehaviorMiddleware:
    """模拟人类行为中间件"""
    def process_request(self, request, spider):
        """随机延迟"""
        if request.meta.get('human_behavior'):
            delay = random.uniform(1, 3)
            time.sleep(delay)

6. 验证码处理

遇到验证码时,可以使用第三方服务(如超级鹰、打码兔):

# middlewares.py
import requests
class CaptchaMiddleware:
    """验证码处理中间件"""
    def __init__(self, captcha_api):
        self.captcha_api = captcha_api
    def process_request(self, request, spider):
        """遇到验证码时自动识别"""
        if request.meta.get('captcha'):
            captcha_image = self._get_captcha_image(request)
            captcha_code = self._recognize_captcha(captcha_image)
            request.meta['captcha_code'] = captcha_code

监控与告警

长期运行的爬虫需要完善的监控和告警系统。

1. 内置统计信息

Scrapy自带统计功能:

# 查看统计信息
scrapy crawl books --stats
# 保存统计信息到文件
scrapy crawl books --statsfile=stats.json

2. 自定义监控中间件

# middlewares.py
import time
from scrapy import signals
class MonitorMiddleware:
    """监控中间件"""
    def __init__(self):
        self.start_time = None
        self.request_count = 0
        self.error_count = 0
        self.item_count = 0
    @classmethod
    def from_crawler(cls, crawler):
        middleware = cls()
        crawler.signals.connect(middleware.spider_opened, signals.spider_opened)
        crawler.signals.connect(middleware.spider_closed, signals.spider_closed)
        crawler.signals.connect(middleware.item_scraped, signals.item_scraped)
        crawler.signals.connect(middleware.spider_error, signals.spider_error)
        return middleware
    def spider_opened(self, spider):
        """爬虫启动时记录时间"""
        self.start_time = time.time()
        spider.logger.info(f"Spider started: {spider.name}")
    def spider_closed(self, spider):
        """爬虫关闭时输出统计"""
        duration = time.time() - self.start_time
        spider.logger.info(f"""
            Spider finished: {spider.name}
            Duration: {duration:.2f}s
            Requests: {self.request_count}
            Items: {self.item_count}
            Errors: {self.error_count}
            Speed: {self.item_count / duration:.2f} items/s
        """)
    def item_scraped(self, item, response, spider):
        """每爬取一个item计数"""
        self.item_count += 1
    def spider_error(self, failure, response, spider):
        """每次错误计数"""
        self.error_count += 1

3. 告警通知

# pipelines.py
import requests
class AlertPipeline:
    """告警管道"""
    def __init__(self, webhook_url):
        self.webhook_url = webhook_url
    def process_item(self, item, spider):
        """每爬取1000条发送通知"""
        if spider.crawler.stats.get_value('item_scraped_count', 0) % 1000 == 0:
            self._send_alert(f"已爬取 {spider.crawler.stats.get_value('item_scraped_count')} 条数据")
        return item
    def _send_alert(self, message):
        """发送告警(钉钉/企业微信)"""
        data = {
            "msgtype": "text",
            "text": {
                "content": message
            }
        }
        requests.post(self.webhook_url, json=data)

4. 日志集中管理

使用ELK(Elasticsearch + Logstash + Kibana)或Sentry进行日志管理。

常见问题解决

问题1:内存溢出

现象:爬取一段时间后程序崩溃,提示"MemoryError"

原因:

        • 请求队列太大

        • Pipeline缓存了太多数据

        • 未正确清理资源

解决方案:

# 1. 限制并发数
CONCURRENT_REQUESTS = 16
# 2. 分批处理数据
# 详见BatchMySQLPipeline
# 3. 及时清理资源
def close_spider(self, spider):
    self.items.clear()
    self.cursor.close()

问题2:反爬严重

现象:频繁403/404,IP被封

解决方案:

        • 使用代理IP池

        • 降低请求频率

        • 随机化UA和Cookie

        • 模拟人类行为

问题3:数据重复

现象:数据库中有大量重复数据

解决方案:

        • 使用Pipeline去重

        • 在数据库表加唯一索引

        • 启用Scrapy的URL去重(Redis)

问题4:速度太慢

现象:爬取速度远低于预期

解决方案:

        • 提高并发数

        • 降低下载延迟

        • 使用批量插入

        • 检查网络连接

完整实战案例:分布式电商爬虫

现在我们把前面的知识点整合起来,做一个完整的分布式电商爬虫。

项目结构

ecommerce_crawler/
├── scrapy.cfg
├── requirements.txt
└── ecommerce_crawler/
    ├── __init__.py
    ├── items.py
    ├── middlewares.py
    ├── pipelines.py
    ├── settings.py
    └── spiders/
        ├── jd_spider.py
        ├── taobao_spider.py
        └── pinduoduo_spider.py

requirements.txt

scrapy==2.11.0
scrapy-redis==0.7.3
pymysql==1.1.0
fake-useragent==1.4.0
playwright==1.40.0
redis==5.0.1

items.py

import scrapy
class ProductItem(scrapy.Item):
    """商品数据模型"""
    title = scrapy.Field()
    price = scrapy.Field()
    sales = scrapy.Field()
    rating = scrapy.Field()
    shop = scrapy.Field()
    url = scrapy.Field()
    image_url = scrapy.Field()
    platform = scrapy.Field()  # JD/Taobao/Pinduoduo
    crawl_time = scrapy.Field()

middlewares.py

整合所有中间件:

from fake_useragent import UserAgent
import random
class ComprehensiveMiddleware:
    """综合中间件:UA+代理+限速+动态渲染"""
    def __init__(self, proxy_list):
        self.ua = UserAgent()
        self.proxy_list = proxy_list
    @classmethod
    def from_crawler(cls, crawler):
        return cls(
            proxy_list=crawler.settings.get('PROXY_LIST', [])
        )
    def process_request(self, request, spider):
        # 随机UA
        request.headers['User-Agent'] = self.ua.random
        # 随机代理
        if self.proxy_list:
            request.meta['proxy'] = random.choice(self.proxy_list)
        # 动态渲染(京东/淘宝需要)
        if request.meta.get('render'):
            request.meta['download_timeout'] = 30

pipelines.py

多存储+批量插入+验证+告警:

import pymysql
from scrapy.exceptions import DropItem
class ProductPipeline:
    """综合管道:验证+批量插入+告警"""
    def __init__(self):
        self.batch = []
        self.batch_size = 100
        self.conn = None
        self.cursor = None
    def open_spider(self, spider):
        self.conn = pymysql.connect(host='localhost', user='root', password='password', database='ecommerce')
        self.cursor = self.conn.cursor()
    def process_item(self, item, spider):
        # 验证
        if not item.get('title') or not item.get('price'):
            raise DropItem("Invalid item")
        # 批量处理
        self.batch.append(item)
        if len(self.batch) >= self.batch_size:
            self._insert_batch()
        # 告警(每1000条)
        if spider.crawler.stats.get_value('item_scraped_count', 0) % 1000 == 0:
            spider.logger.info(f"已爬取 {spider.crawler.stats.get_value('item_scraped_count')} 条数据")
        return item
    def close_spider(self, spider):
        if self.batch:
            self._insert_batch()
        self.cursor.close()
        self.conn.close()
    def _insert_batch(self):
        sql = 'INSERT INTO products (title, price, sales, rating, shop, url, platform) VALUES (%s, %s, %s, %s, %s, %s, %s)'
        values = [(item['title'], item['price'], item.get('sales'), item.get('rating'), item.get('shop'), item['url'], item['platform']) for item in self.batch]
        self.cursor.executemany(sql, values)
        self.conn.commit()
        self.batch.clear()

settings.py

# 基础配置
BOT_NAME = 'ecommerce_crawler'
SPIDER_MODULES = ['ecommerce_crawler.spiders']
NEWSPIDER_MODULE = 'ecommerce_crawler.spiders'
# Scrapy-Redis配置
SCHEDULER = "scrapy_redis.scheduler.Scheduler"
DUPEFILTER_CLASS = "scrapy_redis.dupefilter.RFPDupeFilter"
SCHEDULER_PERSIST = True
# 并发配置
CONCURRENT_REQUESTS = 64
CONCURRENT_REQUESTS_PER_DOMAIN = 32
DOWNLOAD_DELAY = 0.5
# 中间件
DOWNLOADER_MIDDLEWARES = {
    'ecommerce_crawler.middlewares.ComprehensiveMiddleware': 400,
}
# 管道
ITEM_PIPELINES = {
    'ecommerce_crawler.pipelines.ProductPipeline': 300,
}
# Redis
REDIS_HOST = 'localhost'
REDIS_PORT = 6379
# 代理列表(示例)
PROXY_LIST = [
    'http://proxy1.example.com:8080',
    'http://proxy2.example.com:8080',
]

JD爬虫示例

import scrapy
from scrapy_redis.spiders import RedisSpider
from ecommerce_crawler.items import ProductItem
class JDSpider(RedisSpider):
    name = 'jd'
    redis_key = 'jd:start_urls'
    def parse(self, response):
        """解析京东商品列表"""
        products = response.css('.gl-item')
        for product in products:
            item = ProductItem()
            item['title'] = product.css('.p-name em::text').get()
            item['price'] = product.css('.p-price i::text').get()
            item['sales'] = product.css('.p-commit::text').get()
            item['url'] = response.urljoin(product.css('.p-name a::attr(href)').get())
            item['platform'] = 'JD'
            import datetime
            item['crawl_time'] = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
            yield item
        # 下一页
        next_page = response.css('.pn-next::attr(href)').get()
        if next_page:
            yield scrapy.Request(response.urljoin(next_page), callback=self.parse)

启动分布式爬虫

# 启动Redis
redis-server
# 推送起始URL
redis-cli LPUSH jd:start_urls "https://search.jd.com/Search?keyword=手机"
# 启动爬虫(多台机器同时运行)
scrapy crawl jd

学习资源与扩展

官方文档

        • Scrapy官方文档: https://docs.scrapy.org/

        • Scrapy-Redis文档: https://github.com/rmax/scrapy-redis

推荐书籍

        • 《Python网络数据采集》Ryan Mitchell

        • 《精通Scrapy网络爬虫》刘伟

实战项目

        • GitHub上的优秀爬虫项目: https://github.com/topics/scrapy

        • 各种电商爬虫开源项目

进阶方向

        1. 机器学习:用ML识别验证码、预测反爬策略

        2. 大数据:结合Spark处理海量数据

        3. 流处理:使用Kafka进行实时数据处理

        4. 可视化:构建爬虫监控Dashboard

小结

下篇到这里就结束了。我们深入讲解了:

        • 中间件的高级用法(UA、代理、动态渲染)

        • 管道的高级用法(多存储、批量处理、验证、告警)

        • Scrapy-Redis分布式架构

        • 性能优化策略

        • 反爬应对方案

        • 监控与告警系统

        • 完整的分布式电商爬虫实战

总结

上下两篇完整覆盖了Scrapy从入门到精通的全流程:

上篇(入门):

        • Scrapy的核心概念和优势

        • 环境搭建和项目创建

        • 基础爬虫开发流程

        • 实战案例:电商商品爬取

        • 调试技巧

下篇(进阶):

        • 中间件开发(UA、代理、动态渲染)

        • 管道高级用法(多存储、批量处理)

        • Scrapy-Redis分布式爬虫

        • 性能优化策略

        • 反爬应对方案

        • 监控与告警

        • 完整的分布式爬虫实战

掌握了这些内容,你就可以独立开发百万级的爬虫系统了。记住,爬虫开发的核心不是技术,而是理解和解决问题的能力。

祝你在爬虫的道路上越走越远!

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐