Some improvements to Item Pipeline (closes #195):

* Made Item Pipeline Manager a subclass of scrapy.middleware.MiddlewareManager
* Added open_spider/close_spider methods with support for returning deferreds from them
* Inverted the process_item() arguments to be more friendly with deferred
  callbacks (backwards compatibility kept through arguments introspection)
* Updated documentation with new methods and process_item() arguments change
This commit is contained in:
Pablo Hoffman 2010-08-12 10:48:37 -03:00
parent 2167bfb20d
commit 43d47e5d9b
14 changed files with 82 additions and 92 deletions

View File

@ -156,7 +156,7 @@ extracted item into a file using `pickle`_::
import pickle
class StoreItemPipeline(object):
def process_item(self, spider, item):
def process_item(self, item, spider):
torrent_id = item['url'].split('/')[-1]
f = open("torrent-%s.pickle" % torrent_id, "w")
pickle.dump(item, f)

View File

@ -442,7 +442,7 @@ creation step, it's in ``dmoz/pipelines.py`` and looks like this::
# Define your item pipelines here
class DmozPipeline(object):
def process_item(self, spider, item):
def process_item(self, item, spider):
return item
We have to override the ``process_item`` method in order to store our Items
@ -458,7 +458,7 @@ separated values) file using the standard library `csv module`_::
def __init__(self):
self.csvwriter = csv.writer(open('items.csv', 'wb'))
def process_item(self, spider, item):
def process_item(self, item, spider):
self.csvwriter.writerow([item['title'][0], item['link'][0], item['desc'][0]])
return item

View File

@ -62,7 +62,7 @@ Exporter to export scraped items to different files, one per spider::
file = self.files.pop(spider)
file.close()
def process_item(self, spider, item):
def process_item(self, item, spider):
self.exporter.export_item(item)
return item

View File

@ -4,9 +4,6 @@
Item Pipeline
=============
.. module:: scrapy.contrib.pipeline
:synopsis: Item Pipeline manager and built-in pipelines
After an item has been scraped by a spider it is sent to the Item Pipeline
which process it through several components that are executed sequentially.
@ -22,20 +19,36 @@ Writing your own item pipeline
==============================
Writing your own item pipeline is easy. Each item pipeline component is a
single Python class that must define the following method:
single Python class that must implement the following method:
.. method:: process_item(spider, item)
.. method:: process_item(item, spider)
:param spider: the spider which scraped the item
:type spider: :class:`~scrapy.spider.BaseSpider` object
This method is called for every item pipeline component and must either return
a :class:`~scrapy.item.Item` (or any descendant class) object or raise a
:exc:`~scrapy.exceptions.DropItem` exception. Dropped items are no longer
processed by further pipeline components.
:param item: the item scraped
:type item: :class:`~scrapy.item.Item` object
This method is called for every item pipeline component and must either return
a :class:`~scrapy.item.Item` (or any descendant class) object or raise a
:exc:`~scrapy.exceptions.DropItem` exception. Dropped items are no longer
processed by further pipeline components.
:param spider: the spider which scraped the item
:type spider: :class:`~scrapy.spider.BaseSpider` object
Additionally, they may also implement the following methods:
.. method:: open_spider(spider)
This method is called when the spider is opened.
:param spider: the spider which was opened
:type spider: :class:`~scrapy.spider.BaseSpider` object
.. method:: close_spider(spider)
This method is called when the spider is closed.
:param spider: the spider which was closed
:type spider: :class:`~scrapy.spider.BaseSpider` object
Item pipeline example
@ -51,7 +64,7 @@ attribute), and drops those items which don't contain a price::
vat_factor = 1.15
def process_item(self, spider, item):
def process_item(self, item, spider):
if item['price']:
if item['price_excludes_vat']:
item['price'] = item['price'] * self.vat_factor
@ -97,7 +110,7 @@ spider returns multiples items with the same id::
def spider_closed(self, spider):
del self.duplicates[spider]
def process_item(self, spider, item):
def process_item(self, item, spider):
if item['id'] in self.duplicates[spider]:
raise DropItem("Duplicate item found: %s" % item)
else:
@ -107,6 +120,9 @@ spider returns multiples items with the same id::
Built-in Item Pipelines reference
=================================
.. module:: scrapy.contrib.pipeline
:synopsis: Item Pipeline manager and built-in pipelines
Here is a list of item pipelines bundled with Scrapy.
.. _file-export-pipeline:

View File

@ -14,7 +14,7 @@ class FilterWordsPipeline(object):
# put all words in lowercase
words_to_filter = ['politics', 'religion']
def process_item(self, spider, item):
def process_item(self, item, spider):
for word in self.words_to_filter:
if word in unicode(item['description']).lower():
raise DropItem("Contains forbidden word: %s" % word)

View File

@ -4,5 +4,5 @@
# See: http://doc.scrapy.org/topics/item-pipeline.html
class ImdbPipeline(object):
def process_item(self, spider, item):
def process_item(self, item, spider):
return item

View File

@ -7,7 +7,7 @@ class FilterWordsPipeline(object):
# put all words in lowercase
words_to_filter = ['politics', 'religion']
def process_item(self, spider, item):
def process_item(self, item, spider):
for word in self.words_to_filter:
if word in unicode(item['description']).lower():
raise DropItem("Contains forbidden word: %s" % word)

View File

@ -5,61 +5,32 @@ See documentation in docs/item-pipeline.rst
"""
from scrapy import log
from scrapy.exceptions import NotConfigured
from scrapy.item import BaseItem
from scrapy.utils.misc import load_object
from scrapy.utils.defer import defer_succeed, mustbe_deferred
from scrapy.conf import settings
from scrapy.middleware import MiddlewareManager
class ItemPipelineManager(object):
class ItemPipelineManager(MiddlewareManager):
def __init__(self):
self.loaded = False
self.enabled = {}
self.disabled = {}
self.pipeline = []
self.load()
component_name = 'item pipeline'
def load(self):
"""
Load pipelines stages defined in settings module
"""
self.enabled.clear()
self.disabled.clear()
for pipepath in settings.getlist('ITEM_PIPELINES'):
cls = load_object(pipepath)
if cls:
try:
pipe = cls()
self.pipeline.append(pipe)
self.enabled[cls.__name__] = pipe
except NotConfigured, e:
self.disabled[cls.__name__] = pipepath
if e.args:
log.msg(e)
log.msg("Enabled item pipelines: %s" % ", ".join(self.enabled.keys()),
level=log.DEBUG)
self.loaded = True
@classmethod
def _get_mwlist_from_settings(cls, settings):
return settings.getlist('ITEM_PIPELINES')
def open_spider(self, spider):
pass
# FIXME: remove in Scrapy 0.11
def _wrap_old_process_item(self, old):
def new(item, spider):
return old(spider, item)
return new
def close_spider(self, spider):
pass
def _add_middleware(self, pipe):
super(ItemPipelineManager, self)._add_middleware(pipe)
if hasattr(pipe, 'process_item'):
# FIXME: remove in Scrapy 0.11
from scrapy.utils.python import get_func_args
if get_func_args(pipe.process_item.im_func)[1] == 'spider':
log.msg("Update %s.process_item() method to receive (item, spider) instead of (spider, item) or they will stop working on Scrapy 0.11" % pipe.__class__.__name__, log.WARNING)
pipe.process_item = self._wrap_old_process_item(pipe.process_item)
self.methods['process_item'].append(pipe.process_item)
def process_item(self, item, spider):
if not self.pipeline:
return defer_succeed(item)
def next_stage(item, stages_left):
assert isinstance(item, BaseItem), \
'Item pipelines must return a BaseItem, got %s' % type(item).__name__
if not stages_left:
return item
current_stage = stages_left.pop(0)
d = mustbe_deferred(current_stage.process_item, spider, item)
d.addCallback(next_stage, stages_left)
return d
deferred = mustbe_deferred(next_stage, item, self.pipeline[:])
return deferred
return self._process_chain('process_item', item, spider)

View File

@ -17,7 +17,7 @@ class FileExportPipeline(object):
self.exporter.start_exporting()
dispatcher.connect(self.engine_stopped, signals.engine_stopped)
def process_item(self, spider, item):
def process_item(self, item, spider):
self.exporter.export_item(item)
return item

View File

@ -1,9 +1,7 @@
from scrapy.xlib.pydispatch import dispatcher
from twisted.internet.defer import Deferred, DeferredList
from scrapy.utils.defer import mustbe_deferred, defer_result
from scrapy import log
from scrapy import signals
from scrapy.core.manager import scrapymanager
from scrapy.utils.request import request_fingerprint
from scrapy.utils.misc import arg_to_iter
@ -23,16 +21,14 @@ class MediaPipeline(object):
def __init__(self):
self.spiderinfo = {}
dispatcher.connect(self.spider_opened, signals.spider_opened)
dispatcher.connect(self.spider_closed, signals.spider_closed)
def spider_opened(self, spider):
def open_spider(self, spider):
self.spiderinfo[spider] = self.SpiderInfo(spider)
def spider_closed(self, spider):
def close_spider(self, spider):
del self.spiderinfo[spider]
def process_item(self, spider, item):
def process_item(self, item, spider):
info = self.spiderinfo[spider]
requests = arg_to_iter(self.get_media_requests(item, info))
dlist = [self._enqueue(r, info) for r in requests]

View File

@ -27,6 +27,7 @@ class ExecutionEngine(object):
def __init__(self):
self.configured = False
self.closing = {} # dict (spider -> reason) of spiders being closed
self.closing_dfds = {} # dict (spider -> deferred) of spiders being closed
self.running = False
self.killed = False
self.paused = False
@ -208,13 +209,14 @@ class ExecutionEngine(object):
dwld.addBoth(_on_complete)
return dwld
@defer.inlineCallbacks
def open_spider(self, spider):
assert self.has_capacity(), "No free spider slots when opening %r" % \
spider.name
log.msg("Spider opened", spider=spider)
self.scheduler.open_spider(spider)
self.downloader.open_spider(spider)
self.scraper.open_spider(spider)
yield self.scraper.open_spider(spider)
stats.open_spider(spider)
send_catch_log(signals.spider_opened, spider=spider)
self.next_request(spider)
@ -245,6 +247,7 @@ class ExecutionEngine(object):
self.closing[spider] = reason
self.scheduler.clear_pending_requests(spider)
dfd = self.downloader.close_spider(spider)
self.closing_dfds[spider] = dfd
dfd.addBoth(lambda _: self.scheduler.close_spider(spider))
dfd.addErrback(log.err, "Unhandled error in scheduler.close_spider()", \
spider=spider)
@ -258,6 +261,7 @@ class ExecutionEngine(object):
def _close_all_spiders(self):
dfds = [self.close_spider(s, reason='shutdown') for s in self.open_spiders]
dfds += self.closing_dfds.values()
dlist = defer.DeferredList(dfds)
return dlist
@ -275,6 +279,7 @@ class ExecutionEngine(object):
dfd.addErrback(log.err, "Unhandled error in spiders.close_spider()",
spider=spider)
dfd.addBoth(lambda _: log.msg("Spider closed (%s)" % reason, spider=spider))
dfd.addBoth(lambda _: self.closing_dfds.pop(spider).callback(spider))
dfd.addBoth(lambda _: self._spider_closed_callback(spider))
return dfd

View File

@ -63,22 +63,24 @@ class Scraper(object):
def __init__(self, engine):
self.sites = {}
self.spidermw = SpiderMiddlewareManager()
self.itemproc = load_object(settings['ITEM_PROCESSOR'])()
itemproc_cls = load_object(settings['ITEM_PROCESSOR'])
self.itemproc = itemproc_cls.from_settings(settings)
self.concurrent_items = settings.getint('CONCURRENT_ITEMS')
self.engine = engine
@defer.inlineCallbacks
def open_spider(self, spider):
"""Open the given spider for scraping and allocate resources for it"""
assert spider not in self.sites, "Spider already opened: %s" % spider
self.sites[spider] = SpiderInfo()
self.itemproc.open_spider(spider)
yield self.itemproc.open_spider(spider)
def close_spider(self, spider):
"""Close a spider being scraped and release its resources"""
assert spider in self.sites, "Spider not opened: %s" % spider
site = self.sites[spider]
site.closing = defer.Deferred()
self.itemproc.close_spider(spider)
site.closing.addCallback(self.itemproc.close_spider)
self._check_if_closing(spider, site)
return site.closing
@ -89,7 +91,7 @@ class Scraper(object):
def _check_if_closing(self, spider, site):
if site.closing and site.is_idle():
del self.sites[spider]
site.closing.callback(None)
site.closing.callback(spider)
def enqueue_scrape(self, response, request, spider):
site = self.sites.get(spider, None)

View File

@ -4,5 +4,5 @@
# See: http://doc.scrapy.org/topics/item-pipeline.html
class ${ProjectName}Pipeline(object):
def process_item(self, spider, item):
def process_item(self, item, spider):
return item

View File

@ -31,15 +31,15 @@ class MediaPipelineTestCase(unittest.TestCase):
def setUp(self):
self.spider = BaseSpider('media.com')
self.pipe = self.pipeline_class()
self.pipe.spider_opened(self.spider)
self.pipe.open_spider(self.spider)
def tearDown(self):
self.pipe.spider_closed(self.spider)
self.pipe.close_spider(self.spider)
@defer.inlineCallbacks
def test_return_item_by_default(self):
item = dict(name='sofa')
new_item = yield self.pipe.process_item(self.spider, item)
new_item = yield self.pipe.process_item(item, self.spider)
assert new_item is item
@defer.inlineCallbacks
@ -48,7 +48,7 @@ class MediaPipelineTestCase(unittest.TestCase):
info = self.pipe.spiderinfo[self.spider]
req = Request('http://media.com/2.gif')
item = dict(requests=req) # pass a single item
new_item = yield self.pipe.process_item(self.spider, item)
new_item = yield self.pipe.process_item(item, self.spider)
assert new_item is item
assert request_fingerprint(req) in info.downloaded
@ -56,7 +56,7 @@ class MediaPipelineTestCase(unittest.TestCase):
req1 = Request('http://media.com/1.gif')
req2 = Request('http://media.com/1.jpg')
item = dict(requests=iter([req1, req2]))
new_item = yield self.pipe.process_item(self.spider, item)
new_item = yield self.pipe.process_item(item, self.spider)
assert new_item is item
assert info.downloaded.get(request_fingerprint(req1)) is None
assert info.downloaded.get(request_fingerprint(req2)) is None
@ -67,7 +67,7 @@ class MediaPipelineTestCase(unittest.TestCase):
response = Response('http://media.com/2.gif')
request = Request('http://media.com/2.gif', meta=dict(response=response), callback=collected.append)
item = dict(requests=request) # pass a single item
yield self.pipe.process_item(self.spider, item)
yield self.pipe.process_item(item, self.spider)
assert collected == [response]
@defer.inlineCallbacks
@ -76,5 +76,5 @@ class MediaPipelineTestCase(unittest.TestCase):
fail = failure.Failure(Exception())
req = Request('http://media.com/2.gif', meta=dict(response=fail), callback=lambda _:_, errback=collected.append)
item = dict(requests=req) # pass a single item
yield self.pipe.process_item(self.spider, item)
yield self.pipe.process_item(item, self.spider)
assert collected == [fail]