From 7728a23e9936e387ee5dca01d937f6e512e5d7ab Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Fri, 6 Nov 2009 13:46:36 -0200 Subject: [PATCH] Changed item pipeline API to pass spider references (instead of domain names) to process_item() method --- docs/intro/overview.rst | 2 +- docs/intro/tutorial.rst | 16 +++++------ docs/topics/exporters.rst | 10 +++---- docs/topics/item-pipeline.rst | 24 ++++++++-------- scrapy/contrib/itemsampler.py | 36 ++++++++++++------------ scrapy/contrib/pipeline/__init__.py | 2 +- scrapy/contrib/pipeline/fileexport.py | 2 +- scrapy/contrib/pipeline/media.py | 15 +++++----- scrapy/contrib_exp/pipeline/shoveitem.py | 22 +++++++-------- 9 files changed, 64 insertions(+), 65 deletions(-) diff --git a/docs/intro/overview.rst b/docs/intro/overview.rst index 15b4efd45..7a38d4fcd 100644 --- a/docs/intro/overview.rst +++ b/docs/intro/overview.rst @@ -153,7 +153,7 @@ extracted item into a file using `pickle`_:: import pickle class StoreItemPipeline(object): - def process_item(self, domain, response, item): + def process_item(self, spider, item): 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 b5934fa0f..3fb965110 100644 --- a/docs/intro/tutorial.rst +++ b/docs/intro/tutorial.rst @@ -156,17 +156,17 @@ will get an output similar to this:: [dmoz] INFO: Enabled downloader middlewares: ... [dmoz] INFO: Enabled spider middlewares: ... [dmoz] INFO: Enabled item pipelines: ... - [dmoz.org] INFO: Domain opened + [dmoz.org] INFO: Spider opened [dmoz.org] DEBUG: Crawled from [dmoz.org] DEBUG: Crawled from - [dmoz.org] INFO: Domain closed (finished) + [dmoz.org] INFO: Spider closed (finished) [-] Main loop terminated. Pay attention to the lines containing ``[dmoz.org]``, which corresponds to -our spider (identified by the domain "dmoz.org"). You can see a log line for each -URL defined in ``start_urls``. Because these URLs are the starting ones, they -have no referrers, which is shown at the end of the log line, where it says -``from ``. +our spider (identified by the domain ``"dmoz.org"``). You can see a log line +for each URL defined in ``start_urls``. Because these URLs are the starting +ones, they have no referrers, which is shown at the end of the log line, +where it says ``from ``. But more interesting, as our ``parse`` method instructs, two files have been created: *Books* and *Resources*, with the content of both URLs. @@ -445,7 +445,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, domain, item): + def process_item(self, spider, item): return item We have to override the ``process_item`` method in order to store our Items @@ -461,7 +461,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, domain, item): + def process_item(self, spider, item): 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 b9182223f..966724f89 100644 --- a/docs/topics/exporters.rst +++ b/docs/topics/exporters.rst @@ -51,19 +51,17 @@ Exporter to export scraped items to different files, one per spider:: self.files = {} def spider_opened(self, spider): - domain = spider.domain_name - file = open('%s_products.xml' % domain, 'w+b') - self.files[domain] = file + file = open('%s_products.xml' % spider.domain_name, 'w+b') + self.files[spider] = file self.exporter = XmlItemExporter(file) self.exporter.start_exporting() def spider_closed(self, spider): - domain = spider.domain_name self.exporter.finish_exporting() - file = self.files.pop(domain) + file = self.files.pop(spider) file.close() - def process_item(self, domain, item): + def process_item(self, spider, item): self.exporter.export_item(item) return item diff --git a/docs/topics/item-pipeline.rst b/docs/topics/item-pipeline.rst index 3fb90ce39..ce465577f 100644 --- a/docs/topics/item-pipeline.rst +++ b/docs/topics/item-pipeline.rst @@ -24,11 +24,13 @@ 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: -.. method:: process_item(domain, item) +.. method:: process_item(spider, item) -``domain`` is a string with the domain of the spider which scraped the item +:param spider: the spider which scraped the item +:type spider: :class:`~scrapy.spider.BaseSpider` object -``item`` is a :class:`~scrapy.item.Item` with the item scraped +: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 @@ -49,7 +51,7 @@ attribute), and drops those items which don't contain a price:: vat_factor = 1.15 - def process_item(self, domain, item): + def process_item(self, spider, item): if item['price']: if item['price_excludes_vat']: item['price'] = item['price'] * self.vat_factor @@ -68,11 +70,11 @@ To activate an Item Pipeline component you must add its class to the 'myproject.pipeline.PricePipeline', ] -Item pipeline example with resources per domain +Item pipeline example with resources per spider =============================================== Sometimes you need to keep resources about the items processed grouped per -domain, and delete those resource when a domain finish. +spider, and delete those resource when a spider finish. An example is a filter that looks for duplicate items, and drops those items that were already processed. Let say that our items has an unique id, but our @@ -90,16 +92,16 @@ spider returns multiples items with the same id:: dispatcher.connect(self.spider_closed, signals.spider_closed) def spider_opened(self, spider): - self.duplicates[spider.domain_name] = set() + self.duplicates[spider] = set() def spider_closed(self, spider): - del self.duplicates[spider.domain_name] + del self.duplicates[spider] - def process_item(self, domain, item): - if item.id in self.duplicates[domain]: + def process_item(self, spider, item): + if item.id in self.duplicates[spider]: raise DropItem("Duplicate item found: %s" % item) else: - self.duplicates[domain].add(item.id) + self.duplicates[spider].add(item.id) return item Built-in Item Pipelines reference diff --git a/scrapy/contrib/itemsampler.py b/scrapy/contrib/itemsampler.py index 026b4288a..60cd403f3 100644 --- a/scrapy/contrib/itemsampler.py +++ b/scrapy/contrib/itemsampler.py @@ -36,8 +36,8 @@ from scrapy.http import Request from scrapy import log from scrapy.conf import settings -items_per_domain = settings.getint('ITEMSAMPLER_COUNT', 1) -close_domain = settings.getbool('ITEMSAMPLER_CLOSE_DOMAIN', False) +items_per_spider = settings.getint('ITEMSAMPLER_COUNT', 1) +close_spider = settings.getbool('ITEMSAMPLER_CLOSE_SPIDER', False) max_response_size = settings.getbool('ITEMSAMPLER_MAX_RESPONSE_SIZE', ) class ItemSamplerPipeline(object): @@ -47,20 +47,19 @@ class ItemSamplerPipeline(object): if not self.filename: raise NotConfigured self.items = {} - self.domains_count = 0 + self.spiders_count = 0 self.empty_domains = set() dispatcher.connect(self.spider_closed, signal=signals.spider_closed) dispatcher.connect(self.engine_stopped, signal=signals.engine_stopped) - def process_item(self, item, spider): - domain = spider.domain_name - sampled = stats.get_value("items_sampled", 0, domain=domain) - if sampled < items_per_domain: + def process_item(self, spider, item): + sampled = stats.get_value("items_sampled", 0, domain=spider.domain_name) + if sampled < items_per_spider: self.items[item.guid] = item sampled += 1 - stats.set_value("items_sampled", sampled, domain=domain) - log.msg("Sampled %s" % item, domain=domain, level=log.INFO) - if close_domain and sampled == items_per_domain: + stats.set_value("items_sampled", sampled, domain=spider.domain_name) + log.msg("Sampled %s" % item, spider=spider, level=log.INFO) + if close_spider and sampled == items_per_spider: scrapyengine.close_spider(spider) return item @@ -68,14 +67,15 @@ class ItemSamplerPipeline(object): with open(self.filename, 'w') as f: pickle.dump(self.items, f) if self.empty_domains: - log.msg("No products sampled for: %s" % " ".join(self.empty_domains), level=log.WARNING) + log.msg("No products sampled for: %s" % " ".join(self.empty_domains), \ + level=log.WARNING) def spider_closed(self, spider, reason): - domain = spider.domain_name - if reason == 'finished' and not stats.get_value("items_sampled", domain=domain): - self.empty_domains.add(domain) - self.domains_count += 1 - log.msg("Sampled %d domains so far (%d empty)" % (self.domains_count, len(self.empty_domains)), level=log.INFO) + if reason == 'finished' and not stats.get_value("items_sampled", domain=spider.domain_name): + self.empty_domains.add(spider.domain_name) + self.spiders_count += 1 + log.msg("Sampled %d domains so far (%d empty)" % (self.spiders_count, \ + len(self.empty_domains)), level=log.INFO) class ItemSamplerMiddleware(object): @@ -87,7 +87,7 @@ class ItemSamplerMiddleware(object): raise NotConfigured def process_spider_input(self, response, spider): - if stats.get_value("items_sampled", domain=spider.domain_name) >= items_per_domain: + if stats.get_value("items_sampled", domain=spider.domain_name) >= items_per_spider: return [] elif max_response_size and max_response_size > len(response_httprepr(response)): return [] @@ -100,7 +100,7 @@ class ItemSamplerMiddleware(object): else: items.append(r) - if stats.get_value("items_sampled", domain=spider.domain_name) >= items_per_domain: + if stats.get_value("items_sampled", domain=spider.domain_name) >= items_per_spider: return [] else: # TODO: this needs some revision, as keeping only the first item diff --git a/scrapy/contrib/pipeline/__init__.py b/scrapy/contrib/pipeline/__init__.py index 02d6b83cd..5801e86fc 100644 --- a/scrapy/contrib/pipeline/__init__.py +++ b/scrapy/contrib/pipeline/__init__.py @@ -57,7 +57,7 @@ class ItemPipelineManager(object): if not stages_left: return item current_stage = stages_left.pop(0) - d = mustbe_deferred(current_stage.process_item, spider.domain_name, item) + d = mustbe_deferred(current_stage.process_item, spider, item) d.addCallback(next_stage, stages_left) return d diff --git a/scrapy/contrib/pipeline/fileexport.py b/scrapy/contrib/pipeline/fileexport.py index 461d6e492..d20b83c2a 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, domain, item): + def process_item(self, spider, item): self.exporter.export_item(item) return item diff --git a/scrapy/contrib/pipeline/media.py b/scrapy/contrib/pipeline/media.py index 16e11bef6..96c804be3 100644 --- a/scrapy/contrib/pipeline/media.py +++ b/scrapy/contrib/pipeline/media.py @@ -12,27 +12,26 @@ from scrapy.utils.misc import arg_to_iter class MediaPipeline(object): DOWNLOAD_PRIORITY = 1000 - class DomainInfo(object): + class SpiderInfo(object): def __init__(self, spider): - self.domain = spider.domain_name self.spider = spider self.downloading = {} self.downloaded = {} self.waiting = {} def __init__(self): - self.domaininfo = {} + self.spiderinfo = {} dispatcher.connect(self.spider_opened, signals.spider_opened) dispatcher.connect(self.spider_closed, signals.spider_closed) def spider_opened(self, spider): - self.domaininfo[spider.domain_name] = self.DomainInfo(spider) + self.spiderinfo[spider] = self.SpiderInfo(spider) def spider_closed(self, spider): - del self.domaininfo[spider.domain_name] + del self.spiderinfo[spider] - def process_item(self, domain, item): - info = self.domaininfo[domain] + def process_item(self, spider, item): + info = self.spiderinfo[spider] requests = arg_to_iter(self.get_media_requests(item, info)) dlist = [] for request in requests: @@ -83,7 +82,7 @@ class MediaPipeline(object): info.downloading[fp] = (request, dwld) # fill downloading state data dwld.addBoth(_downloaded) # append post-download hook - dwld.addErrback(log.err, domain=info.domain) + dwld.addErrback(log.err, spider=info.spider) # declare request in downloading state (None is used as place holder) info.downloading[fp] = None diff --git a/scrapy/contrib_exp/pipeline/shoveitem.py b/scrapy/contrib_exp/pipeline/shoveitem.py index 334c07993..869c62c58 100644 --- a/scrapy/contrib_exp/pipeline/shoveitem.py +++ b/scrapy/contrib_exp/pipeline/shoveitem.py @@ -27,11 +27,11 @@ class ShoveItemPipeline(object): dispatcher.connect(self.spider_opened, signal=signals.spider_opened) dispatcher.connect(self.spider_closed, signal=signals.spider_closed) - def process_item(self, domain, item): + def process_item(self, spider, item): guid = str(item.guid) - if guid in self.stores[domain]: - if self.stores[domain][guid] == item: + if guid in self.stores[spider]: + if self.stores[spider][guid] == item: status = 'old' else: status = 'upd' @@ -39,17 +39,17 @@ class ShoveItemPipeline(object): status = 'new' if not status == 'old': - self.stores[domain][guid] = item - self.log(domain, item, status) + self.stores[spider][guid] = item + self.log(spider, item, status) return item def spider_opened(self, spider): - domain = spider.domain_name - uri = Template(self.uritpl).substitute(domain=domain) - self.stores[domain] = Shove(uri, **self.opts) + uri = Template(self.uritpl).substitute(domain=spider.domain_name) + self.stores[spider] = Shove(uri, **self.opts) def spider_closed(self, spider): - self.stores[spider.domain_name].sync() + self.stores[spider].sync() - def log(self, domain, item, status): - log.msg("Shove (%s): Item guid=%s" % (status, item.guid), level=log.DEBUG, domain=domain) + def log(self, spider, item, status): + log.msg("Shove (%s): Item guid=%s" % (status, item.guid), level=log.DEBUG, \ + spider=spider)