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 文档维护者。本文为中文翻译,代码及命令保留原文。











暂无评论内容