From d7daf836d5091bf2190dc4711339aed3af5f83a0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 6 Dec 2012 10:42:57 -0200 Subject: [PATCH] Altering delay is enough to auto throttle downloads --- docs/topics/autothrottle.rst | 20 ++--- scrapy/contrib/throttle.py | 127 ++++++++++------------------- scrapy/core/downloader/__init__.py | 20 +++-- 3 files changed, 65 insertions(+), 102 deletions(-) diff --git a/docs/topics/autothrottle.rst b/docs/topics/autothrottle.rst index 086c9574d..8b25e6c3f 100644 --- a/docs/topics/autothrottle.rst +++ b/docs/topics/autothrottle.rst @@ -37,12 +37,6 @@ This adjusts download delays and concurrency based on the following rules: :setting:`AUTOTHROTTLE_START_DELAY` 2. when a response is received, the download delay is adjusted to the average of previous download delay and the latency of the response. -3. after :setting:`AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD` responses have - passed, the average latency of this period is checked against the previous - one and: - - * if the latency remained constant (within standard deviation limits), it is increased - * if the latency has increased (beyond standard deviation limits) and the concurrency is higher than 1, the concurrency is decreased .. note:: The AutoThrottle extension honours the standard Scrapy settings for concurrency and delay. This means that it will never set a download delay @@ -55,11 +49,11 @@ The settings used to control the AutoThrottle extension are: * :setting:`AUTOTHROTTLE_ENABLED` * :setting:`AUTOTHROTTLE_START_DELAY` -* :setting:`AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD` +* :setting:`AUTOTHROTTLE_MAX_DELAY` * :setting:`AUTOTHROTTLE_DEBUG` -* :setting:`DOWNLOAD_DELAY` * :setting:`CONCURRENT_REQUESTS_PER_DOMAIN` * :setting:`CONCURRENT_REQUESTS_PER_IP` +* :setting:`DOWNLOAD_DELAY` For more information see :ref:`autothrottle-algorithm`. @@ -81,14 +75,14 @@ Default: ``5.0`` The initial download delay (in seconds). -.. setting:: AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD +.. setting:: AUTOTHROTTLE_MAX_DELAY -AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD -~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +AUTOTHROTTLE_MAX_DELAY +~~~~~~~~~~~~~~~~~~~~~~ -Default: ``10`` +Default: ``60.0`` -How many responses should pass to perform concurrency adjustments. +The maximum download delay (in seconds) to be set in case of high latencies. .. setting:: AUTOTHROTTLE_DEBUG diff --git a/scrapy/contrib/throttle.py b/scrapy/contrib/throttle.py index d4d7762f1..13c07a509 100644 --- a/scrapy/contrib/throttle.py +++ b/scrapy/contrib/throttle.py @@ -1,106 +1,69 @@ +import logging from scrapy.exceptions import NotConfigured from scrapy import signals -from scrapy.utils.httpobj import urlparse_cached -from scrapy.resolver import dnscache + class AutoThrottle(object): def __init__(self, crawler): - settings = crawler.settings - if not settings.getbool('AUTOTHROTTLE_ENABLED'): - raise NotConfigured self.crawler = crawler - crawler.signals.connect(self.spider_opened, signal=signals.spider_opened) - crawler.signals.connect(self.response_received, signal=signals.response_received) - self.START_DELAY = settings.getfloat("AUTOTHROTTLE_START_DELAY", 5.0) - self.CONCURRENCY_CHECK_PERIOD = settings.getint("AUTOTHROTTLE_CONCURRENCY_CHECK_PERIOD", 10) - self.MAX_CONCURRENCY = self._max_concurency(settings) - self.MIN_DOWNLOAD_DELAY = self._min_download_delay(settings) - self.DEBUG = settings.getbool("AUTOTHROTTLE_DEBUG") - self.last_latencies = [self.START_DELAY] - self.last_lat = self.START_DELAY, 0.0 + if not crawler.settings.getbool('AUTOTHROTTLE_ENABLED'): + raise NotConfigured - def _min_download_delay(self, settings): - return max(settings.getfloat("AUTOTHROTTLE_MIN_DOWNLOAD_DELAY"), - settings.getfloat("DOWNLOAD_DELAY")) - - def _max_concurency(self, settings): - delay = self._min_download_delay(settings) - if delay == 0: - candidates = ["AUTOTHROTTLE_MAX_CONCURRENCY", - "CONCURRENT_REQUESTS_PER_DOMAIN", "CONCURRENT_REQUESTS_PER_IP"] - candidates = [settings.getint(x) for x in candidates] - candidates = [x for x in candidates if x > 0] - if candidates: - return min(candidates) - return 1 + self.debug = crawler.settings.getbool("AUTOTHROTTLE_DEBUG") + crawler.signals.connect(self._spider_opened, signal=signals.spider_opened) + crawler.signals.connect(self._response_downloaded, signal=signals.response_downloaded) @classmethod def from_crawler(cls, crawler): return cls(crawler) - def spider_opened(self, spider): - if hasattr(spider, "download_delay"): - self.MIN_DOWNLOAD_DELAY = spider.download_delay - spider.download_delay = self.START_DELAY - if hasattr(spider, "max_concurrent_requests"): - self.MAX_CONCURRENCY = spider.max_concurrent_requests - # override in order to avoid to initialize slot with concurrency > 1 - spider.max_concurrent_requests = 1 + def _spider_opened(self, spider): + self.mindelay = self._min_delay(spider) + self.maxdelay = self._max_delay(spider) + spider.download_delay = self._start_delay(spider) - def response_received(self, response, spider): - key, slot = self._get_slot(response.request) - latency = response.meta.get('download_latency') - - if not latency or not slot: + def _min_delay(self, spider): + s = self.crawler.settings + return getattr(spider, 'download_delay', 0.0) or \ + s.getfloat('AUTOTHROTTLE_MIN_DOWNLOAD_DELAY') or \ + s.getfloat('DOWNLOAD_DELAY') + + def _max_delay(self, spider): + return self.crawler.settings.getfloat('AUTOTHROTTLE_MAX_DELAY', 60.0) + + def _start_delay(self, spider): + return max(self.mindelay, self.crawler.settings.getfloat('AUTOTHROTTLE_START_DELAY', 5.0)) + + def _response_downloaded(self, response, request, spider): + key, slot = self._get_slot(request, spider) + latency = request.meta.get('download_latency') + if latency is None or slot is None: return + olddelay = slot.delay self._adjust_delay(slot, latency, response) - self._check_concurrency(slot, latency) + if self.debug: + diff = slot.delay - olddelay + size = len(response.body) + conc = len(slot.transferring) + msg = "slot: %s | conc:%2d | delay:%5d ms (%+d) | latency:%5d ms | size:%6d bytes" % \ + (key, conc, slot.delay * 1000, diff * 1000, latency * 1000, size) + spider.log(msg, level=logging.INFO) - if self.DEBUG: - spider.log("slot: %s | conc:%2d | delay:%5d ms | latency:%5d ms | size:%6d bytes" % \ - (key, 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 key, downloader.slots.get(key) or downloader.inactive_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 - self.last_lat = curavg, curdev - del latencies[:] - if curavg > preavg + predev: - if slot.concurrency > 1: - slot.concurrency -= 1 - elif slot.concurrency < self.MAX_CONCURRENCY: - slot.concurrency += 1 + def _get_slot(self, request, spider): + key = request.meta.get('download_slotkey') + return key, self.crawler.engine.downloader.slots.get(key) 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 + # If latency is bigger than old delay, then use latency instead of mean. + # It works better with problematic sites + new_delay = min(max(self.mindelay, latency, (slot.delay + latency) / 2.0), self.maxdelay) - if new_delay < self.MIN_DOWNLOAD_DELAY: - new_delay = self.MIN_DOWNLOAD_DELAY - - # 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. + # 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""" - avg = sum(lst)/len(lst) - sdsq = sum((x-avg) ** 2 for x in lst) - stdev = (sdsq / (len(lst) -1)) ** 0.5 - return avg, stdev diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index a9808e7c8..25573fa6e 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -76,6 +76,7 @@ class Downloader(object): def fetch(self, request, spider): key, slot = self._get_slot(request, spider) + request.meta['download_slotkey'] = key self.active.add(request) slot.active.add(request) @@ -112,12 +113,7 @@ class Downloader(object): return key, self.slots[key] def _enqueue_request(self, request, spider, slot): - def _downloaded(response): - self.signals.send_catch_log(signal=signals.response_downloaded, \ - response=response, request=request, spider=spider) - return response - - deferred = defer.Deferred().addCallback(_downloaded) + deferred = defer.Deferred() slot.queue.append((request, deferred)) self._process_queue(spider, slot) return deferred @@ -152,7 +148,17 @@ class Downloader(object): # 1. Create the download deferred dfd = mustbe_deferred(self.handlers.download_request, request, spider) - # 2. After response arrives, remove the request from transferring + # 2. Notify response_downloaded listeners about the recent download + # before querying queue for next request + def _downloaded(response): + self.signals.send_catch_log(signal=signals.response_downloaded, + response=response, + request=request, + spider=spider) + return response + dfd.addCallback(_downloaded) + + # 3. After response arrives, remove the request from transferring # state to free up the transferring slot so it can be used by the # following requests (perhaps those which came from the downloader # middleware itself)