Scrapy 中间件与 Pipeline

Scrapy 的骨架只负责「调度—下载—解析」,真正决定爬虫能否长期稳定运行的是中间件与 Pipeline:前者管请求怎么发出去、响应怎么收回来,后者管数据怎么清洗、去重、落地。

中间件的两层与执行顺序

类型配置项处理对象典型用途
下载器中间件DOWNLOADER_MIDDLEWARESRequest / Response随机 UA、代理注入、重试、异常兜底
Spider 中间件SPIDER_MIDDLEWARESItem / Request过滤站外链接、补字段、控制深度

process_request 按优先级从小到大执行,process_responseprocess_exception 从大到小:

Engine → SpiderMW(小→大) → DownloaderMW(小→大) → Downloader
Engine ← SpiderMW(大→小) ← DownloaderMW(大→小) ← Downloader
process_request   返回 None 放行;返回 Request/Response 提前短路
process_response  返回 Response 放行;抛 IgnoreRequest 则丢弃响应
process_exception 仅前序中间件抛异常时被调用,返回 None 交给下一个

内置中间件的默认数字是固定的:RobotsTxtMiddleware: 100DefaultHeadersMiddleware: 400UserAgentMiddleware: 500RetryMiddleware: 550RedirectMiddleware: 600HttpProxyMiddleware: 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
最后一批数据没入库只在条数达标时 flushclose_spider 里补一次 flush
日志刷屏但数据为 0DropItem 未打日志raise DropItem(...) 里带上关键字段

小结:中间件解决「请求与响应」的横切问题,Pipeline 解决「数据」的横切问题;记住 process_request 小数字先走、process_response 大数字先走,把清洗、去重、入库拆成独立 Pipeline,入库用「攒 N 条或 T 秒」的缓冲并在 close_spider 里兜底冲刷。

笔记加载中…