Keep track of pending next-request calls in the engine and downloader, to cancel them properly when the spider is closed. This also avoids keeping spider references alive, after they are closed.

This commit is contained in:
Pablo Hoffman 2009-11-26 17:00:16 -02:00
parent c0f1c8de04
commit fd50891113
2 changed files with 22 additions and 8 deletions

View File

@ -36,6 +36,7 @@ class SiteInfo(object):
self.transferring = set()
self.closing = False
self.lastseen = 0
self.next_request_calls = set()
def free_transfer_slots(self):
return self.max_concurrent_requests - len(self.transferring)
@ -44,6 +45,10 @@ class SiteInfo(object):
# use self.active to include requests in the downloader middleware
return len(self.active) > 2 * self.max_concurrent_requests
def cancel_request_calls(self):
for call in self.next_request_calls:
call.cancel()
self.next_request_calls.clear()
class Downloader(object):
"""Mantain many concurrent downloads and provide an HTTP abstraction.
@ -97,7 +102,11 @@ class Downloader(object):
if site.download_delay:
penalty = site.download_delay - now + site.lastseen
if penalty > 0:
reactor.callLater(penalty, self._process_queue, spider=spider)
d = defer.Deferred()
d.addCallback(self._process_queue)
call = reactor.callLater(penalty, d.callback, spider)
site.next_request_calls.add(call)
d.addBoth(lambda x: site.next_request_calls.remove(call))
return
site.lastseen = now
@ -153,12 +162,13 @@ class Downloader(object):
def close_spider(self, spider):
"""Free any resources associated with the given spider"""
domain = spider.domain_name
site = self.sites.get(spider)
if not site or site.closing:
raise RuntimeError('Downloader spider already closed: %s' % domain)
raise RuntimeError('Downloader spider already closed: %s' % \
spider.domain_name)
site.closing = True
site.cancel_request_calls()
self._process_queue(spider)
def has_capacity(self):

View File

@ -32,7 +32,7 @@ class ExecutionEngine(object):
self.running = False
self.killed = False
self.paused = False
self._next_request_pending = set()
self._next_request_calls = {}
self._mainloop_task = task.LoopingCall(self._mainloop)
self._crawled_logline = load_object(settings['LOG_FORMATTER_CRAWLED'])
@ -107,10 +107,11 @@ class ExecutionEngine(object):
The spider is closed if there are no more pages to scrape.
"""
if now:
self._next_request_pending.discard(spider)
elif spider not in self._next_request_pending:
self._next_request_pending.add(spider)
return reactor.callLater(0, self.next_request, spider, now=True)
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:
return
@ -307,6 +308,9 @@ class ExecutionEngine(object):
send_catch_log(signal=signals.spider_closed, sender=self.__class__, \
spider=spider, reason=reason)
stats.close_spider(spider, reason=reason)
call = self._next_request_calls.pop(spider, None)
if call and call.active():
call.cancel()
dfd = defer.maybeDeferred(spiders.close_spider, spider)
dfd.addBoth(log.msg, "Spider closed (%s)" % reason, spider=spider)
reactor.callLater(0, self._mainloop)