diff --git a/scrapy/contrib/throttle.py b/scrapy/contrib/throttle.py index 228af2021..a1afd2c5f 100644 --- a/scrapy/contrib/throttle.py +++ b/scrapy/contrib/throttle.py @@ -1,11 +1,8 @@ -# TODO: this extension is currently broken and needs to be ported after the -# downloader refactoring introduced in r2732 - from scrapy.xlib.pydispatch import dispatcher -from scrapy.utils.python import setattr_default -from scrapy.conf import settings from scrapy.exceptions import NotConfigured from scrapy import signals +from scrapy.utils.httpobj import urlparse_cached +from scrapy.resolver import dnscache class AutoThrottle(object): """ @@ -61,56 +58,74 @@ class AutoThrottle(object): """ - # TODO: convert these to settings - START_DELAY = 5.0 - MAX_CONCURRENCY = 8 - CONCURRENCY_CHECK_PERIOD = 10 - - DEBUG = 0 - - def __init__(self): + def __init__(self, crawler): + settings = crawler.settings if not settings.getbool('AUTOTHROTTLE_ENABLED'): raise NotConfigured + self.crawler = crawler dispatcher.connect(self.spider_opened, signal=signals.spider_opened) - dispatcher.connect(self.spider_closed, signal=signals.spider_closed) dispatcher.connect(self.response_received, signal=signals.response_received) - self.last_latencies = {} - self.last_lat = {} + self.START_DELAY = settings.getfloat("AUTOTHROTTLE_START_DELAY", 5.0) + self.CONCURRENCY_CHECK_PERIOD = settings.getint("AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD", 10) + self.MAX_CONCURRENCY = settings.getint("AUTOTHROTTLE_MAX_CONCURRENCY", 8) + self.DEBUG = settings.getint("AUTOTHROTTLE_DEBUG", False) + + @classmethod + def from_crawler(cls, crawler): + return cls(crawler) def spider_opened(self, spider): spider.download_delay = self.START_DELAY spider.max_concurrent_requests = 1 - self.last_latencies[spider] = [self.START_DELAY] - self.last_lat[spider] = self.START_DELAY, 0.0 - - def spider_closed(self, spider): - del self.last_latencies[spider] - del self.last_lat[spider] - + self.last_latencies = [self.START_DELAY] + self.last_lat = self.START_DELAY, 0.0 + def response_received(self, response, spider): + slot = self._get_slot(response.request) latency = response.meta.get('download_latency') - if not latency: + + if not latency or not slot: return - spider.download_delay = (spider.download_delay + latency) / 2.0 - self._check_concurrency(spider, latency) - if self.DEBUG: - print "conc:%2d | delay:%5d ms | latency:%5d ms | size:%6d bytes" % \ - (spider.max_concurrent_requests, spider.download_delay*1000, \ - latency*1000, len(response.body)) - def _check_concurrency(self, spider, latency): - latencies = self.last_latencies[spider] + self._adjust_delay(slot, latency, response) + self._check_concurrency(slot, latency) + + if self.DEBUG: + spider.log("conc:%2d | delay:%5d ms | latency:%5d ms | size:%6d bytes" % \ + (slot.concurrency, slot.delay*1000, \ + latency*1000, len(response.body))) + + def _get_slot(self, request): + downloader = self.crawler.engine.downloader + key = urlparse_cached(request).hostname or '' + if downloader.ip_concurrency: + key = dnscache.get(key, key) + return downloader.slots.get(key) + + def _check_concurrency(self, slot, latency): + latencies = self.last_latencies latencies.append(latency) if len(latencies) == self.CONCURRENCY_CHECK_PERIOD: curavg, curdev = avg_stdev(latencies) - preavg, predev = self.last_lat[spider] - self.last_lat[spider] = curavg, curdev + preavg, predev = self.last_lat + self.last_lat = curavg, curdev del latencies[:] if curavg > preavg + predev: - if spider.max_concurrent_requests > 1: - spider.max_concurrent_requests -= 1 - elif spider.max_concurrent_requests < self.MAX_CONCURRENCY: - spider.max_concurrent_requests += 1 + if slot.concurrency > 1: + slot.concurrency -= 1 + elif slot.concurrency < self.MAX_CONCURRENCY: + slot.concurrency += 1 + + def _adjust_delay(self, slot, latency, response): + """Define delay adjustment policy""" + # if latency is bigger than old delay, then use latency instead of mean. Works better with problematic sites + new_delay = (slot.delay + latency) / 2.0 if latency < slot.delay else latency + + # dont adjust delay if response status != 200 and new delay is smaller than old one, + # as error pages (and redirections) are usually small and so tend to reduce latency, thus provoking a positive feedback + # by reducing delay instead of increase. + if response.status == 200 or new_delay > slot.delay: + slot.delay = new_delay def avg_stdev(lst): """Return average and standard deviation of the given list"""