Changed item pipeline API to pass spider references (instead of domain names) to process_item() method

This commit is contained in:
Pablo Hoffman 2009-11-06 13:46:36 -02:00
parent a432c1ee40
commit 7728a23e99
9 changed files with 64 additions and 65 deletions

View File

@ -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)

View File

@ -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 <http://www.dmoz.org/Computers/Programming/Languages/Python/Resources/> from <None>
[dmoz.org] DEBUG: Crawled <http://www.dmoz.org/Computers/Programming/Languages/Python/Books/> from <None>
[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 <None>``.
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 <None>``.
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

View File

@ -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

View File

@ -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

View File

@ -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

View File

@ -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

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, domain, item):
def process_item(self, spider, item):
self.exporter.export_item(item)
return item

View File

@ -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

View File

@ -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)