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)

This commit is contained in:
Pablo Hoffman 2009-07-06 15:31:50 -03:00
parent 4f1d388733
commit 31b3d7ce1e
6 changed files with 87 additions and 26 deletions

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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