more core refactoring including moving engine next request call logic to a separate class

This commit is contained in:
Pablo Hoffman 2011-07-25 10:46:00 -03:00
parent 209ecdf471
commit ea3bf6d95d
4 changed files with 52 additions and 49 deletions

View File

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

View File

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

View File

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

View File

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