From 31b3d7ce1e60a5c91d8b9bfd1bd3d49113b0a73b Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Mon, 6 Jul 2009 15:31:50 -0300 Subject: [PATCH] Added flow control mechanism to new Scraper component, to prevent cases where memory fills because of requests being downloaded much faster than they can be processed (by the spiders) --- .../contrib/downloadermiddleware/robotstxt.py | 2 +- scrapy/contrib/pipeline/media.py | 2 +- scrapy/contrib/pipeline/s3images.py | 2 +- scrapy/contrib/web/service.py | 2 +- scrapy/core/engine.py | 16 ++-- scrapy/core/scraper.py | 89 +++++++++++++++---- 6 files changed, 87 insertions(+), 26 deletions(-) diff --git a/scrapy/contrib/downloadermiddleware/robotstxt.py b/scrapy/contrib/downloadermiddleware/robotstxt.py index b7948a46c..8e4295804 100644 --- a/scrapy/contrib/downloadermiddleware/robotstxt.py +++ b/scrapy/contrib/downloadermiddleware/robotstxt.py @@ -45,7 +45,7 @@ class RobotsTxtMiddleware(object): self._parsers[urldomain] = None robotsurl = "%s://%s/robots.txt" % parsedurl[0:2] robotsreq = Request(robotsurl, priority=self.DOWNLOAD_PRIORITY) - dfd = scrapyengine.schedule(robotsreq, spiders.fromdomain(spiderdomain)) + dfd = scrapyengine.download(robotsreq, spiders.fromdomain(spiderdomain)) dfd.addCallbacks(callback=self._parse_robots, callbackArgs=[urldomain]) self._spiderdomains[spiderdomain].add(urldomain) diff --git a/scrapy/contrib/pipeline/media.py b/scrapy/contrib/pipeline/media.py index 02ebe9f82..7fbdddbf7 100644 --- a/scrapy/contrib/pipeline/media.py +++ b/scrapy/contrib/pipeline/media.py @@ -123,7 +123,7 @@ class MediaPipeline(object): """ request.priority = self.DOWNLOAD_PRIORITY - return scrapyengine.schedule(request, info.spider) + return scrapyengine.download(request, info.spider) def media_to_download(self, request, info): """ Ongoing request hook pre-cache diff --git a/scrapy/contrib/pipeline/s3images.py b/scrapy/contrib/pipeline/s3images.py index 0bde3be67..2a50a9ca1 100644 --- a/scrapy/contrib/pipeline/s3images.py +++ b/scrapy/contrib/pipeline/s3images.py @@ -94,7 +94,7 @@ class S3ImagesPipeline(BaseImagesPipeline): """ if self.AmazonS3Spider: - return scrapyengine.schedule(request, self.AmazonS3Spider) + return scrapyengine.download(request, self.AmazonS3Spider) return self.download(request, info) diff --git a/scrapy/contrib/web/service.py b/scrapy/contrib/web/service.py index b2b51bc37..7859a6e4e 100644 --- a/scrapy/contrib/web/service.py +++ b/scrapy/contrib/web/service.py @@ -68,7 +68,7 @@ def url_to_guid(httprequest): deferred = defer.Deferred().addCallbacks(_on_success, _on_error) request = Request(url=url, callback=deferred, dont_filter=True) - schd = scrapyengine.schedule(request, spider) + schd = scrapyengine.download(request, spider) schd.chainDeferred(deferred) return deferred diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index fbc3ec010..da0b60916 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -113,7 +113,8 @@ class ExecutionEngine(object): self.running = False for domain in self.open_domains: spider = spiders.fromdomain(domain) - signals.send_catch_log(signal=signals.domain_closed, sender=self.__class__, domain=domain, spider=spider, reason='shutdown') + signals.send_catch_log(signal=signals.domain_closed, sender=self.__class__, \ + domain=domain, spider=spider, reason='shutdown') for tsk, _, _ in self.tasks: # stop looping calls if tsk.running: tsk.stop() @@ -162,7 +163,8 @@ class ExecutionEngine(object): if not self.running or \ self.domain_is_closed(domain) or \ - self.downloader.sites[domain].needs_backout(): + self.downloader.sites[domain].needs_backout() or \ + self.scraper.sites[domain].needs_backout(): return # Next pending request from scheduler @@ -176,10 +178,10 @@ class ExecutionEngine(object): self._domain_idle(domain) def domain_is_idle(self, domain): - scraping = domain in self.scraper.sites and bool(self.scraper.sites[domain]) + scraper_idle = domain in self.scraper.sites and self.scraper.sites[domain].is_idle() pending = self.scheduler.domain_has_pending_requests(domain) downloading = domain in self.downloader.sites and self.downloader.sites[domain].active - return not (pending or downloading or scraping) + return scraper_idle and not (pending or downloading) def domain_is_closed(self, domain): """Return True if the domain is fully closed (ie. not even in the @@ -197,7 +199,7 @@ class ExecutionEngine(object): def crawl(self, request, spider): schd = mustbe_deferred(self.schedule, request, spider) - schd.addBoth(self.scraper.scrape, request, spider) + schd.addBoth(self.scraper.enqueue_scrape, request, spider) schd.addErrback(log.err) schd.addBoth(lambda _: self.next_request(spider.domain_name)) @@ -346,7 +348,9 @@ class ExecutionEngine(object): "len(self.downloader.sites[domain].transferring)", "self.downloader.sites[domain].closing", "self.downloader.sites[domain].lastseen", - "len(self.scraper.sites[domain])", + "len(self.scraper.sites[domain].queue)", + "len(self.scraper.sites[domain].processing)", + "self.scraper.sites[domain].backlog_size", ] for test in global_tests: diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index ad239493a..b3b6a4ef4 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -1,8 +1,11 @@ -"""Extract information from pages""" +"""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 @@ -13,6 +16,44 @@ from scrapy.spider.middleware import SpiderMiddlewareManager from scrapy import log from scrapy.stats import stats +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): + self.queue = [] + self.processing = set() + self.backlog_size = 0 + self.max_backlog_size = max_backlog_size + + 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) + else: + self.backlog_size += self.FAILURE_SIZE + return deferred + + def next_response_request_deferred(self): + response, request, deferred = self.queue.pop(0) + self.processing.add(response) + return response, request, deferred + + def finish_response(self, response): + self.processing.remove(response) + if isinstance(response, Response): + self.backlog_size -= len(response.body) + else: + self.backlog_size -= self.FAILURE_SIZE + + def is_idle(self): + return not (self.queue or self.processing) + + def needs_backout(self): + return self.backlog_size > self.max_backlog_size + class Scraper(object): def __init__(self, engine): @@ -24,7 +65,7 @@ class Scraper(object): """Open the given domain for scraping and allocate resources for it""" if domain in self.sites: raise RuntimeError('Scraper domain already opened: %s' % domain) - self.sites[domain] = set() + self.sites[domain] = SiteInfo() def close_domain(self, domain): """Close a domain being scraped and release its resources""" @@ -36,27 +77,41 @@ class Scraper(object): """Return True if there isn't any more spiders to process""" return not self.sites - def scrape(self, response, request, spider): + def enqueue_scrape(self, response, request, spider): + site = self.sites[spider.domain_name] + dfd = site.add_response_request(response, request) + def finish_scraping(_): + site.finish_response(response) + return _ + dfd.addBoth(finish_scraping) + dfd.addErrback(log.err, 'Scraper bug processing %s' % request, \ + domain=spider.domain_name) + self.scrape_next(spider) + return dfd + + def scrape_next(self, spider): + site = self.sites.get(spider.domain_name) + if not site: + return + + # Process responses in queue + while site.queue: + response, request, deferred = site.next_response_request_deferred() + self._scrape(response, request, spider).chainDeferred(deferred) + + def _scrape(self, response, request, spider): """Handle the downloaded response or failure trough the spider callback/errback""" assert isinstance(response, (Response, Failure)) - domain = spider.domain_name - self.sites[domain].add(response) - def _finish_scraping(_): - self.sites[domain].remove(response) - return _ - - dfd = self._scrape(response, request, spider) # returns spiders processed output + 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.addBoth(_finish_scraping) - dfd.addErrback(log.err, 'Scraper bug processing %s' % request, \ - domain=domain) return dfd - def _scrape(self, request_result, request, spider): - """Handle the diferent cases of request's result been a Response or a Failure""" + 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, \ request_result, request, spider) @@ -74,13 +129,15 @@ class Scraper(object): referer = request.headers.get('Referer', None) msg = "SPIDER BUG processing <%s> from <%s>: %s" % (request.url, referer, _failure) log.msg(msg, log.ERROR, domain=spider.domain_name) - stats.incpath("%s/spider_exceptions/%s" % (spider.domain_name, _failure.value.__class__.__name__)) + 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 process_spider_output(self, output, request, spider): + # TODO: keep closing state internally instead of checking engine if spider.domain_name in self.engine.closing: return elif isinstance(output, Request):