Added new ItemProcessor component to Scraper component

This commit is contained in:
Pablo Hoffman 2009-07-08 23:48:06 -03:00
parent 42b86a385f
commit 8b26e49636
10 changed files with 128 additions and 193 deletions

View File

@ -229,9 +229,17 @@ CONCURRENT_DOMAINS
Default: ``8``
Number of domains to scrape concurrently in one process. This doesn't affect
the number of domains scraped concurrently by the Scrapy cluster which spawns a
new process per domain.
Maximum number of domains to scrape in parallel.
.. setting:: CONCURRENT_ITEMS
CONCURRENT_ITEMS
----------------
Default: ``100``
Maximum number of concurrent items to process in parallel (per domain) in the
Item Processor (aka. Item Pipeline).
.. setting:: COOKIES_DEBUG

View File

@ -115,9 +115,9 @@ order.
:type response: :class:`~scrapy.http.Response` object
.. signal:: item_passed
.. function:: item_passed(item, spider, response, pipe_output)
.. function:: item_passed(item, spider, response, output)
Sent after an item has passed al the :ref:`topics-item-pipeline` stages without
Sent after an item has passed all the :ref:`topics-item-pipeline` stages without
being dropped.
:param item: the item which passed the pipeline
@ -129,7 +129,7 @@ order.
:param response: the response from which the item was scraped
:type response: :class:`~scrapy.http.Response` object
:param pipe_output: the output of the item pipeline. This is typically the
:param output: the output of the item pipeline. This is typically the
same :class:`~scrapy.item.ScrapedItem` object received in the ``item``
parameter, unless some pipeline stage created a new item.

View File

@ -38,7 +38,9 @@ CLUSTER_WORKER_PORT = 8789
COMMANDS_MODULE = ''
COMMANDS_SETTINGS_MODULE = ''
CONCURRENT_DOMAINS = 8 # number of domains to scrape in parallel
CONCURRENT_DOMAINS = 8
CONCURRENT_ITEMS = 100
COOKIES_DEBUG = False
@ -111,6 +113,8 @@ HTTPCACHE_IGNORE_MISSING = False
HTTPCACHE_SECTORIZE = True
HTTPCACHE_EXPIRATION_SECS = 0
ITEM_PROCESSOR = 'scrapy.item.pipeline.ItemPipelineManager'
# Item pipelines are typically set in specific commands settings
ITEM_PIPELINES = []
@ -172,7 +176,6 @@ SPIDER_MIDDLEWARES = {}
SPIDER_MIDDLEWARES_BASE = {
# Engine side
'scrapy.contrib.spidermiddleware.itempipeline.ItemPipelineMiddleware': 30,
'scrapy.contrib.spidermiddleware.httperror.HttpErrorMiddleware': 50,
'scrapy.contrib.itemsampler.ItemSamplerMiddleware': 100,
'scrapy.contrib.spidermiddleware.requestlimit.RequestLimitMiddleware': 200,

View File

@ -1,75 +0,0 @@
"""
ItemPipelineMiddleware: feed item pipeline with scraped items
"""
from pydispatch import dispatcher
from twisted.python.failure import Failure
from scrapy.core import signals
from scrapy.core.exceptions import DontCloseDomain
from scrapy.item.pipeline import ItemPipelineManager
from scrapy.item import ScrapedItem
from scrapy.conf import settings
from scrapy import log
class ItemPipelineMiddleware(object):
"""SpiderMiddleware that sends items through a pipeline"""
# The type of items to process by pipeline
ScrapedItem = ScrapedItem
# The Pipeline Manager to use for processing these item
ItemPipelineManager = ItemPipelineManager
# Maximum number of items to process in parallel by this pipeline
concurrent_limit = settings.getint('ITEMPIPELINE_CONCURRENTLIMIT', 0)
def __init__(self):
self.pipeline = self.ItemPipelineManager()
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
dispatcher.connect(self.domain_idle, signal=signals.domain_idle)
def domain_opened(self, domain):
self.pipeline.open_domain(domain)
def domain_closed(self, domain):
self.pipeline.close_domain(domain)
def domain_idle(self, domain):
if not self.pipeline.domain_is_idle(domain):
raise DontCloseDomain
def process_spider_output(self, response, result, spider):
domain = spider.domain_name
info = self.pipeline.domaininfo[domain]
for item_or_request in result:
# return to engine until pipeline frees up some slots
# TODO: this is ugly, a proper flow control mechanism should be
# added instead
while 0 < self.concurrent_limit <= len(info):
yield None
if isinstance(item_or_request, self.ScrapedItem):
log.msg("Scraped %s in <%s>" % (item_or_request, response.request.url), \
domain=domain)
signals.send_catch_log(signal=signals.item_scraped, sender=self.__class__, \
item=item_or_request, spider=spider, response=response)
self.pipeline.pipe(item_or_request, spider).addBoth(self._pipeline_finished, \
item_or_request, spider)
# yielding here breaks the loop and allows the engine to run
# other tasks, such as attending IO (very important)
yield None
else:
yield item_or_request
def _pipeline_finished(self, pipe_result, item, spider):
# exception can only be of DropItem type here, since other exceptions
# are caught in the Item Pipeline (item/pipeline.py)
if isinstance(pipe_result, Failure):
signals.send_catch_log(signal=signals.item_dropped, \
sender=self.__class__, item=item, spider=spider, exception=pipe_result.value)
else:
signals.send_catch_log(signal=signals.item_passed, \
sender=self.__class__, item=item, spider=spider, pipe_output=pipe_result)

View File

@ -36,13 +36,13 @@ class ExecutionEngine(object):
self.control_reactor = True
self._next_request_pending = set()
def configure(self, scheduler=None, downloader=None):
def configure(self):
"""
Configure execution engine with the given scheduling policy and downloader.
"""
self.scheduler = scheduler or Scheduler()
self.scheduler = load_object(settings['SCHEDULER'])()
self.domain_scheduler = load_object(settings['DOMAIN_SCHEDULER'])()
self.downloader = downloader or Downloader()
self.downloader = Downloader()
self.scraper = Scraper(self)
self.configured = True
@ -161,9 +161,8 @@ class ExecutionEngine(object):
if self.paused:
return reactor.callLater(5, self.next_request, domain)
if not self.running or \
self.domain_is_closed(domain) or \
self.downloader.sites[domain].needs_backout() or \
if not self.running or self.domain_is_closed(domain) or \
self.downloader.sites[domain].needs_backout() or \
self.scraper.sites[domain].needs_backout():
return
@ -222,7 +221,6 @@ class ExecutionEngine(object):
if not self.running or self.paused:
return
# main domain starter loop
while self.running and self.downloader.has_capacity():
if not self.next_domain():
return self._stop_if_idle()
@ -349,8 +347,10 @@ class ExecutionEngine(object):
"self.downloader.sites[domain].closing",
"self.downloader.sites[domain].lastseen",
"len(self.scraper.sites[domain].queue)",
"len(self.scraper.sites[domain].processing)",
"self.scraper.sites[domain].backlog_size",
"len(self.scraper.sites[domain].active)",
"self.scraper.sites[domain].active_size",
"self.scraper.sites[domain].itemproc_size",
"self.scraper.sites[domain].needs_backout()",
]
for test in global_tests:

View File

@ -7,7 +7,7 @@ from scrapy import log
from scrapy.http import Request
from scrapy.core.engine import scrapyengine
from scrapy.spider import spiders
from scrapy.utils.misc import load_object, arg_to_iter
from scrapy.utils.misc import arg_to_iter
from scrapy.utils.url import is_url
from scrapy.conf import settings
@ -39,8 +39,7 @@ class ExecutionManager(object):
log.msg("Enabled extensions: %s" % ", ".join(extensions.enabled.iterkeys()),
level=log.DEBUG)
scheduler = load_object(settings['SCHEDULER'])()
scrapyengine.configure(scheduler=scheduler)
scrapyengine.configure()
def crawl(self, *args):
"""Schedule the given args for crawling. args is a list of urls or domains"""

View File

@ -1,64 +1,66 @@
"""This module implements the Scraper component which parses responses and
extracts information from them"""
from itertools import imap
from twisted.internet import task
from twisted.python.failure import Failure
from twisted.internet import defer
from scrapy.utils.defer import defer_result
from scrapy.utils.misc import arg_to_iter
from scrapy.core.exceptions import IgnoreRequest
from scrapy.utils.defer import defer_result, defer_succeed, parallel
from scrapy.utils.misc import arg_to_iter, load_object
from scrapy.core.exceptions import IgnoreRequest, DropItem
from scrapy.core import signals
from scrapy.http import Request, Response
from scrapy.item import ScrapedItem
from scrapy.spider.middleware import SpiderMiddlewareManager
from scrapy import log
from scrapy.stats import stats
from scrapy.conf import settings
class SiteInfo(object):
"""Object for holding data of the responses being scraped"""
FAILURE_SIZE = 1024 # make failures equivalent to 1K responses in size
def __init__(self, max_backlog_size=5000000):
def __init__(self, max_active_size=5000000):
self.max_active_size = max_active_size
self.queue = []
self.processing = set()
self.backlog_size = 0
self.max_backlog_size = max_backlog_size
self.active = set()
self.active_size = 0
self.itemproc_size = 0
def add_response_request(self, response, request):
deferred = defer.Deferred()
self.queue.append((response, request, deferred))
if isinstance(response, Response):
self.backlog_size += len(response.body)
self.active_size += len(response.body)
else:
self.backlog_size += self.FAILURE_SIZE
self.active_size += self.FAILURE_SIZE
return deferred
def next_response_request_deferred(self):
response, request, deferred = self.queue.pop(0)
self.processing.add(response)
self.active.add(response)
return response, request, deferred
def finish_response(self, response):
self.processing.remove(response)
self.active.remove(response)
if isinstance(response, Response):
self.backlog_size -= len(response.body)
self.active_size -= len(response.body)
else:
self.backlog_size -= self.FAILURE_SIZE
self.active_size -= self.FAILURE_SIZE
def is_idle(self):
return not (self.queue or self.processing)
return not (self.queue or self.active)
def needs_backout(self):
return self.backlog_size > self.max_backlog_size
return self.active_size > self.max_active_size
class Scraper(object):
def __init__(self, engine):
self.sites = {}
self.middleware = SpiderMiddlewareManager()
self.spidermw = SpiderMiddlewareManager()
self.itemproc = load_object(settings['ITEM_PROCESSOR'])()
self.concurrent_items = settings.getint('CONCURRENT_ITEMS')
self.engine = engine
def open_domain(self, domain):
@ -66,12 +68,14 @@ class Scraper(object):
if domain in self.sites:
raise RuntimeError('Scraper domain already opened: %s' % domain)
self.sites[domain] = SiteInfo()
self.itemproc.open_domain(domain)
def close_domain(self, domain):
"""Close a domain being scraped and release its resources"""
if domain not in self.sites:
raise RuntimeError('Scraper domain already closed: %s' % domain)
del self.sites[domain]
self.sites.pop(domain)
self.itemproc.open_domain(domain)
def is_idle(self):
"""Return True if there isn't any more spiders to process"""
@ -82,19 +86,15 @@ class Scraper(object):
dfd = site.add_response_request(response, request)
def finish_scraping(_):
site.finish_response(response)
self._scrape_next(spider, site)
return _
dfd.addBoth(finish_scraping)
dfd.addErrback(log.err, 'Scraper bug processing %s' % request, \
domain=spider.domain_name)
self.scrape_next(spider)
self._scrape_next(spider, site)
return dfd
def scrape_next(self, spider):
site = self.sites.get(spider.domain_name)
if not site:
return
# Process responses in queue
def _scrape_next(self, spider, site):
while site.queue:
response, request, deferred = site.next_response_request_deferred()
self._scrape(response, request, spider).chainDeferred(deferred)
@ -106,14 +106,14 @@ class Scraper(object):
dfd = self._scrape2(response, request, spider) # returns spiders processed output
dfd.addErrback(self.handle_spider_error, request, spider)
dfd.addCallback(self.handle_spider_output, request, spider)
dfd.addCallback(self.handle_spider_output, request, response, spider)
return dfd
def _scrape2(self, request_result, request, spider):
"""Handle the diferent cases of request's result been a Response or a
Failure"""
if not isinstance(request_result, Failure):
return self.middleware.scrape_response(self.call_spider, \
return self.spidermw.scrape_response(self.call_spider, \
request_result, request, spider)
else:
# FIXME: don't ignore errors in spider middleware
@ -132,22 +132,39 @@ class Scraper(object):
stats.incpath("%s/spider_exceptions/%s" % (spider.domain_name, \
_failure.value.__class__.__name__))
def handle_spider_output(self, result, request, spider):
func = lambda o: self.process_spider_output(o, request, spider)
return task.coiterate(imap(func, result or []))
def handle_spider_output(self, result, request, response, spider):
domain = spider.domain_name
if not result:
return defer_succeed(None)
dfd = parallel(iter(result), self.concurrent_items,
self._process_spidermw_output, request, response, spider)
return dfd
def process_spider_output(self, output, request, spider):
def _process_spidermw_output(self, output, request, response, spider):
"""Process each Request/Item (given in the output parameter) returned
from the given spider
"""
# TODO: keep closing state internally instead of checking engine
if spider.domain_name in self.engine.closing:
domain = spider.domain_name
if domain in self.engine.closing:
return
elif isinstance(output, Request):
signals.send_catch_log(signal=signals.request_received, request=output, spider=spider)
signals.send_catch_log(signal=signals.request_received, request=output, \
spider=spider)
self.engine.crawl(request=output, spider=spider)
elif isinstance(output, ScrapedItem):
log.msg("Scraped %s in <%s>" % (output, request.url), domain=domain)
signals.send_catch_log(signal=signals.item_scraped, sender=self.__class__, \
item=output, spider=spider, response=response)
self.sites[domain].itemproc_size += 1
dfd = self.itemproc.process_item(output, spider)
dfd.addBoth(self._itemproc_finished, output, response, spider)
return dfd
elif output is None:
pass # may be next time.
pass
else:
log.msg("Spider must return Request, ScrapedItem or None, got '%s' while processing %s" \
% (type(output).__name__, request), log.WARNING, domain=spider.domain_name)
log.msg("Spider must return Request, ScrapedItem or None, got '%s' in %s" % \
(type(output).__name__, request), log.ERROR, domain=domain)
def _check_propagated_failure(self, spider_failure, propagated_failure, request, spider):
"""Log and silence the bugs raised outside of spiders, but still allow
@ -162,3 +179,22 @@ class Scraper(object):
return # stop propagating this error
else:
return spider_failure # exceptions raised in the spider code
def _itemproc_finished(self, output, item, response, spider):
"""ItemProcessor finished for the given ``item`` and returned ``output``
"""
domain = spider.domain_name
self.sites[domain].itemproc_size -= 1
if isinstance(output, Failure):
ex = output.value
if isinstance(ex, DropItem):
log.msg("Dropped %s - %s" % (item, str(ex)), log.DEBUG, domain=domain)
signals.send_catch_log(signal=signals.item_dropped, \
sender=self.__class__, item=item, spider=spider, exception=output.value)
else:
log.msg('Error processing %s - %s' % (item, output), \
log.ERROR, domain=domain)
else:
signals.send_catch_log(signal=signals.item_passed, \
sender=self.__class__, item=item, spider=spider, output=output)

View File

@ -1,5 +1,5 @@
from scrapy import log
from scrapy.core.exceptions import DropItem, NotConfigured
from scrapy.core.exceptions import NotConfigured
from scrapy.item import ScrapedItem
from scrapy.utils.misc import load_object
from scrapy.utils.defer import defer_succeed, mustbe_deferred
@ -10,7 +10,6 @@ class ItemPipelineManager(object):
def __init__(self):
self.loaded = False
self.pipeline = []
self.domaininfo = {}
self.load()
def load(self):
@ -30,69 +29,24 @@ class ItemPipelineManager(object):
self.loaded = True
def open_domain(self, domain):
self.domaininfo[domain] = set()
pass
def close_domain(self, domain):
del self.domaininfo[domain]
pass
def is_idle(self):
return not self.domaininfo
def domain_is_idle(self, domain):
return not self.domaininfo.get(domain)
def pipe(self, item, spider):
"""
item pipelines are instanceable classes that defines a `pipeline` method
that takes ScrapedItem as input and returns ScrapedItem.
The output from one stage is the input of the next.
Raising DropItem stops pipeline.
This pipeline is configurable with the ITEM_PIPELINES setting
"""
domain = spider.domain_name
if not self.pipeline or domain not in self.domaininfo:
def process_item(self, item, spider):
if not self.pipeline:
return defer_succeed(item)
pipeline = self.pipeline[:]
current_stage = pipeline[0]
info = self.domaininfo[domain]
info.add(item)
def _next_stage(item):
def next_stage(item, stages_left):
assert isinstance(item, ScrapedItem), \
'Pipeline stages must return a ScrapedItem or raise DropItem, got %s' % type(item).__name__
if not pipeline:
'Item pipelines must return a ScrapedItem, got %s' % type(item).__name__
if not stages_left:
return item
current_stage = pipeline.pop(0)
log.msg("_%s_ Pipeline stage: %s" % (item, type(current_stage).__name__), log.TRACE, domain=domain)
d = mustbe_deferred(current_stage.process_item, domain, item)
d.addCallback(_next_stage)
current_stage = stages_left.pop(0)
d = mustbe_deferred(current_stage.process_item, spider.domain_name, item)
d.addCallback(next_stage, stages_left)
return d
def _ondrop(_failure):
ex = _failure.value
if isinstance(ex, DropItem):
# TODO: current_stage is not working, check why
#log.msg("%s: Dropped %s - %s" % (type(current_stage).__name__, item, str(ex)), log.DEBUG, domain=domain)
log.msg("Dropped %s - %s" % (item, str(ex)), log.DEBUG, domain=domain)
return _failure
else:
# TODO: current_stage is not working, check why
#log.msg('%s: Error processing %s - %s' % (type(current_stage).__name__, item, _failure), log.ERROR, domain=domain)
log.msg('Error processing %s - %s' % (item, _failure), log.ERROR, domain=domain)
def _pipeline_finished(_):
log.msg("_%s_ Pipeline finished" % item, log.TRACE, domain=domain)
info.remove(item)
return _
deferred = mustbe_deferred(_next_stage, item)
deferred.addErrback(_ondrop)
deferred.addBoth(_pipeline_finished)
deferred = mustbe_deferred(next_stage, item, self.pipeline[:])
return deferred

View File

@ -41,7 +41,7 @@ class CoreStats(object):
stats.incpath('%s/item_scraped_count' % spider.domain_name)
stats.incpath('_global/item_scraped_count')
def item_passed(self, item, spider, pipe_output):
def item_passed(self, item, spider):
stats.incpath('%s/item_passed_count' % spider.domain_name)
stats.incpath('_global/item_passed_count')

View File

@ -2,7 +2,7 @@
Helper functions for dealing with Twisted deferreds
"""
from twisted.internet import defer, reactor
from twisted.internet import defer, reactor, task
from twisted.python import failure
def defer_fail(_failure):
@ -37,3 +37,13 @@ def mustbe_deferred(f, *args, **kw):
def chain_deferred(d1, d2):
return d1.chainDeferred(d2).addBoth(lambda _:d2)
def parallel(iterable, count, callable, *args, **named):
"""Execute a callable over the objects in the given iterable, in parallel,
using no more than ``count`` concurrent calls.
Taken from: http://jcalderone.livejournal.com/24285.html
"""
coop = task.Cooperator()
work = (callable(elem, *args, **named) for elem in iterable)
return defer.DeferredList([coop.coiterate(work) for i in xrange(count)])