Scrapy Pipeline与中间件深度解析:从入门到实战

前言

在Scrapy框架中,Pipeline(管道)和Middleware(中间件)是构建健壮爬虫的核心组件。Pipeline负责数据的后处理(清洗、去重、持久化),而中间件则控制请求/响应的生命周期(UA伪装、代理切换、异常处理)。本文将结合实战场景,详细剖析两者的用法、返回值机制及工作流程,并对比爬虫中间件与下载中间件的异同。


一、Pipeline管道:数据清洗与持久化的“最后一公里”

1.1 Pipeline的工作流程

Pipeline在Spider生成Item后被依次调用。每个Pipeline类实现特定方法,返回Item或DropItem来控制数据流向。默认按settings.pyITEM_PIPELINES字典的优先级顺序执行(数值越小越优先)。

# settings.py
ITEM_PIPELINES = {
'myproject.pipelines.DataCleanPipeline': 300,
'myproject.pipelines.DuplicatesPipeline': 400,
'myproject.pipelines.MysqlPipeline': 500,
'myproject.pipelines.RedisPipeline': 600,
}

1.2 核心方法与返回值

方法触发时机返回值作用
process_item(self, item, spider)每个Item经过Pipeline时必须返回Item或DropItem异常处理数据,可修改item
open_spider(self, spider)Spider开启时初始化资源(如数据库连接)
close_spider(self, spider)Spider关闭时释放资源(如关闭连接)
from_crawler(cls, crawler)类方法,用于获取配置返回Pipeline实例从settings读取参数

关键返回值说明

  • 返回Item:继续传递给下一个Pipeline。
  • 抛出DropItem异常:丢弃该Item,后续Pipeline不再执行。
  • 返回None:Scrapy会忽略该Item(不推荐,建议显式返回Item或DropItem)。

1.3 实战案例:数据清洗与去重

# pipelines.py
import re
from scrapy.exceptions import DropItem

class DataCleanPipeline:
    """清洗数据:去除HTML标签、统一日期格式"""
    def process_item(self, item, spider):
        # 清洗标题中的空白字符
        item['title'] = item['title'].strip()
        # 去除HTML标签
        item['content'] = re.sub(r'<[^>]+>', '', item['content'])
        # 日期格式统一为YYYY-MM-DD
        if item.get('date'):
            item['date'] = item['date'].replace('/', '-')
        return item

class DuplicatesPipeline:
    """基于Redis的去重(需安装redis)"""
    def __init__(self):
        import redis
        self.redis = redis.Redis(host='localhost', port=6379, db=0)
        self.dup_key = 'scrapy:duplicates'

    def process_item(self, item, spider):
        # 使用item的url作为唯一标识
        url = item.get('url')
        if url and self.redis.sismember(self.dup_key, url):
            raise DropItem(f"Duplicate item found: {url}")
        else:
            self.redis.sadd(self.dup_key, url)
            return item

1.4 自动写入MySQL/Redis

MySQL写入示例(使用pymysql):

import pymysql

class MysqlPipeline:
    def open_spider(self, spider):
        self.conn = pymysql.connect(
            host='localhost', user='root', password='123456',
            database='scrapy_db', charset='utf8mb4'
        )
        self.cursor = self.conn.cursor()

    def process_item(self, item, spider):
        sql = "INSERT INTO articles(title, content, date) VALUES(%s, %s, %s)"
        self.cursor.execute(sql, (item['title'], item['content'], item['date']))
        self.conn.commit()
        return item

    def close_spider(self, spider):
        self.cursor.close()
        self.conn.close()

Redis写入示例(存储为Hash):

import redis

class RedisPipeline:
    def open_spider(self, spider):
        self.r = redis.Redis(host='localhost', port=6379, db=1)

    def process_item(self, item, spider):
        # 以url为key,存储整个item为hash
        self.r.hset(f"article:{item['url']}", mapping=item)
        return item

1.5 数据持久化落地(文件存储)

import json

class JsonPipeline:
    def open_spider(self, spider):
        self.file = open('items.json', 'a', encoding='utf-8')

    def process_item(self, item, spider):
        line = json.dumps(dict(item), ensure_ascii=False) + '\n'
        self.file.write(line)
        return item

    def close_spider(self, spider):
        self.file.close()

二、下载中间件:请求与响应的“魔法工厂”

2.1 下载中间件的工作流程

下载中间件位于Engine和Downloader之间,处理每个Request和Response。其核心方法有四个,按执行顺序如下:

  1. process_request(request, spider) → 返回Request/Response/None/异常
  2. process_response(request, response, spider) → 返回Response/Request/异常
  3. process_exception(request, exception, spider) → 返回Request/Response/None
  4. from_crawler(cls, crawler) → 获取配置

执行顺序
process_request(按优先级从小到大)→ 下载器 → process_response(按优先级从大到小)

2.2 自定义UA池随机切换

# middlewares.py
import random
from scrapy import signals
from scrapy.downloadermiddlewares.useragent import UserAgentMiddleware

class RandomUserAgentMiddleware(UserAgentMiddleware):
    def __init__(self, user_agent_list):
        self.user_agent_list = user_agent_list

    @classmethod
    def from_crawler(cls, crawler):
        # 从settings获取UA列表
        ua_list = crawler.settings.get('USER_AGENT_LIST', [])
        return cls(ua_list)

    def process_request(self, request, spider):
        ua = random.choice(self.user_agent_list)
        request.headers['User-Agent'] = ua
        # 注意:不要调用父类方法,否则可能覆盖

settings.py中配置:

USER_AGENT_LIST = [
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36...',
'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7)...',
# 更多UA...
]
DOWNLOADER_MIDDLEWARES = {
'myproject.middlewares.RandomUserAgentMiddleware': 400,
}

2.3 代理IP池配置

class ProxyMiddleware:
    def __init__(self, proxy_list):
        self.proxy_list = proxy_list

    @classmethod
    def from_crawler(cls, crawler):
        return cls(crawler.settings.get('PROXY_LIST', []))

    def process_request(self, request, spider):
        if self.proxy_list:
            proxy = random.choice(self.proxy_list)
            request.meta['proxy'] = proxy  # 格式如 'http://user:pass@ip:port'

注意:代理IP需要定期验证有效性,可结合process_exception处理代理失效的情况。

2.4 请求异常捕获与重试

from scrapy.downloadermiddlewares.retry import RetryMiddleware
from scrapy.utils.response import response_status_message

class CustomRetryMiddleware(RetryMiddleware):
    def process_exception(self, request, exception, spider):
        # 捕获超时、连接错误等异常
        if isinstance(exception, (TimeoutError, ConnectionError)):
            spider.logger.warning(f"Retry due to {exception}")
            # 更换代理后重试
            new_proxy = self.get_new_proxy()
            request.meta['proxy'] = new_proxy
            return self._retry(request, exception, spider)
        return None

    def get_new_proxy(self):
        # 从代理池获取新代理
        return random.choice(spider.settings.get('PROXY_LIST'))

2.5 响应篡改(如修改状态码、注入Cookie)

class ResponseModifyMiddleware:
    def process_response(self, request, response, spider):
        # 如果遇到403,尝试修改响应状态码欺骗Spider
        if response.status == 403:
            # 可以修改response对象,但注意不要破坏原始结构
            response.status = 200
            # 或者注入自定义header
            response.headers['X-Custom'] = 'modified'
        return response

注意:篡改响应需谨慎,可能导致数据不一致。更推荐在Spider中处理异常状态。


三、爬虫中间件 vs 下载中间件:异同对比

3.1 相同点

  • 都是通过settings.py中的字典配置优先级。
  • 都提供from_crawler类方法获取配置。
  • 都可以通过process_*方法拦截并修改数据流。
  • 都支持抛出IgnoreRequest异常来丢弃请求/响应。

3.2 不同点

维度爬虫中间件 (Spider Middleware)下载中间件 (Downloader Middleware)
作用位置Engine与Spider之间Engine与Downloader之间
处理对象Request(从Engine到Spider)和Response(从Spider到Engine)Request(从Engine到Downloader)和Response(从Downloader到Engine)
核心方法process_spider_input(response)process_spider_output(response, result)process_spider_exception(response, exception)process_start_requests(start_requests)process_request(request, spider)process_response(request, response, spider)process_exception(request, exception, spider)
典型用途修改Spider的输入输出(如过滤响应、注入额外请求、修改Item)修改请求/响应本身(如UA、代理、Cookie、重试)
返回值影响返回Request/Item/None等,影响Spider的后续处理返回Request/Response/None,影响下载器或Engine
优先级顺序process_spider_input按优先级从小到大;process_spider_output按优先级从大到小process_request按优先级从小到大;process_response按优先级从大到小

3.3 实战选择建议

  • 需要修改请求头、代理、Cookie → 下载中间件。
  • 需要修改Spider接收到的响应内容(如解析前预处理HTML)→ 爬虫中间件(process_spider_input)。
  • 需要动态生成新的请求(如翻页)→ 爬虫中间件(process_spider_output)。
  • 需要全局重试机制 → 下载中间件(process_exception)。

四、完整工作流程时序图(文字描述)

  1. Engine从Spider获取初始Request。
  2. Engine将Request发送给爬虫中间件process_start_requests可修改初始请求)。
  3. 爬虫中间件处理后,Request进入下载中间件process_request链。
  4. 下载中间件修改Request(如添加UA、代理),然后交给Downloader
  5. Downloader下载网页,返回Response。
  6. Response进入下载中间件process_response链(可修改响应、重试等)。
  7. 处理后的Response进入爬虫中间件process_spider_input
  8. Spider解析Response,生成Item或新的Request。
  9. Spider的输出(Item/Request)进入爬虫中间件process_spider_output
  10. 最终Item进入Pipeline链,Request重新进入步骤1。

五、总结

  • Pipeline是数据后处理的核心,通过process_item的返回值控制数据流向,适合清洗、去重、持久化。
  • 下载中间件是请求/响应的“过滤器”,通过process_requestprocess_response实现UA随机、代理切换、异常重试等。
  • 爬虫中间件更贴近Spider逻辑,适合修改Spider的输入输出,但日常开发中下载中间件使用频率更高。
  • 合理配置优先级(数值越小越优先)和返回值,是构建稳定爬虫的关键。
Logo

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

更多推荐