Some enhancements to Scrapy core:

- added graceful shutdown (with one ^C) and forced shutdown (two ^C)
- added optional loglevel to IgnoreRequest exception
- calling Spider Manager close_domain() function from engine to avoid race
  conditions caused by signal dispatch order
- moved twisted reactor controlling logic outside the engine and into the
  scrapy manager
- made log.err function respect the current log level filter
This commit is contained in:
Pablo Hoffman 2009-09-03 08:27:48 -03:00
parent 51eb641fc7
commit a7f5c6f878
8 changed files with 196 additions and 139 deletions

View File

@ -129,7 +129,7 @@ def execute(argv=None):
del args[0] # remove command name from args
_save_command_executed(cmdname, cmd, args, opts)
from scrapy.core.manager import scrapymanager
scrapymanager.configure()
scrapymanager.configure(control_reactor=True)
ret = _run_command(cmd, args, opts)
if ret is False:
parser.print_help()

View File

@ -9,8 +9,6 @@ import urlparse
from twisted.plugin import getCache
from twisted.python.rebuild import rebuild
from scrapy.xlib.pydispatch import dispatcher
from scrapy.core import signals
from scrapy.spider.models import ISpider
from scrapy import log
from scrapy.conf import settings
@ -25,10 +23,9 @@ class TwistedPluginSpiderManager(object):
self.default_domain = None
self.force_domain = None
self.spider_modules = None
dispatcher.connect(self._domain_closed, signal=signals.domain_closed)
def fromdomain(self, domain_name):
return self.asdict().get(domain_name)
def fromdomain(self, domain):
return self.asdict().get(domain)
def fromurl(self, url):
if self.force_domain:
@ -112,10 +109,11 @@ class TwistedPluginSpiderManager(object):
sys.stderr.write("Interrupted while loading Scrapy spiders\n")
sys.exit(2)
def _domain_closed(self, domain, spider):
def close_domain(self, domain):
"""Reload spider module to release any resources held on to by the
spider
"""
spider = self._spiders[domain]
module_name = spider.__module__
module = sys.modules[module_name]
if hasattr(module, 'SPIDER'):

View File

@ -5,6 +5,7 @@ Download web pages using asynchronous IO
from time import time
from twisted.internet import reactor, defer
from twisted.python.failure import Failure
from scrapy.core.exceptions import IgnoreRequest
from scrapy.spider import spiders
@ -31,8 +32,8 @@ class SiteInfo(object):
else:
self.max_concurrent_requests = max_concurrent_requests
self.queue = []
self.active = set()
self.queue = []
self.transferring = set()
self.closing = False
self.lastseen = 0
@ -56,7 +57,7 @@ class Downloader(object):
self.concurrent_domains = settings.getint('CONCURRENT_DOMAINS')
def fetch(self, request, spider):
""" Main method to use to request a download
"""Main method to use to request a download
This method includes middleware mangling. Middleware can returns a
Response object, then request never reach downloader queue, and it will
@ -64,7 +65,7 @@ class Downloader(object):
"""
site = self.sites[spider.domain_name]
if site.closing:
raise IgnoreRequest('Can\'t fetch on a closing domain')
raise IgnoreRequest('Cannot fetch on a closing domain')
site.active.add(request)
def _deactivate(_):
@ -72,25 +73,21 @@ class Downloader(object):
self._close_if_idle(spider.domain_name)
return _
return self.middleware.download(self.enqueue, request, spider).addBoth(_deactivate)
dfd = self.middleware.download(self.enqueue, request, spider)
return dfd.addBoth(_deactivate)
def enqueue(self, request, spider):
"""Enqueue a Request for a effective download from site"""
deferred = defer.Deferred()
site = self.sites[spider.domain_name]
if site.closing:
raise IgnoreRequest
deferred = defer.Deferred()
site.queue.append((request, deferred))
self.process_queue(spider)
self._process_queue(spider)
return deferred
def process_queue(self, spider):
try:
self._process_queue(spider)
except:
log.exc('Downloader process queue bug')
def _process_queue(self, spider):
""" Effective download requests from site queue
"""
"""Effective download requests from site queue"""
domain = spider.domain_name
site = self.sites.get(domain)
if not site:
@ -101,14 +98,18 @@ 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)
reactor.callLater(penalty, self._process_queue, spider=spider)
return
site.lastseen = now
# Process requests in queue if there are free slots to transfer for this site
# Process enqueued requests if there are free slots to transfer for this site
while site.queue and site.free_transfer_slots() > 0:
request, deferred = site.queue.pop(0)
self._download(site, request, spider).chainDeferred(deferred)
if site.closing:
dfd = defer.fail(Failure(IgnoreRequest()))
else:
dfd = self._download(site, request, spider)
dfd.chainDeferred(deferred)
self._close_if_idle(domain)
@ -130,7 +131,13 @@ class Downloader(object):
site.transferring.add(request)
def finish_transferring(_):
site.transferring.remove(request)
self.process_queue(spider)
self._process_queue(spider)
# avoid partially downloaded responses from propagating to the
# downloader middleware, to speed-up the closing process
if site.closing:
log.msg("Crawled while closing domain: %s" % request, \
level=log.DEBUG)
raise IgnoreRequest
return _
return dfd.addBoth(finish_transferring)
@ -153,7 +160,7 @@ class Downloader(object):
site.closing = True
spider = spiders.fromdomain(domain)
self.process_queue(spider)
self._process_queue(spider)
def has_capacity(self):
"""Does the downloader have capacity to handle more domains"""

View File

@ -6,7 +6,7 @@ For more information see docs/topics/architecture.rst
"""
from datetime import datetime
from twisted.internet import reactor, task
from twisted.internet import reactor, task, defer
from twisted.python.failure import Failure
from scrapy.xlib.pydispatch import dispatcher
@ -30,8 +30,8 @@ class ExecutionEngine(object):
self.keep_alive = False
self.closing = {} # dict (domain -> reason) of spiders being closed
self.running = False
self.killed = False
self.paused = False
self.control_reactor = True
self._next_request_pending = set()
self._mainloop_task = task.LoopingCall(self._mainloop)
@ -45,34 +45,34 @@ class ExecutionEngine(object):
self.scraper = Scraper(self)
self.configured = True
def start(self, control_reactor=True):
def start(self):
"""Start the execution engine"""
if self.running:
return
self.control_reactor = control_reactor
self.start_time = datetime.utcnow()
send_catch_log(signal=signals.engine_started, sender=self.__class__)
self._mainloop_task.start(5.0, now=True)
reactor.callWhenRunning(self._mainloop)
self.running = True
if control_reactor:
reactor.run() # blocking call
def stop(self):
"""Stop the execution engine"""
"""Stop the execution engine gracefully"""
if not self.running:
return
self.running = False
for domain in self.open_domains:
spider = spiders.fromdomain(domain)
send_catch_log(signal=signals.domain_closed, sender=self.__class__, \
domain=domain, spider=spider, reason='shutdown')
stats.close_domain(domain, reason='shutdown')
reactor.addSystemEventTrigger('before', 'shutdown', \
self.close_domain, domain, reason='shutdown')
if self._mainloop_task.running:
self._mainloop_task.stop()
if self.control_reactor and reactor.running:
reactor.stop()
send_catch_log(signal=signals.engine_stopped, sender=self.__class__)
def kill(self):
"""Forces shutdown without waiting for pending transfers to finish.
stop() must have been called first
"""
if self.running:
return
self.killed = True
def pause(self):
"""Pause the execution engine"""
@ -83,7 +83,8 @@ class ExecutionEngine(object):
self.paused = False
def is_idle(self):
return self.scheduler.is_idle() and self.downloader.is_idle() and self.scraper.is_idle()
return self.scheduler.is_idle() and self.downloader.is_idle() and \
self.scraper.is_idle()
def next_domain(self):
domain = self.domain_scheduler.next_domain()
@ -130,12 +131,15 @@ class ExecutionEngine(object):
spider = spiders.fromdomain(domain)
dwld = mustbe_deferred(self.download, request, spider)
dwld.chainDeferred(deferred).addBoth(lambda _: deferred)
return dwld.addErrback(log.err)
dwld.addErrback(log.err, "Unhandled error on engine._next_request")
return dwld
def domain_is_idle(self, domain):
scraper_idle = domain in self.scraper.sites and self.scraper.sites[domain].is_idle()
scraper_idle = domain in self.scraper.sites \
and self.scraper.sites[domain].is_idle()
pending = self.scheduler.domain_has_pending_requests(domain)
downloading = domain in self.downloader.sites and self.downloader.sites[domain].active
downloading = domain in self.downloader.sites \
and self.downloader.sites[domain].active
return scraper_idle and not (pending or downloading)
def domain_is_closed(self, domain):
@ -155,7 +159,7 @@ class ExecutionEngine(object):
def crawl(self, request, spider):
schd = mustbe_deferred(self.schedule, request, spider)
schd.addBoth(self.scraper.enqueue_scrape, request, spider)
schd.addErrback(log.err)
schd.addErrback(log.err, "Unhandled error on engine.crawl()")
schd.addBoth(lambda _: self.next_request(spider.domain_name))
def schedule(self, request, spider):
@ -190,8 +194,8 @@ class ExecutionEngine(object):
assert isinstance(response, (Response, Request))
if isinstance(response, Response):
response.request = request # tie request to response received
log.msg("Crawled %s (referer: <%s>)" % (response, referer), level=log.DEBUG, \
domain=domain)
log.msg("Crawled %s (referer: <%s>)" % (response, referer), \
level=log.DEBUG, domain=domain)
return response
elif isinstance(response, Request):
newrequest = response
@ -201,12 +205,17 @@ class ExecutionEngine(object):
def _on_error(_failure):
"""handle an error processing a page"""
ex = _failure.value
errmsg = str(_failure) if not isinstance(ex, IgnoreRequest) \
else _failure.getErrorMessage()
log.msg("Downloading <%s> (referer: <%s>): %s" % (request.url, referer, errmsg), \
log.ERROR, domain=domain)
return Failure(IgnoreRequest(str(ex)))
exc = _failure.value
if isinstance(exc, IgnoreRequest):
errmsg = _failure.getErrorMessage()
level = exc.level
else:
errmsg = str(_failure)
level = log.ERROR
if errmsg:
log.msg("Downloading <%s> (referer: <%s>): %s" % (request.url, \
referer, errmsg), level=level, domain=domain)
return Failure(IgnoreRequest(str(exc)))
def _on_complete(_):
self.next_request(domain)
@ -249,7 +258,7 @@ class ExecutionEngine(object):
self.next_request(domain)
return
except:
log.err(_why="Exception catched on domain_idle signal dispatch")
log.err("Exception catched on domain_idle signal dispatch")
if self.domain_is_idle(domain):
self.close_domain(domain, reason='finished')
@ -265,14 +274,19 @@ class ExecutionEngine(object):
self.closing[domain] = reason
self.downloader.close_domain(domain)
self.scheduler.clear_pending_requests(domain)
self._finish_closing_domain_if_idle(domain)
return self._finish_closing_domain_if_idle(domain)
return defer.succeed(None)
def _finish_closing_domain_if_idle(self, domain):
"""Call _finish_closing_domain if domain is idle"""
if self.domain_is_idle(domain):
if self.domain_is_idle(domain) or self.killed:
self._finish_closing_domain(domain)
else:
reactor.callLater(5, self._finish_closing_domain_if_idle, domain)
dfd = defer.Deferred()
dfd.addCallback(self._finish_closing_domain_if_idle)
delay = 5 if self.running else 1
reactor.callLater(delay, dfd.callback, domain)
return dfd
def _finish_closing_domain(self, domain):
"""This function is called after the domain has been closed"""
@ -284,6 +298,10 @@ class ExecutionEngine(object):
domain=domain, spider=spider, reason=reason)
stats.close_domain(domain, reason=reason)
log.msg("Domain closed (%s)" % reason, domain=domain)
self._mainloop()
spiders.close_domain(domain)
if self.running:
self._mainloop()
elif not self.open_domains:
send_catch_log(signal=signals.engine_stopped, sender=self.__class__)
scrapyengine = ExecutionEngine()

View File

@ -5,6 +5,8 @@ These exceptions are documented in docs/topics/exceptions.rst. Please don't add
new exceptions here without documenting them there.
"""
from scrapy import log
# Internal
class NotConfigured(Exception):
@ -15,7 +17,13 @@ class NotConfigured(Exception):
class IgnoreRequest(Exception):
"""Indicates a decision was made not to process a request"""
pass
def __init__(self, msg='', level=log.ERROR):
self.msg = msg
self.level = level
def __str__(self):
return self.msg
class DontCloseDomain(Exception):
"""Request the domain not to be closed yet"""

View File

@ -9,7 +9,48 @@ from scrapy.core.engine import scrapyengine
from scrapy.spider import spiders
from scrapy.utils.misc import arg_to_iter
from scrapy.utils.url import is_url
from scrapy.conf import settings
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
def _parse_args(args):
"""Parse crawl arguments and return a dict of domains -> list of requests"""
requests, urls, sites = set(), set(), set()
for a in args:
if isinstance(a, Request):
requests.add(a)
elif is_url(a):
urls.add(a)
else:
sites.add(a)
perdomain = {}
# sites
for domain in sites:
spider = spiders.fromdomain(domain)
if not spider:
log.msg('Could not find spider for %s' % domain, log.ERROR)
continue
reqs = spider.start_requests()
perdomain.setdefault(domain, []).extend(reqs)
# urls
for url in urls:
spider = spiders.fromurl(url)
if spider:
for req in arg_to_iter(spider.make_requests_from_url(url)):
perdomain.setdefault(spider.domain_name, []).append(req)
else:
log.msg('Could not find spider for <%s>' % url, log.ERROR)
# requests
for request in requests:
spider = spiders.fromurl(request.url)
if not spider:
log.msg('Could not find spider for %s' % request, log.ERROR)
continue
perdomain.setdefault(spider.domain_name, []).append(request)
return perdomain
class ExecutionManager(object):
"""Process a list of sites or urls.
@ -27,8 +68,8 @@ class ExecutionManager(object):
def configure(self, control_reactor=True):
self.control_reactor = control_reactor
if control_reactor:
self._install_signals()
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
install_shutdown_handlers(self._signal_shutdown)
reactor.addSystemEventTrigger('before', 'shutdown', scrapyengine.stop)
if not log.started:
log.start()
@ -45,7 +86,7 @@ class ExecutionManager(object):
def crawl(self, *args):
"""Schedule the given args for crawling. args is a list of urls or domains"""
requests = self._parse_args(args)
requests = _parse_args(args)
# schedule initial requests to be scraped at engine start
for domain in requests or ():
spider = spiders.fromdomain(domain)
@ -54,88 +95,40 @@ class ExecutionManager(object):
def runonce(self, *args):
"""Run the engine until it finishes scraping all domains and then exit"""
if not self.configured:
self.configure()
assert self.configured, "Scrapy Manger not yet configured"
self.crawl(*args)
scrapyengine.start()
if self.control_reactor:
reactor.run(installSignalHandlers=False)
def start(self, control_reactor=True):
def start(self):
"""Start the scrapy server, without scheduling any domains"""
self.configure(control_reactor)
assert self.configured, "Scrapy Manger not yet configured"
scrapyengine.keep_alive = True
scrapyengine.start(control_reactor=control_reactor)
if control_reactor:
self.stop()
scrapyengine.start()
if self.control_reactor:
reactor.run(installSignalHandlers=False)
def stop(self):
"""Stop the scrapy server, shutting down the execution engine"""
self.interrupted = True
scrapyengine.stop()
log.log_level = -999 # disable logging
if self.control_reactor:
signal.signal(signal.SIGTERM, signal.SIG_IGN)
signal.signal(signal.SIGINT, signal.SIG_IGN)
if hasattr(signal, "SIGBREAK"):
signal.signal(signal.SIGBREAK, signal.SIG_IGN)
if self.control_reactor and reactor.running:
reactor.stop()
def reload_spiders(self):
"""Reload all enabled spiders except for the ones that are currently
running.
"""
spiders.reload(skip_domains=scrapyengine.open_domains)
def _signal_shutdown(self, signum, _):
signame = signal_names[signum]
log.msg("Received %s, shutting down gracefully. Send again to force " \
"unclean shutdown" % signame, level=log.INFO)
reactor.callFromThread(self.stop)
install_shutdown_handlers(self._signal_kill)
def _install_signals(self):
def sig_handler_terminate(signalinfo, param):
log.msg('Received shutdown request, waiting for deferreds to finish...', log.INFO)
self.stop()
signal.signal(signal.SIGTERM, sig_handler_terminate)
# only handle SIGINT if there isn't already a handler (e.g. for Pdb)
if signal.getsignal(signal.SIGINT) == signal.default_int_handler:
signal.signal(signal.SIGINT, sig_handler_terminate)
# Catch Ctrl-Break in windows
if hasattr(signal, "SIGBREAK"):
signal.signal(signal.SIGBREAK, sig_handler_terminate)
def _parse_args(self, args):
"""Parse crawl arguments and return a dict of domains -> list of
requests"""
requests, urls, sites = set(), set(), set()
for a in args:
if isinstance(a, Request):
requests.add(a)
elif is_url(a):
urls.add(a)
else:
sites.add(a)
perdomain = {}
# sites
for domain in sites:
spider = spiders.fromdomain(domain)
if not spider:
log.msg('Could not find spider for %s' % domain, log.ERROR)
continue
reqs = spider.start_requests()
perdomain.setdefault(domain, []).extend(reqs)
# urls
for url in urls:
spider = spiders.fromurl(url)
if spider:
for req in arg_to_iter(spider.make_requests_from_url(url)):
perdomain.setdefault(spider.domain_name, []).append(req)
else:
log.msg('Could not find spider for <%s>' % url, log.ERROR)
# requests
for request in requests:
spider = spiders.fromurl(request.url)
if not spider:
log.msg('Could not find spider for %s' % request, log.ERROR)
continue
perdomain.setdefault(spider.domain_name, []).append(request)
return perdomain
def _signal_kill(self, signum, _):
signame = signal_names[signum]
log.msg('Received %s twice, forcing unclean shutdown' % signame, \
level=log.INFO)
log.log_level = log.SILENT # disable logging of confusing tracebacks
reactor.callFromThread(scrapyengine.kill)
install_shutdown_handlers(signal.SIG_IGN)
scrapymanager = ExecutionManager()

View File

@ -56,19 +56,24 @@ def start(logfile=None, loglevel=None, logstdout=None):
def msg(message, level=INFO, component=BOT_NAME, domain=None):
"""Log message according to the level"""
if level > log_level:
return
dispatcher.send(signal=logmessage_received, message=message, level=level, \
domain=domain)
system = domain if domain else component
if level <= log_level:
msg_txt = unicode_to_str("%s: %s" % (level_names[level], message))
log.msg(msg_txt, system=system)
msg_txt = unicode_to_str("%s: %s" % (level_names[level], message))
log.msg(msg_txt, system=system)
def exc(message, level=ERROR, component=BOT_NAME, domain=None):
message = message + '\n' + format_exc()
msg(message, level, component, domain)
def err(*args, **kwargs):
def err(_stuff=None, _why=None, **kwargs):
if ERROR > log_level:
return
domain = kwargs.pop('domain', None)
component = kwargs.pop('component', BOT_NAME)
kwargs['system'] = domain if domain else component
log.err(*args, **kwargs)
if _why:
_why = unicode_to_str("ERROR: %s" % _why)
log.err(_stuff, _why, **kwargs)

28
scrapy/utils/ossignal.py Normal file
View File

@ -0,0 +1,28 @@
from __future__ import absolute_import
from twisted.internet import reactor
import signal
signal_names = {}
for signame in dir(signal):
if signame.startswith("SIG"):
signum = getattr(signal, signame)
if isinstance(signum, int):
signal_names[signum] = signame
def install_shutdown_handlers(function, override_sigint=True):
"""Install the given function as a signal handler for all common shutdown
signals (such as SIGINT, SIGTERM, etc). If override_sigint is ``False`` the
SIGINT handler won't be install if there is already a handler in place
(e.g. Pdb)
"""
reactor._handleSignals()
signal.signal(signal.SIGTERM, function)
if signal.getsignal(signal.SIGINT) == signal.default_int_handler or \
override_sigint:
signal.signal(signal.SIGINT, function)
# Catch Ctrl-Break in windows
if hasattr(signal, "SIGBREAK"):
signal.signal(signal.SIGBREAK, function)