【爬虫入门第12讲】Scrapy 框架深度解析(三)Pipeline与中间件
·
Scrapy Pipeline与中间件深度解析:从入门到实战
前言
在Scrapy框架中,Pipeline(管道)和Middleware(中间件)是构建健壮爬虫的核心组件。Pipeline负责数据的后处理(清洗、去重、持久化),而中间件则控制请求/响应的生命周期(UA伪装、代理切换、异常处理)。本文将结合实战场景,详细剖析两者的用法、返回值机制及工作流程,并对比爬虫中间件与下载中间件的异同。
一、Pipeline管道:数据清洗与持久化的“最后一公里”
1.1 Pipeline的工作流程
Pipeline在Spider生成Item后被依次调用。每个Pipeline类实现特定方法,返回Item或DropItem来控制数据流向。默认按settings.py中ITEM_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。其核心方法有四个,按执行顺序如下:
process_request(request, spider)→ 返回Request/Response/None/异常process_response(request, response, spider)→ 返回Response/Request/异常process_exception(request, exception, spider)→ 返回Request/Response/Nonefrom_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)。
四、完整工作流程时序图(文字描述)
- Engine从Spider获取初始Request。
- Engine将Request发送给爬虫中间件(
process_start_requests可修改初始请求)。 - 爬虫中间件处理后,Request进入下载中间件的
process_request链。 - 下载中间件修改Request(如添加UA、代理),然后交给Downloader。
- Downloader下载网页,返回Response。
- Response进入下载中间件的
process_response链(可修改响应、重试等)。 - 处理后的Response进入爬虫中间件的
process_spider_input。 - Spider解析Response,生成Item或新的Request。
- Spider的输出(Item/Request)进入爬虫中间件的
process_spider_output。 - 最终Item进入Pipeline链,Request重新进入步骤1。
五、总结
- Pipeline是数据后处理的核心,通过
process_item的返回值控制数据流向,适合清洗、去重、持久化。 - 下载中间件是请求/响应的“过滤器”,通过
process_request和process_response实现UA随机、代理切换、异常重试等。 - 爬虫中间件更贴近Spider逻辑,适合修改Spider的输入输出,但日常开发中下载中间件使用频率更高。
- 合理配置优先级(数值越小越优先)和返回值,是构建稳定爬虫的关键。
更多推荐
所有评论(0)