项目管道
当一个爬取项被爬虫抓取后,它会被发送到项目管道,项目管道会通过顺序执行的多个组件来处理它。
每个项目管道组件(有时也简称为“项目管道”)是一个实现了简单方法的 Python 类。它们接收一个爬取项并对其执行操作,同时决定该爬取项是应继续通过管道,还是应被丢弃且不再处理。
项目管道的典型用途包括:
清洗 HTML 数据
验证抓取的数据(检查爬取项是否包含特定字段)
检查重复项(并丢弃它们)
将抓取的爬取项存储到数据库
编写自己的项目管道
每个项目管道都是一个组件,必须实现以下方法:
- process_item(self, 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