Scrapy 中间件与 Pipeline
Scrapy 的骨架只负责「调度—下载—解析」,真正决定爬虫能否长期稳定运行的是中间件与 Pipeline:前者管请求怎么发出去、响应怎么收回来,后者管数据怎么清洗、去重、落地。
中间件的两层与执行顺序
| 类型 | 配置项 | 处理对象 | 典型用途 |
|---|---|---|---|
| 下载器中间件 | DOWNLOADER_MIDDLEWARES | Request / Response | 随机 UA、代理注入、重试、异常兜底 |
| Spider 中间件 | SPIDER_MIDDLEWARES | Item / Request | 过滤站外链接、补字段、控制深度 |
process_request 按优先级从小到大执行,process_response 与 process_exception 从大到小:
Engine → SpiderMW(小→大) → DownloaderMW(小→大) → Downloader
Engine ← SpiderMW(大→小) ← DownloaderMW(大→小) ← Downloader
process_request 返回 None 放行;返回 Request/Response 提前短路
process_response 返回 Response 放行;抛 IgnoreRequest 则丢弃响应
process_exception 仅前序中间件抛异常时被调用,返回 None 交给下一个
内置中间件的默认数字是固定的:RobotsTxtMiddleware: 100、DefaultHeadersMiddleware: 400、UserAgentMiddleware: 500、RetryMiddleware: 550、RedirectMiddleware: 600、HttpProxyMiddleware: 750,自己的中间件避开这些数字即可。
下载器中间件:换 UA、注代理、兜异常
# myproj/middlewares.py
import logging, random
logger = logging.getLogger(__name__)
UA_POOL = ["Mozilla/5.0 (Windows NT 10.0; Win64; x64) Chrome/122.0 Safari/537.36",
"Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) Version/17.0 Safari/605.1.15"]
class RandomUserAgentMiddleware:
"""每个请求随机挑一个 UA,避免整站同指纹高频访问"""
def __init__(self, ua_pool):
self.ua_pool = ua_pool
@classmethod
def from_crawler(cls, crawler):
return cls(crawler.settings.getlist("UA_POOL") or UA_POOL)
def process_request(self, request, spider):
request.headers["User-Agent"] = random.choice(self.ua_pool)
return None # None 表示继续走后续中间件
class ProxyMiddleware:
"""从代理池取一个代理;池子为空则放行,不阻断采集"""
def __init__(self, proxy_pool):
self.proxy_pool = proxy_pool
@classmethod
def from_crawler(cls, crawler):
return cls(crawler.settings.getlist("PROXY_POOL"))
def process_request(self, request, spider):
if self.proxy_pool and not request.meta.get("proxy"):
request.meta["proxy"] = random.choice(self.proxy_pool)
def process_exception(self, request, exception, spider):
logger.warning("请求异常 url=%s err=%s", request.url, exception)
return None # 交给 RetryMiddleware 决定是否重试
注意:process_response 如果返回 Request,Scrapy 会把它重新放回调度器,而不是继续往后的中间件。
Pipeline 三件套:清洗、去重、入库
process_item 返回 Item 交给下一个 Pipeline,抛 DropItem 表示丢弃;open_spider / close_spider 各执行一次,是建连接与冲刷缓冲的正确位置。
# myproj/pipelines.py
import time
from itemadapter import ItemAdapter
from scrapy.exceptions import DropItem
class CleanPipeline:
"""去空白、丢空值,保证后面拿到的字段干净"""
def process_item(self, item, spider):
adapter = ItemAdapter(item)
adapter["title"] = (adapter.get("title") or "").strip()
if not adapter["title"]:
raise DropItem(f"标题为空,丢弃 {adapter.get('url')}")
return item
class BufferMysqlPipeline:
"""攒够 200 条或超过 30 秒才写一次库"""
BATCH_SIZE, FLUSH_INTERVAL = 200, 30.0
def open_spider(self, spider):
import pymysql
self.conn = pymysql.connect(host="127.0.0.1", user="crawler", password="***",
database="spider", charset="utf8mb4", autocommit=False)
self.buffer, self.last_flush = [], time.time()
def process_item(self, item, spider):
self.buffer.append(ItemAdapter(item).asdict())
if len(self.buffer) >= self.BATCH_SIZE or time.time() - self.last_flush > self.FLUSH_INTERVAL:
self.flush()
return item
def flush(self):
if not self.buffer:
return
sql = ("INSERT INTO news(title, url, pub_time) VALUES(%s, %s, %s) "
"ON DUPLICATE KEY UPDATE title=VALUES(title)")
with self.conn.cursor() as cur:
cur.executemany(sql, [(d["title"], d["url"], d["pub_time"]) for d in self.buffer])
self.conn.commit()
self.buffer.clear()
self.last_flush = time.time()
def close_spider(self, spider):
self.flush() # 收尾必须冲一次,否则最后一批数据丢失
self.conn.close()
settings.py 启用与优先级含义
DOWNLOADER_MIDDLEWARES = {
"myproj.middlewares.RandomUserAgentMiddleware": 400,
"myproj.middlewares.ProxyMiddleware": 543,
"myproj.middlewares.RetryOrDropMiddleware": 560, # 整数区间 0~1000,越小越靠近引擎
}
ITEM_PIPELINES = {
"myproj.pipelines.CleanPipeline": 300,
"myproj.pipelines.DedupPipeline": 400,
"myproj.pipelines.BufferMysqlPipeline": 900,
}
CONCURRENT_REQUESTS = 8
CONCURRENT_REQUESTS_PER_DOMAIN = 2
DOWNLOAD_DELAY = 1.0 # 同域名两次请求的平均间隔
AUTOTHROTTLE_ENABLED = True # 被限流时自动放慢
常见坑
| 现象 | 原因 | 处理 |
|---|---|---|
| 中间件完全不生效 | 路径写错或数字被内置项覆盖 | 确认可 import,避开 100/400/500/550/600 |
| 重试请求被去重丢掉 | dont_filter 未置 True | 重新入队的请求设 dont_filter=True |
| 最后一批数据没入库 | 只在条数达标时 flush | close_spider 里补一次 flush |
| 日志刷屏但数据为 0 | DropItem 未打日志 | raise DropItem(...) 里带上关键字段 |
小结:中间件解决「请求与响应」的横切问题,Pipeline 解决「数据」的横切问题;记住 process_request 小数字先走、process_response 大数字先走,把清洗、去重、入库拆成独立 Pipeline,入库用「攒 N 条或 T 秒」的缓冲并在 close_spider 里兜底冲刷。