From 43d47e5d9bdc3af1be56f5195a137755b61efcc9 Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Thu, 12 Aug 2010 10:48:37 -0300 Subject: [PATCH] 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 --- docs/intro/overview.rst | 2 +- docs/intro/tutorial.rst | 4 +- docs/topics/exporters.rst | 2 +- docs/topics/item-pipeline.rst | 42 +++++++---- .../googledir/googledir/pipelines.py | 2 +- examples/experimental/imdb/imdb/pipelines.py | 2 +- examples/googledir/googledir/pipelines.py | 2 +- scrapy/contrib/pipeline/__init__.py | 73 ++++++------------- scrapy/contrib/pipeline/fileexport.py | 2 +- scrapy/contrib/pipeline/media.py | 10 +-- scrapy/core/engine.py | 7 +- scrapy/core/scraper.py | 10 ++- .../project/module/pipelines.py.tmpl | 2 +- scrapy/tests/test_pipeline_media.py | 14 ++-- 14 files changed, 82 insertions(+), 92 deletions(-) diff --git a/docs/intro/overview.rst b/docs/intro/overview.rst index 29300fbfe..6f589cfc9 100644 --- a/docs/intro/overview.rst +++ b/docs/intro/overview.rst @@ -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) diff --git a/docs/intro/tutorial.rst b/docs/intro/tutorial.rst index bd091c445..ac8de521e 100644 --- a/docs/intro/tutorial.rst +++ b/docs/intro/tutorial.rst @@ -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 diff --git a/docs/topics/exporters.rst b/docs/topics/exporters.rst index a7f216a37..c41b63df6 100644 --- a/docs/topics/exporters.rst +++ b/docs/topics/exporters.rst @@ -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 diff --git a/docs/topics/item-pipeline.rst b/docs/topics/item-pipeline.rst index ff0d3f165..968b8d0ec 100644 --- a/docs/topics/item-pipeline.rst +++ b/docs/topics/item-pipeline.rst @@ -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: diff --git a/examples/experimental/googledir/googledir/pipelines.py b/examples/experimental/googledir/googledir/pipelines.py index 77472bcee..51e4da0e1 100644 --- a/examples/experimental/googledir/googledir/pipelines.py +++ b/examples/experimental/googledir/googledir/pipelines.py @@ -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) diff --git a/examples/experimental/imdb/imdb/pipelines.py b/examples/experimental/imdb/imdb/pipelines.py index e60714159..23740eba9 100644 --- a/examples/experimental/imdb/imdb/pipelines.py +++ b/examples/experimental/imdb/imdb/pipelines.py @@ -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 diff --git a/examples/googledir/googledir/pipelines.py b/examples/googledir/googledir/pipelines.py index 7131144a4..b842bb5cf 100644 --- a/examples/googledir/googledir/pipelines.py +++ b/examples/googledir/googledir/pipelines.py @@ -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) diff --git a/scrapy/contrib/pipeline/__init__.py b/scrapy/contrib/pipeline/__init__.py index 9e839beec..947112a06 100644 --- a/scrapy/contrib/pipeline/__init__.py +++ b/scrapy/contrib/pipeline/__init__.py @@ -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) diff --git a/scrapy/contrib/pipeline/fileexport.py b/scrapy/contrib/pipeline/fileexport.py index 76adf5150..44081c379 100644 --- a/scrapy/contrib/pipeline/fileexport.py +++ b/scrapy/contrib/pipeline/fileexport.py @@ -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 diff --git a/scrapy/contrib/pipeline/media.py b/scrapy/contrib/pipeline/media.py index b524ef6ea..5a302d8c1 100644 --- a/scrapy/contrib/pipeline/media.py +++ b/scrapy/contrib/pipeline/media.py @@ -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] diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index dad1b8149..78c6457ee 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -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 diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index c3791d61b..ec4a82cc1 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -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) diff --git a/scrapy/templates/project/module/pipelines.py.tmpl b/scrapy/templates/project/module/pipelines.py.tmpl index e3f89342d..f62608b76 100644 --- a/scrapy/templates/project/module/pipelines.py.tmpl +++ b/scrapy/templates/project/module/pipelines.py.tmpl @@ -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 diff --git a/scrapy/tests/test_pipeline_media.py b/scrapy/tests/test_pipeline_media.py index 01724c57e..d1bd62aac 100644 --- a/scrapy/tests/test_pipeline_media.py +++ b/scrapy/tests/test_pipeline_media.py @@ -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]