mirror of https://github.com/scrapy/scrapy.git
fix autothrottle extension for working with new downloader
This commit is contained in:
parent
509b05db57
commit
dcc50201e3
|
|
@ -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"""
|
||||
|
|
|
|||
Loading…
Reference in New Issue