diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 464668391..002bd7a1f 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -6,7 +6,7 @@ For more information see docs/topics/architecture.rst """ from time import time -from twisted.internet import reactor, defer +from twisted.internet import defer from twisted.python.failure import Failure from scrapy import log, signals @@ -18,14 +18,17 @@ from scrapy.http import Response, Request from scrapy.utils.misc import load_object from scrapy.utils.signal import send_catch_log, send_catch_log_deferred from scrapy.utils.defer import mustbe_deferred +from scrapy.utils.reactor import CallLaterOnce + class Slot(object): - def __init__(self, start_requests, close_if_idle): + def __init__(self, start_requests, close_if_idle, nextcall): self.closing = False self.inprogress = set() # requests in progress self.start_requests = iter(start_requests) self.close_if_idle = close_if_idle + self.nextcall = nextcall def add_request(self, request): self.inprogress.add(request) @@ -41,6 +44,8 @@ class Slot(object): def _maybe_fire_closing(self): if self.closing and not self.inprogress: + if self.nextcall: + self.nextcall.cancel() self.closing.callback(None) @@ -51,7 +56,6 @@ class ExecutionEngine(object): self.slots = {} self.running = False self.paused = False - self._next_request_calls = {} self.scheduler = load_object(settings['SCHEDULER'])() self.downloader = Downloader() self.scraper = Scraper(self, self.settings) @@ -84,31 +88,19 @@ class ExecutionEngine(object): return self.scheduler.is_idle() and self.downloader.is_idle() and \ self.scraper.is_idle() - def next_request(self, spider, now=False): - """Scrape the next request for the spider passed. - - The next request to be scraped is retrieved from the scheduler and - requested from the downloader. - - The spider is closed if there are no more pages to scrape. - """ - if now: - self._next_request_calls.pop(spider, None) - elif spider not in self._next_request_calls: - call = reactor.callLater(0, self.next_request, spider, now=True) - self._next_request_calls[spider] = call - return call - else: + def _next_request(self, spider): + try: + slot = self.slots[spider] + except KeyError: return if self.paused: - return reactor.callLater(5, self.next_request, spider) + return slot.nextcall.schedule(5) while not self._needs_backout(spider): - if not self._next_request(spider): + if not self._next_request_from_scheduler(spider): break - slot = self.slots[spider] if slot.start_requests and not self._needs_backout(spider): try: request = slot.start_requests.next() @@ -127,7 +119,7 @@ class ExecutionEngine(object): or self.downloader.slots[spider].needs_backout() \ or self.scraper.slots[spider].needs_backout() - def _next_request(self, spider): + def _next_request_from_scheduler(self, spider): request = self.scheduler.next_request(spider) if not request: return @@ -137,7 +129,8 @@ class ExecutionEngine(object): slot = self.slots[spider] d.addBoth(lambda _: slot.remove_request(request)) d.addErrback(log.msg, spider=spider) - d.addBoth(lambda _: self.next_request(spider)) + d.addBoth(lambda _: slot.nextcall.schedule()) + d.addErrback(log.msg, spider=spider) return d def _handle_downloader_output(self, response, request, spider): @@ -147,13 +140,8 @@ class ExecutionEngine(object): self.crawl(response, spider) return # response is a Response or Failure - d = defer.Deferred() - d.addBoth(self.scraper.enqueue_scrape, request, spider) + d = self.scraper.enqueue_scrape(response, request, spider) d.addErrback(log.err, spider=spider) - if isinstance(response, Failure): - d.errback(response) - else: - d.callback(response) return d def spider_is_idle(self, spider): @@ -171,7 +159,7 @@ class ExecutionEngine(object): @property def open_spiders(self): - return self.downloader.slots.keys() + return self.slots.keys() def has_capacity(self): """Does the engine have capacity to handle more spiders""" @@ -181,7 +169,7 @@ class ExecutionEngine(object): assert spider in self.open_spiders, \ "Spider %r not opened when crawling: %s" % (spider.name, request) self.schedule(request, spider) - self.next_request(spider) + self.slots[spider].nextcall.schedule() def schedule(self, request, spider): return self.scheduler.enqueue_request(spider, request) @@ -202,7 +190,6 @@ class ExecutionEngine(object): slot = self.slots[spider] slot.add_request(request) def _on_success(response): - """handle the result of a page download""" assert isinstance(response, (Response, Request)) if isinstance(response, Response): response.request = request # tie request to response received @@ -213,7 +200,7 @@ class ExecutionEngine(object): return response def _on_complete(_): - self.next_request(spider) + slot.nextcall.schedule() return _ dwld = mustbe_deferred(self.downloader.fetch, request, spider) @@ -226,13 +213,15 @@ class ExecutionEngine(object): assert self.has_capacity(), "No free spider slots when opening %r" % \ spider.name log.msg("Spider opened", spider=spider) - self.slots[spider] = Slot(start_requests or (), close_if_idle) + nextcall = CallLaterOnce(self._next_request, spider) + slot = Slot(start_requests or (), close_if_idle, nextcall) + self.slots[spider] = slot yield self.scheduler.open_spider(spider) self.downloader.open_spider(spider) yield self.scraper.open_spider(spider) stats.open_spider(spider) yield send_catch_log_deferred(signals.spider_opened, spider=spider) - self.next_request(spider) + slot.nextcall.schedule() def _spider_idle(self, spider): """Called when a spider gets idle. This function is called when there @@ -246,8 +235,7 @@ class ExecutionEngine(object): spider=spider, dont_log=DontCloseSpider) if any(isinstance(x, Failure) and isinstance(x.value, DontCloseSpider) \ for _, x in res): - reactor.callLater(5, self.next_request, spider) - return + return self.slots[spider].nextcall.schedule(5) if self.spider_is_idle(spider): self.close_spider(spider, reason='finished') @@ -273,9 +261,6 @@ class ExecutionEngine(object): dfd.addBoth(lambda _: self.scheduler.close_spider(spider)) dfd.addErrback(log.err, spider=spider) - dfd.addBoth(lambda _: self._cancel_next_call(spider)) - dfd.addErrback(log.err, spider=spider) - dfd.addBoth(lambda _: send_catch_log_deferred(signal=signals.spider_closed, \ spider=spider, reason=reason)) dfd.addErrback(log.err, spider=spider) @@ -292,11 +277,6 @@ class ExecutionEngine(object): return dfd - def _cancel_next_call(self, spider): - call = self._next_request_calls.pop(spider, None) - if call and call.active: - call.cancel() - def _close_all_spiders(self): dfds = [self.close_spider(s, reason='shutdown') for s in self.open_spiders] dlist = defer.DeferredList(dfds) diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 521766f96..67187bd85 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -95,9 +95,7 @@ class Scraper(object): slot.closing.callback(spider) def enqueue_scrape(self, response, request, spider): - slot = self.slots.get(spider, None) - if slot is None: - return + slot = self.slots[spider] dfd = slot.add_response_request(response, request) def finish_scraping(_): slot.finish_response(response) diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py index 90cddbefe..43c98cc09 100644 --- a/scrapy/utils/engine.py +++ b/scrapy/utils/engine.py @@ -23,6 +23,7 @@ def get_engine_status(engine=None): spider_tests = [ "engine.spider_is_idle(spider)", "engine.slots[spider].closing", + "len(engine.slots[spider].inprogress)", "len(engine.scheduler.pending_requests[spider])", "len(engine.downloader.slots[spider].queue)", "len(engine.downloader.slots[spider].active)", @@ -41,7 +42,7 @@ def get_engine_status(engine=None): status['global'] += [(test, eval(test))] except Exception, e: status['global'] += [(test, "%s (exception)" % type(e).__name__)] - for spider in set(engine.downloader.slots.keys() + engine.scraper.slots.keys()): + for spider in engine.slots.keys(): x = [] for test in spider_tests: try: diff --git a/scrapy/utils/reactor.py b/scrapy/utils/reactor.py index 88bb15eb4..a99063a61 100644 --- a/scrapy/utils/reactor.py +++ b/scrapy/utils/reactor.py @@ -15,3 +15,27 @@ def listen_tcp(portrange, host, factory): except error.CannotListenError: if x == portrange[1]: raise + + +class CallLaterOnce(object): + """Schedule a function to be called in the next reactor loop, but only if + it hasn't been already scheduled since the last time it run. + """ + + def __init__(self, func, *a, **kw): + self._func = func + self._a = a + self._kw = kw + self._call = None + + def schedule(self, delay=0): + if self._call is None: + self._call = reactor.callLater(delay, self) + + def cancel(self): + if self._call: + self._call.cancel() + + def __call__(self): + self._call = None + return self._func(*self._a, **self._kw)