Scrapy Item Pipeline 数据处理管道

Spider 抓取出 item 后,会将它交给 Item Pipeline,由多个组件依次处理。

每个组件都是实现简单方法的 Python 类。它接收 item,执行操作,并决定让它继续流向后续组件,还是丢弃、不再处理。

典型用途包括清理 HTML、验证数据是否包含指定字段、检测并丢弃重复项,以及写入数据库。

编写自己的管道

每个组件必须实现 process_item(self, item)。参数 item 是 item 对象,相关类型兼容说明见 Supporting All Item Types。

方法必须返回 item,或者抛出 DropItem。被丢弃的 item 不再进入后续组件。

还可以实现:

  • open_spider(self):Spider 打开时调用。Scrapy 2.18.0 起支持抛出 CloseSpider,在抓取前关闭 Spider,例如所需资源不可用时。
  • close_spider(self):Spider 关闭时、发送 spider_closed 信号之前调用。

这些方法都可以定义为 async def 协程。

管道示例

验证价格,丢弃无价格项目

下面对 price_excludes_vat 为真的项目调整价格,加入增值税,并丢弃没有价格的 item:

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

下面把全部 Spider 的 item 写入同一个 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

注意:这个例子只是演示如何编写管道。实际需要把所有 item 保存为 JSON 文件时,应使用 Feed exports。

写入 MongoDB

下面通过 pymongo 写入 MongoDB。服务器地址和数据库名来自 Scrapy 设置,集合名定义为类属性。

重点是演示如何取得 crawler,以及正确清理资源:

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

为 item 截图

这个例子展示 process_item() 的协程语法。管道请求本地 Splash,为 item URL 渲染截图。响应下载后,把截图保存为文件,并将文件名写入 item:

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 = "http://localhost: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

过滤重复项

假设 item 有唯一 id,但 Spider 可能多次返回同一 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 范围的数值。

完整示例

前面的例子是独立组件。实际项目由四部分协同工作:Spider 生成的 item、产出它的 Spider、处理它的管道,以及启用管道的 ITEM_PIPELINES 设置。

下面复用 PricePipeline,验证从 books.toscrape.com 抓取的图书价格。

在 myproject/items.py 定义 item:

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,在到达 Feed exports 或其他输出之前,都先经过 PricePipeline。

测试管道

让单个 URL 的 item 经过管道,可以使用带 --pipelines 的 parse:

scrapy parse --pipelines "https://books.toscrape.com/"

要测试特定数据,添加一个根据关键字参数构造 item 的回调:

class BooksSpider(scrapy.Spider):
    # ...

    def parse_item(self, response, **fields):
        yield BookItem(**fields)

然后通过命令行传入参数:

scrapy parse --pipelines -c parse_item --cbkwargs '{"title": "Test", "price": 10}' "https://books.toscrape.com/"

URL 可以是 Spider 能处理的任意地址。即使回调不使用响应,它仍会下载。

常见问题

管道未运行

只有类列入 ITEM_PIPELINES 才会运行,该设置通常位于项目 settings.py。把管道加入 Spider 或其他位置没有作用。

抓取日志开头应出现类似:

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

缺少时,核对导入路径与 ITEM_PIPELINES 条目是否一致,以及是否被 custom_settings 或 settings.py 中重复定义覆盖。

忘记返回 item

process_item() 必须返回 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.

应该返回 item,让下一个组件以及 Scrapy 后续流程继续处理:

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

原文:Item Pipeline。作者/维护方:Scrapy 文档维护者。本文为中文翻译,代码及命令保留原文。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容