mirror of https://github.com/scrapy/scrapy.git
126 lines
4.5 KiB
Python
126 lines
4.5 KiB
Python
import signal
|
|
|
|
from twisted.internet import reactor, defer
|
|
|
|
from scrapy.xlib.pydispatch import dispatcher
|
|
from scrapy.queue import ExecutionQueue
|
|
from scrapy.core.engine import ExecutionEngine
|
|
from scrapy.extension import ExtensionManager
|
|
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
|
|
from scrapy.utils.misc import load_object
|
|
from scrapy import log, signals
|
|
|
|
|
|
class Crawler(object):
|
|
|
|
def __init__(self, settings):
|
|
self.configured = False
|
|
self.settings = settings
|
|
|
|
def install(self):
|
|
import scrapy.project
|
|
assert not hasattr(scrapy.project, 'crawler'), "crawler already installed"
|
|
scrapy.project.crawler = self
|
|
|
|
def uninstall(self):
|
|
import scrapy.project
|
|
assert hasattr(scrapy.project, 'crawler'), "crawler not installed"
|
|
del scrapy.project.crawler
|
|
|
|
def configure(self):
|
|
if self.configured:
|
|
return
|
|
self.configured = True
|
|
self.extensions = ExtensionManager.from_settings(self.settings)
|
|
spman_cls = load_object(self.settings['SPIDER_MANAGER_CLASS'])
|
|
self.spiders = spman_cls.from_settings(self.settings)
|
|
spq_cls = load_object(self.settings['SPIDER_QUEUE_CLASS'])
|
|
spq = spq_cls.from_settings(self.settings)
|
|
keepalive = self.settings.getbool('KEEP_ALIVE')
|
|
pollint = self.settings.getfloat('QUEUE_POLL_INTERVAL')
|
|
self.queue = ExecutionQueue(self.spiders, spq, poll_interval=pollint,
|
|
keep_alive=keepalive)
|
|
self.engine = ExecutionEngine(self.settings, self._spider_closed)
|
|
|
|
@defer.inlineCallbacks
|
|
def _start_next_spider(self):
|
|
spider, requests = yield defer.maybeDeferred(self.queue.get_next)
|
|
if spider:
|
|
self._start_spider(spider, requests)
|
|
if self.engine.has_capacity() and not self._nextcall.active():
|
|
self._nextcall = reactor.callLater(self.queue.poll_interval, \
|
|
self._spider_closed)
|
|
|
|
@defer.inlineCallbacks
|
|
def _start_spider(self, spider, requests):
|
|
"""Don't call this method. Use self.queue to start new spiders"""
|
|
spider.set_crawler(self)
|
|
yield defer.maybeDeferred(self.engine.open_spider, spider)
|
|
for request in requests:
|
|
self.engine.crawl(request, spider)
|
|
|
|
@defer.inlineCallbacks
|
|
def _spider_closed(self, spider=None):
|
|
if not self.engine.open_spiders:
|
|
is_finished = yield defer.maybeDeferred(self.queue.is_finished)
|
|
if is_finished:
|
|
self.stop()
|
|
return
|
|
if self.engine.has_capacity():
|
|
self._start_next_spider()
|
|
|
|
@defer.inlineCallbacks
|
|
def start(self):
|
|
yield defer.maybeDeferred(self.configure)
|
|
yield defer.maybeDeferred(self.engine.start)
|
|
self._nextcall = reactor.callLater(0, self._start_next_spider)
|
|
|
|
@defer.inlineCallbacks
|
|
def stop(self):
|
|
if self._nextcall.active():
|
|
self._nextcall.cancel()
|
|
if self.engine.running:
|
|
yield defer.maybeDeferred(self.engine.stop)
|
|
|
|
|
|
class CrawlerProcess(Crawler):
|
|
"""A class to run a single Scrapy crawler in a process. It provides
|
|
automatic control of the Twisted reactor and installs some convenient
|
|
signals for shutting down the crawl.
|
|
"""
|
|
|
|
def __init__(self, *a, **kw):
|
|
super(CrawlerProcess, self).__init__(*a, **kw)
|
|
dispatcher.connect(self.stop, signals.engine_stopped)
|
|
install_shutdown_handlers(self._signal_shutdown)
|
|
|
|
def start(self):
|
|
super(CrawlerProcess, self).start()
|
|
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
|
|
reactor.run(installSignalHandlers=False) # blocking call
|
|
|
|
def stop(self):
|
|
d = super(CrawlerProcess, self).stop()
|
|
d.addBoth(self._stop_reactor)
|
|
return d
|
|
|
|
def _stop_reactor(self, _=None):
|
|
try:
|
|
reactor.stop()
|
|
except RuntimeError: # raised if already stopped or in shutdown stage
|
|
pass
|
|
|
|
def _signal_shutdown(self, signum, _):
|
|
install_shutdown_handlers(self._signal_kill)
|
|
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)
|
|
|
|
def _signal_kill(self, signum, _):
|
|
install_shutdown_handlers(signal.SIG_IGN)
|
|
signame = signal_names[signum]
|
|
log.msg('Received %s twice, forcing unclean shutdown' % signame, \
|
|
level=log.INFO)
|
|
reactor.callFromThread(self._stop_reactor)
|