项目管道

当一个爬取项被爬虫抓取后,它会被发送到项目管道,项目管道会通过顺序执行的多个组件来处理它。

每个项目管道组件(有时也简称为“项目管道”)是一个实现了简单方法的 Python 类。它们接收一个爬取项并对其执行操作,同时决定该爬取项是应继续通过管道,还是应被丢弃且不再处理。

项目管道的典型用途包括:

  • 清洗 HTML 数据

  • 验证抓取的数据(检查爬取项是否包含特定字段)

  • 检查重复项(并丢弃它们)

  • 将抓取的爬取项存储到数据库

编写自己的项目管道

每个项目管道都是一个组件,必须实现以下方法:

process_item(self, item)

此方法为每个项目管道组件调用。

item 是一个爬取项对象,请参阅支持所有爬取项类型

process_item() 必须返回一个爬取项对象,或者抛出DropItem异常。

被丢弃的爬取项将不再由后续管道组件处理。

参数:

item (爬取项对象) – 被抓取的爬取项

此外,它们还可以实现以下方法:

open_spider(self)

当爬虫开启时调用此方法。

close_spider(self)

当爬虫关闭时调用此方法。

这些方法中的任何一个都可以定义为协程函数(async def)。

项目管道示例

价格验证及丢弃无价格的爬取项

让我们看看下面这个假设的管道,它会调整不含增值税(price_excludes_vat属性)的爬取项的price属性,并丢弃那些不包含价格的爬取项:

from itemadapter import ItemAdapter
from scrapy.exceptions import DropItem


class PricePipeline:
    vat_factor = 1.15

    def process_item(self, item):
        adapter = ItemAdapter(item)
        if adapter.get("price"):
            if adapter.get("price_excludes_vat"):
                adapter["price"] = adapter["price"] * self.vat_factor
            return item
        else:
            raise DropItem("Missing price")

将爬取项写入 JSON Lines 文件

以下管道将所有抓取的爬取项(来自所有爬虫)存储到一个名为items.jsonl的文件中,该文件每行包含一个以 JSON 格式序列化的爬取项:

import json

from itemadapter import ItemAdapter


class JsonWriterPipeline:
    def open_spider(self):
        self.file = open("items.jsonl", "w")

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

    def process_item(self, item):
        line = json.dumps(ItemAdapter(item).asdict()) + "\n"
        self.file.write(line)
        return item

注意

JsonWriterPipeline 的目的是介绍如何编写爬取项管道。如果确实想将所有抓取的爬取项存储到 JSON 文件中,应使用数据导出(Feed exports)

将爬取项写入 MongoDB

在此示例中,我们将使用pymongo将爬取项写入MongoDB。MongoDB 地址和数据库名称在 Scrapy 设置中指定;MongoDB 集合以爬取项类命名。

此示例的重点是展示如何获取爬虫以及如何正确清理资源。

import pymongo
from itemadapter import ItemAdapter


class MongoPipeline:
    collection_name = "scrapy_items"

    def __init__(self, mongo_uri, mongo_db):
        self.mongo_uri = mongo_uri
        self.mongo_db = mongo_db

    @classmethod
    def from_crawler(cls, crawler):
        return cls(
            mongo_uri=crawler.settings.get("MONGO_URI"),
            mongo_db=crawler.settings.get("MONGO_DATABASE", "items"),
        )

    def open_spider(self):
        self.client = pymongo.MongoClient(self.mongo_uri)
        self.db = self.client[self.mongo_db]

    def close_spider(self):
        self.client.close()

    def process_item(self, item):
        self.db[self.collection_name].insert_one(ItemAdapter(item).asdict())
        return item

为爬取项截图

此示例演示了如何在process_item()方法中使用协程语法

此项目管道向本地运行的Splash实例发出请求,以渲染爬取项 URL 的屏幕截图。请求响应下载完成后,项目管道会将屏幕截图保存到文件,并将文件名添加到爬取项中。

import hashlib
from pathlib import Path
from urllib.parse import quote

import scrapy
from itemadapter import ItemAdapter
from scrapy.http.request import NO_CALLBACK


class ScreenshotPipeline:
    """Pipeline that uses Splash to render screenshot of
    every Scrapy item."""

    SPLASH_URL = "https://:8050/render.png?url={}"

    def __init__(self, crawler):
        self.crawler = crawler

    @classmethod
    def from_crawler(cls, crawler):
        return cls(crawler)

    async def process_item(self, item):
        adapter = ItemAdapter(item)
        encoded_item_url = quote(adapter["url"])
        screenshot_url = self.SPLASH_URL.format(encoded_item_url)
        request = scrapy.Request(screenshot_url, callback=NO_CALLBACK)
        response = await self.crawler.engine.download_async(request)

        if response.status != 200:
            # Error happened, return item.
            return item

        # Save screenshot to file, filename will be hash of url.
        url = adapter["url"]
        url_hash = hashlib.md5(url.encode("utf8")).hexdigest()
        filename = f"{url_hash}.png"
        Path(filename).write_bytes(response.body)

        # Store filename in item.
        adapter["screenshot_filename"] = filename
        return item

重复项过滤器

一个用于查找重复爬取项并丢弃已处理爬取项的过滤器。假设我们的爬取项具有唯一的 ID,但我们的爬虫返回了多个具有相同 ID 的爬取项:

from itemadapter import ItemAdapter
from scrapy.exceptions import DropItem


class DuplicatesPipeline:
    def __init__(self):
        self.ids_seen = set()

    def process_item(self, item):
        adapter = ItemAdapter(item)
        if adapter["id"] in self.ids_seen:
            raise DropItem(f"Item ID already seen: {adapter['id']}")
        else:
            self.ids_seen.add(adapter["id"])
            return item

激活项目管道组件

要激活项目管道组件,必须将其类添加到ITEM_PIPELINES设置中,如下例所示:

ITEM_PIPELINES = {
    "myproject.pipelines.PricePipeline": 300,
    "myproject.pipelines.JsonWriterPipeline": 800,
}

在此设置中为类分配的整数值决定了它们的运行顺序:爬取项从较低值类到较高值类依次通过。通常将这些数字定义在 0-1000 的范围内。

一个完整的示例

以上示例展示了独立的爬取项管道组件。在一个项目中,管道是协同工作的四个部分之一:您的爬虫产生的爬取项、生成该爬取项的爬虫、处理该爬取项的管道,以及启用管道的ITEM_PIPELINES设置。

以下示例将这些部分连接起来,以验证从books.toscrape.com抓取的书籍价格,并复用上述价格验证及丢弃无价格的爬取项中的PricePipeline

myproject/items.py中定义爬取项

from dataclasses import dataclass


@dataclass
class BookItem:
    title: str
    price: float

从您的爬虫中生成该爬取项的实例,例如在myproject/spiders/books.py

import scrapy

from myproject.items import BookItem


class BooksSpider(scrapy.Spider):
    name = "books"
    start_urls = ["https://books.toscrape.com/"]

    def parse(self, response):
        for book in response.css("article.product_pod"):
            yield BookItem(
                title=book.css("h3 a::attr(title)").get(),
                price=float(book.css("p.price_color::text").re_first(r"[\d.]+")),
            )

将前面展示的PricePipeline放入myproject/pipelines.py中,并在myproject/settings.py中启用它

ITEM_PIPELINES = {
    "myproject.pipelines.PricePipeline": 300,
}

有了这些组件,BooksSpider生成的每个BookItem都会先通过PricePipeline,然后才到达数据导出或任何其他输出。

常见陷阱

管道未运行

管道组件只有在其类列在ITEM_PIPELINES设置中时才会运行,通常在您项目的settings.py文件中(参见激活项目管道组件)。将其添加到爬虫或其他地方无效。

要确认 Scrapy 已加载您的管道,请在爬取日志的开头查找类似以下内容的行:

[scrapy.middleware] INFO: Enabled item pipelines:
['myproject.pipelines.PricePipeline']

如果您的管道不在该列表中,请检查其导入路径是否与ITEM_PIPELINES条目匹配,并检查该设置是否被覆盖,例如通过custom_settings或在settings.py中重新定义ITEM_PIPELINES

爬取项未返回

process_item() 必须返回爬取项(或抛出DropItem)。一个常见的错误是修改了爬取项却忘记返回它:

def process_item(self, item):
    ItemAdapter(item)["price"] *= 1.15
    # Bug: returns None, so the next component gets None instead of the item.

返回爬取项,以便下一个组件和 Scrapy 的其余部分可以继续处理它。

def process_item(self, item):
    ItemAdapter(item)["price"] *= 1.15
    return item