diff --git a/docs/ref/settings.rst b/docs/ref/settings.rst index 99022b2ff..2221f1270 100644 --- a/docs/ref/settings.rst +++ b/docs/ref/settings.rst @@ -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 diff --git a/docs/ref/signals.rst b/docs/ref/signals.rst index d6f60c4c2..b60936396 100644 --- a/docs/ref/signals.rst +++ b/docs/ref/signals.rst @@ -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. diff --git a/scrapy/conf/default_settings.py b/scrapy/conf/default_settings.py index 71c8375c2..6db530ace 100644 --- a/scrapy/conf/default_settings.py +++ b/scrapy/conf/default_settings.py @@ -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, diff --git a/scrapy/contrib/spidermiddleware/itempipeline.py b/scrapy/contrib/spidermiddleware/itempipeline.py deleted file mode 100644 index 66a62a750..000000000 --- a/scrapy/contrib/spidermiddleware/itempipeline.py +++ /dev/null @@ -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) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index da0b60916..acc933980 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -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: diff --git a/scrapy/core/manager.py b/scrapy/core/manager.py index 4c9dcd068..d8fb44565 100644 --- a/scrapy/core/manager.py +++ b/scrapy/core/manager.py @@ -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""" diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index b3b6a4ef4..79b1786b7 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -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) + diff --git a/scrapy/item/pipeline.py b/scrapy/item/pipeline.py index edb535300..ae18d979e 100644 --- a/scrapy/item/pipeline.py +++ b/scrapy/item/pipeline.py @@ -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 diff --git a/scrapy/stats/corestats.py b/scrapy/stats/corestats.py index 88b0ef2e8..200993daf 100644 --- a/scrapy/stats/corestats.py +++ b/scrapy/stats/corestats.py @@ -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') diff --git a/scrapy/utils/defer.py b/scrapy/utils/defer.py index de24eecc1..12bb8b653 100644 --- a/scrapy/utils/defer.py +++ b/scrapy/utils/defer.py @@ -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)]) +