mirror of https://github.com/scrapy/scrapy.git
Added ExecutionQueue class for feeding spiders and requests to scrape. This
class can (and is meant to) be subclassed by projects that want to use a custom mechanism for feeding spiders to crawl. For example, a queue that pulls spiders to scrape from Amazon SQS (an example will be added soon). Also introduced a rather big core refactoring of Scrapy manager and Scrapy engine.
This commit is contained in:
parent
8c1feb7ae4
commit
cae22930c8
|
|
@ -1,5 +1,6 @@
|
|||
from scrapy import log
|
||||
from scrapy.command import ScrapyCommand
|
||||
from scrapy.core.queue import ExecutionQueue
|
||||
from scrapy.core.manager import scrapymanager
|
||||
from scrapy.conf import settings
|
||||
from scrapy.http import Request
|
||||
|
|
@ -31,23 +32,25 @@ class Command(ScrapyCommand):
|
|||
settings.overrides['CRAWLSPIDER_FOLLOW_LINKS'] = False
|
||||
|
||||
def run(self, args, opts):
|
||||
q = ExecutionQueue()
|
||||
urls, names = self._split_urls_and_names(args)
|
||||
for name in names:
|
||||
scrapymanager.crawl_spider_name(name)
|
||||
q.append_spider_name(name)
|
||||
|
||||
if opts.spider:
|
||||
try:
|
||||
spider = spiders.create(opts.spider)
|
||||
for url in urls:
|
||||
scrapymanager.crawl_url(url, spider)
|
||||
q.append_url(url, spider)
|
||||
except KeyError:
|
||||
log.msg('Could not find spider: %s' % opts.spider, log.ERROR)
|
||||
log.msg('Unable to find spider: %s' % opts.spider, log.ERROR)
|
||||
else:
|
||||
for name, urls in self._group_urls_by_spider(urls):
|
||||
spider = spiders.create(name)
|
||||
for url in urls:
|
||||
scrapymanager.crawl_url(url, spider)
|
||||
q.append_url(url, spider)
|
||||
|
||||
scrapymanager.queue = q
|
||||
scrapymanager.start()
|
||||
|
||||
def _group_urls_by_spider(self, urls):
|
||||
|
|
|
|||
|
|
@ -28,28 +28,27 @@ class Command(ScrapyCommand):
|
|||
parser.add_option("--headers", dest="headers", action="store_true", \
|
||||
help="print response HTTP headers instead of body")
|
||||
|
||||
def _print_response(self, response, opts):
|
||||
if opts.headers:
|
||||
pprint.pprint(response.headers)
|
||||
else:
|
||||
print response.body
|
||||
|
||||
def run(self, args, opts):
|
||||
if len(args) != 1 or not is_url(args[0]):
|
||||
return False
|
||||
responses = [] # to collect downloaded responses
|
||||
request = Request(args[0], callback=responses.append, dont_filter=True)
|
||||
cb = lambda x: self._print_response(x, opts)
|
||||
request = Request(args[0], callback=cb, dont_filter=True)
|
||||
|
||||
spider = None
|
||||
if opts.spider:
|
||||
try:
|
||||
spider = spiders.create(opts.spider)
|
||||
except KeyError:
|
||||
log.msg("Could not find spider: %s" % opts.spider, log.ERROR)
|
||||
else:
|
||||
spider = scrapymanager._create_spider_for_request(request, \
|
||||
BaseSpider('default'))
|
||||
|
||||
scrapymanager.crawl_request(request, spider)
|
||||
scrapymanager.configure()
|
||||
scrapymanager.queue.append_request(request, spider, \
|
||||
default_spider=BaseSpider('default'))
|
||||
scrapymanager.start()
|
||||
|
||||
# display response
|
||||
if responses:
|
||||
if opts.headers:
|
||||
pprint.pprint(responses[0].headers)
|
||||
else:
|
||||
print responses[0].body
|
||||
|
||||
|
|
|
|||
|
|
@ -8,8 +8,6 @@ from scrapy.utils.spider import iterate_spider_output
|
|||
from scrapy.utils.url import is_url
|
||||
from scrapy import log
|
||||
|
||||
from collections import defaultdict
|
||||
|
||||
class Command(ScrapyCommand):
|
||||
|
||||
requires_project = True
|
||||
|
|
@ -75,25 +73,20 @@ class Command(ScrapyCommand):
|
|||
if not len(args) == 1 or not is_url(args[0]):
|
||||
return False
|
||||
|
||||
request = Request(args[0])
|
||||
responses = [] # to collect downloaded responses
|
||||
request = Request(args[0], callback=responses.append)
|
||||
|
||||
if opts.spider:
|
||||
try:
|
||||
spider = spiders.create(opts.spider)
|
||||
except KeyError:
|
||||
log.msg('Could not find spider: %s' % opts.spider, log.ERROR)
|
||||
log.msg('Unable to find spider: %s' % opts.spider, log.ERROR)
|
||||
return
|
||||
else:
|
||||
spider = scrapymanager._create_spider_for_request(request, \
|
||||
log_none=True, log_multiple=True)
|
||||
spider = spiders.create_for_request(request)
|
||||
|
||||
if not spider:
|
||||
return
|
||||
|
||||
responses = [] # to collect downloaded responses
|
||||
request = request.replace(callback=responses.append)
|
||||
|
||||
scrapymanager.crawl_request(request, spider)
|
||||
scrapymanager.configure()
|
||||
scrapymanager.queue.append_request(request, spider)
|
||||
scrapymanager.start()
|
||||
|
||||
if not responses:
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ class Command(ScrapyCommand):
|
|||
module = _import_file(args[0])
|
||||
|
||||
# schedule spider and start engine
|
||||
scrapymanager.crawl_spider(module.SPIDER)
|
||||
scrapymanager.queue.append_spider(module.SPIDER)
|
||||
scrapymanager.start()
|
||||
|
||||
if opts.output:
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
from scrapy.core.queue import KeepAliveExecutionQueue
|
||||
from scrapy.command import ScrapyCommand
|
||||
from scrapy.core.manager import scrapymanager
|
||||
|
||||
|
|
@ -9,4 +10,6 @@ class Command(ScrapyCommand):
|
|||
return "Start the Scrapy manager but don't run any spider (idle mode)"
|
||||
|
||||
def run(self, args, opts):
|
||||
scrapymanager.start(keep_alive=True)
|
||||
q = KeepAliveExecutionQueue()
|
||||
scrapymanager.queue = q
|
||||
scrapymanager.start()
|
||||
|
|
|
|||
|
|
@ -34,6 +34,27 @@ class TwistedPluginSpiderManager(object):
|
|||
return [name for name, spider in self._spiders.iteritems()
|
||||
if url_is_from_spider(request.url, spider)]
|
||||
|
||||
def create_for_request(self, request, default_spider=None, \
|
||||
log_none=False, log_multiple=False, **spider_kwargs):
|
||||
"""Create a spider to handle the given Request.
|
||||
|
||||
This will look for the spiders that can handle the given request (using
|
||||
find_by_request) and return a (new) Spider if (and only if) there is
|
||||
only one Spider able to handle the Request.
|
||||
|
||||
If multiple spiders (or no spider) are found, it will return the
|
||||
default_spider passed. It can optionally log if multiple or no spiders
|
||||
are found.
|
||||
"""
|
||||
snames = self.find_by_request(request)
|
||||
if len(snames) == 1:
|
||||
return self.create(snames[0], **spider_kwargs)
|
||||
if len(snames) > 1 and log_multiple:
|
||||
log.msg('More than one spider found for: %s' % request, log.ERROR)
|
||||
if len(snames) == 0 and log_none:
|
||||
log.msg('Unable to find spider for: %s' % request, log.ERROR)
|
||||
return default_spider
|
||||
|
||||
def list(self):
|
||||
"""Returns list of spiders available."""
|
||||
return self._spiders.keys()
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ from scrapy.core import signals
|
|||
from scrapy.core.manager import scrapymanager
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.management.web import banner
|
||||
from scrapy.management.web import banner, webconsole_discover_module
|
||||
from scrapy.conf import settings
|
||||
|
||||
class Spiderctl(object):
|
||||
|
|
@ -20,8 +20,6 @@ class Spiderctl(object):
|
|||
self.finished = set()
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
from scrapy.management.web import webconsole_discover_module
|
||||
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
|
||||
|
||||
def spider_opened(self, spider):
|
||||
|
|
@ -35,7 +33,7 @@ class Spiderctl(object):
|
|||
if wc_request.args:
|
||||
changes = self.webconsole_control(wc_request)
|
||||
|
||||
self.scheduled = [s.name for s in scrapyengine.spider_scheduler._pending_spiders]
|
||||
self.scheduled = [s[0].name for s in scrapymanager.queue.spider_requests]
|
||||
self.idle = [d for d in self.enabled_spiders if d not in self.scheduled
|
||||
and d not in self.running
|
||||
and d not in self.finished]
|
||||
|
|
@ -126,8 +124,8 @@ class Spiderctl(object):
|
|||
if "remove_pending_spiders" in args:
|
||||
removed = []
|
||||
for name in args["remove_pending_spiders"]:
|
||||
if scrapyengine.spider_scheduler.remove_pending_spider(name):
|
||||
removed.append(name)
|
||||
q = scrapymanager.queue
|
||||
q.spider_requests = [x for x in q.spider_requests if x[0].name != name]
|
||||
if removed:
|
||||
s += "<p>"
|
||||
s += "Removed scheduled spiders: <ul><li>%s</li></ul>" % "</li><li>".join(args["remove_pending_spiders"])
|
||||
|
|
@ -135,14 +133,14 @@ class Spiderctl(object):
|
|||
if "add_pending_spiders" in args:
|
||||
for name in args["add_pending_spiders"]:
|
||||
if name not in scrapyengine.scheduler.pending_requests:
|
||||
scrapymanager.crawl_spider_name(name)
|
||||
scrapymanager.queue.append_spider_name(name)
|
||||
s += "<p>"
|
||||
s += "Scheduled spiders: <ul><li>%s</li></ul>" % "</li><li>".join(args["add_pending_spiders"])
|
||||
s += "</p>"
|
||||
if "rerun_finished_spiders" in args:
|
||||
for name in args["rerun_finished_spiders"]:
|
||||
if name not in scrapyengine.scheduler.pending_requests:
|
||||
scrapymanager.crawl_spider_name(name)
|
||||
scrapymanager.queue.append_spider_name(name)
|
||||
self.finished.remove(name)
|
||||
s += "<p>"
|
||||
s += "Re-scheduled finished spiders: <ul><li>%s</li></ul>" % "</li><li>".join(args["rerun_finished_spiders"])
|
||||
|
|
|
|||
|
|
@ -183,10 +183,6 @@ class Downloader(object):
|
|||
site.cancel_request_calls()
|
||||
self._process_queue(spider)
|
||||
|
||||
def has_capacity(self):
|
||||
"""Does the downloader have capacity to handle more spiders"""
|
||||
return len(self.sites) < self.concurrent_spiders
|
||||
|
||||
def is_idle(self):
|
||||
return not self.sites
|
||||
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@ For more information see docs/topics/architecture.rst
|
|||
"""
|
||||
from time import time
|
||||
|
||||
from twisted.internet import reactor, task, defer
|
||||
from twisted.internet import reactor, defer
|
||||
from twisted.python.failure import Failure
|
||||
from scrapy.xlib.pydispatch import dispatcher
|
||||
|
||||
|
|
@ -27,57 +27,42 @@ class ExecutionEngine(object):
|
|||
|
||||
def __init__(self):
|
||||
self.configured = False
|
||||
self.keep_alive = False
|
||||
self.closing = {} # dict (spider -> reason) of spiders being closed
|
||||
self.running = False
|
||||
self.killed = False
|
||||
self.paused = False
|
||||
self._next_request_calls = {}
|
||||
self._mainloop_task = task.LoopingCall(self._mainloop)
|
||||
self._crawled_logline = load_object(settings['LOG_FORMATTER_CRAWLED'])
|
||||
|
||||
def configure(self):
|
||||
def configure(self, spider_closed_callback):
|
||||
"""
|
||||
Configure execution engine with the given scheduling policy and downloader.
|
||||
"""
|
||||
self.scheduler = load_object(settings['SCHEDULER'])()
|
||||
self.spider_scheduler = load_object(settings['SPIDER_SCHEDULER'])()
|
||||
self.downloader = Downloader()
|
||||
self.scraper = Scraper(self)
|
||||
self.configured = True
|
||||
self._spider_closed_callback = spider_closed_callback
|
||||
|
||||
def start(self):
|
||||
"""Start the execution engine"""
|
||||
if self.running:
|
||||
return
|
||||
assert not self.running, "Engine already running"
|
||||
self.start_time = time()
|
||||
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
|
||||
|
||||
def stop(self):
|
||||
"""Stop the execution engine gracefully"""
|
||||
if not self.running:
|
||||
return
|
||||
assert self.running, "Engine not running"
|
||||
self.running = False
|
||||
def before_shutdown():
|
||||
dfd = self._close_all_spiders()
|
||||
return dfd.addBoth(lambda _: self._finish_stopping_engine())
|
||||
reactor.addSystemEventTrigger('before', 'shutdown', before_shutdown)
|
||||
if self._mainloop_task.running:
|
||||
self._mainloop_task.stop()
|
||||
try:
|
||||
reactor.stop()
|
||||
except RuntimeError: # raised if already stopped or in shutdown stage
|
||||
pass
|
||||
dfd = self._close_all_spiders()
|
||||
return dfd.addBoth(lambda _: self._finish_stopping_engine())
|
||||
|
||||
def kill(self):
|
||||
"""Forces shutdown without waiting for pending transfers to finish.
|
||||
stop() must have been called first
|
||||
"""
|
||||
if self.running:
|
||||
return
|
||||
assert not self.running, "Call engine.stop() before engine.kill()"
|
||||
self.killed = True
|
||||
|
||||
def pause(self):
|
||||
|
|
@ -92,12 +77,6 @@ class ExecutionEngine(object):
|
|||
return self.scheduler.is_idle() and self.downloader.is_idle() and \
|
||||
self.scraper.is_idle()
|
||||
|
||||
def next_spider(self):
|
||||
spider = self.spider_scheduler.next_spider()
|
||||
if spider:
|
||||
self.open_spider(spider)
|
||||
return True
|
||||
|
||||
def next_request(self, spider, now=False):
|
||||
"""Scrape the next request for the spider passed.
|
||||
|
||||
|
|
@ -163,7 +142,13 @@ class ExecutionEngine(object):
|
|||
def open_spiders(self):
|
||||
return self.downloader.sites.keys()
|
||||
|
||||
def has_capacity(self):
|
||||
"""Does the engine have capacity to handle more spiders"""
|
||||
return len(self.downloader.sites) < self.downloader.concurrent_spiders
|
||||
|
||||
def crawl(self, request, spider):
|
||||
assert spider in self.open_spiders, \
|
||||
"Spider %r not opened when crawling: %s" % (spider.name, request)
|
||||
if not request.deferred.callbacks:
|
||||
log.msg("Unable to crawl Request with no callback: %s" % request,
|
||||
level=log.ERROR, spider=spider)
|
||||
|
|
@ -180,25 +165,9 @@ class ExecutionEngine(object):
|
|||
def schedule(self, request, spider):
|
||||
if spider in self.closing:
|
||||
raise IgnoreRequest()
|
||||
if not self.scheduler.spider_is_open(spider):
|
||||
self.scheduler.open_spider(spider)
|
||||
if self.spider_is_closed(spider): # scheduler auto-open
|
||||
self.spider_scheduler.add_spider(spider)
|
||||
self.next_request(spider)
|
||||
return self.scheduler.enqueue_request(spider, request)
|
||||
|
||||
def _mainloop(self):
|
||||
"""Add more spiders to be scraped if the downloader has the capacity.
|
||||
|
||||
If there is nothing else scheduled then stop the execution engine.
|
||||
"""
|
||||
if not self.running or self.paused:
|
||||
return
|
||||
|
||||
while self.running and self.downloader.has_capacity():
|
||||
if not self.next_spider():
|
||||
return self._stop_if_idle()
|
||||
|
||||
def download(self, request, spider):
|
||||
def _on_success(response):
|
||||
"""handle the result of a page download"""
|
||||
|
|
@ -241,14 +210,15 @@ class ExecutionEngine(object):
|
|||
return dwld
|
||||
|
||||
def open_spider(self, spider):
|
||||
assert self.has_capacity(), "No free spider slots when opening %r" % \
|
||||
spider.name
|
||||
log.msg("Spider opened", spider=spider)
|
||||
self.next_request(spider)
|
||||
|
||||
self.scheduler.open_spider(spider)
|
||||
self.downloader.open_spider(spider)
|
||||
self.scraper.open_spider(spider)
|
||||
stats.open_spider(spider)
|
||||
|
||||
send_catch_log(signals.spider_opened, sender=self.__class__, spider=spider)
|
||||
self.next_request(spider)
|
||||
|
||||
def _spider_idle(self, spider):
|
||||
"""Called when a spider gets idle. This function is called when there
|
||||
|
|
@ -270,11 +240,6 @@ class ExecutionEngine(object):
|
|||
if self.spider_is_idle(spider):
|
||||
self.close_spider(spider, reason='finished')
|
||||
|
||||
def _stop_if_idle(self):
|
||||
"""Call the stop method if the system has no outstanding tasks. """
|
||||
if self.is_idle() and not self.keep_alive:
|
||||
self.stop()
|
||||
|
||||
def close_spider(self, spider, reason='cancelled'):
|
||||
"""Close (cancel) spider and clear all its outstanding requests"""
|
||||
if spider not in self.closing:
|
||||
|
|
@ -316,7 +281,7 @@ class ExecutionEngine(object):
|
|||
dfd.addErrback(log.err, "Unhandled error on SpiderManager.close_spider()",
|
||||
spider=spider)
|
||||
dfd.addBoth(lambda _: log.msg("Spider closed (%s)" % reason, spider=spider))
|
||||
reactor.callLater(0, self._mainloop)
|
||||
reactor.callLater(0, self._spider_closed_callback)
|
||||
return dfd
|
||||
|
||||
def _finish_stopping_engine(self):
|
||||
|
|
|
|||
|
|
@ -1,28 +1,26 @@
|
|||
import signal
|
||||
|
||||
from twisted.internet import reactor
|
||||
from twisted.internet import reactor, defer
|
||||
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.core.queue import ExecutionQueue
|
||||
from scrapy.extension import extensions
|
||||
from scrapy import log
|
||||
from scrapy.http import Request
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.utils.misc import arg_to_iter
|
||||
from scrapy.utils.ossignal import install_shutdown_handlers, signal_names
|
||||
|
||||
|
||||
class ExecutionManager(object):
|
||||
|
||||
def __init__(self):
|
||||
self.interrupted = False
|
||||
self.configured = False
|
||||
self.control_reactor = True
|
||||
self.engine = scrapyengine
|
||||
|
||||
def configure(self, control_reactor=True):
|
||||
def configure(self, control_reactor=True, queue=None):
|
||||
self.control_reactor = control_reactor
|
||||
if control_reactor:
|
||||
install_shutdown_handlers(self._signal_shutdown)
|
||||
reactor.addSystemEventTrigger('before', 'shutdown', scrapyengine.stop)
|
||||
|
||||
if not log.started:
|
||||
log.start()
|
||||
|
|
@ -33,68 +31,54 @@ class ExecutionManager(object):
|
|||
log.msg("Enabled extensions: %s" % ", ".join(extensions.enabled.iterkeys()),
|
||||
level=log.DEBUG)
|
||||
|
||||
scrapyengine.configure()
|
||||
self.queue = queue or ExecutionQueue()
|
||||
self.engine.configure(self._spider_closed)
|
||||
self.configured = True
|
||||
|
||||
def crawl_url(self, url, spider=None):
|
||||
"""Schedule given url for crawling."""
|
||||
if spider is None:
|
||||
spider = self._create_spider_for_request(Request(url), log_none=True, \
|
||||
log_multiple=True)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def _start_next_spider(self):
|
||||
spider, requests = yield defer.maybeDeferred(self.queue.get_next)
|
||||
if spider:
|
||||
requests = arg_to_iter(spider.make_requests_from_url(url))
|
||||
self._crawl_requests(requests, spider)
|
||||
self._start_spider(spider, requests)
|
||||
if self.engine.has_capacity() and not self._nextcall.active():
|
||||
self._nextcall = reactor.callLater(self.queue.polling_delay, \
|
||||
self._start_next_spider)
|
||||
|
||||
def crawl_request(self, request, spider=None):
|
||||
"""Schedule request for crawling."""
|
||||
assert self.configured, "Scrapy Manager not yet configured"
|
||||
if spider is None:
|
||||
spider = self._create_spider_for_request(request, log_none=True, \
|
||||
log_multiple=True)
|
||||
if spider:
|
||||
scrapyengine.crawl(request, spider)
|
||||
@defer.inlineCallbacks
|
||||
def _start_spider(self, spider, requests):
|
||||
"""Don't call this method. Use self.queue to start new spiders"""
|
||||
yield defer.maybeDeferred(self.engine.open_spider, spider)
|
||||
for request in requests:
|
||||
self.engine.crawl(request, spider)
|
||||
|
||||
def crawl_spider_name(self, name):
|
||||
"""Schedule given spider by name for crawling."""
|
||||
try:
|
||||
spider = spiders.create(name)
|
||||
except KeyError:
|
||||
log.msg('Could not find spider: %s' % name, log.ERROR)
|
||||
else:
|
||||
self.crawl_spider(spider)
|
||||
@defer.inlineCallbacks
|
||||
def _spider_closed(self):
|
||||
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()
|
||||
|
||||
def crawl_spider(self, spider):
|
||||
"""Schedule spider for crawling."""
|
||||
requests = spider.start_requests()
|
||||
self._crawl_requests(requests, spider)
|
||||
|
||||
def _crawl_requests(self, requests, spider):
|
||||
"""Shortcut to schedule a list of requests"""
|
||||
for req in requests:
|
||||
self.crawl_request(req, spider)
|
||||
|
||||
def start(self, keep_alive=False):
|
||||
"""Start the scrapy server, without scheduling any domains"""
|
||||
scrapyengine.keep_alive = keep_alive
|
||||
scrapyengine.start()
|
||||
@defer.inlineCallbacks
|
||||
def start(self):
|
||||
yield defer.maybeDeferred(self.engine.start)
|
||||
self._nextcall = reactor.callLater(0, self._spider_closed)
|
||||
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
|
||||
if self.control_reactor:
|
||||
reactor.run(installSignalHandlers=False)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def stop(self):
|
||||
"""Stop the scrapy server, shutting down the execution engine"""
|
||||
self.interrupted = True
|
||||
scrapyengine.stop()
|
||||
|
||||
def _create_spider_for_request(self, request, default=None, log_none=False, \
|
||||
log_multiple=False):
|
||||
spider_names = spiders.find_by_request(request)
|
||||
if len(spider_names) == 1:
|
||||
return spiders.create(spider_names[0])
|
||||
if len(spider_names) > 1 and log_multiple:
|
||||
log.msg('More than one spider found for: %s' % request, log.ERROR)
|
||||
if len(spider_names) == 0 and log_none:
|
||||
log.msg('Could not find spider for: %s' % request, log.ERROR)
|
||||
return default
|
||||
if self._nextcall.active():
|
||||
self._nextcall.cancel()
|
||||
if self.engine.running:
|
||||
yield defer.maybeDeferred(self.engine.stop)
|
||||
try:
|
||||
reactor.stop()
|
||||
except RuntimeError: # raised if already stopped or in shutdown stage
|
||||
pass
|
||||
|
||||
def _signal_shutdown(self, signum, _):
|
||||
signame = signal_names[signum]
|
||||
|
|
@ -108,7 +92,7 @@ class ExecutionManager(object):
|
|||
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)
|
||||
reactor.callFromThread(self.engine.kill)
|
||||
install_shutdown_handlers(signal.SIG_IGN)
|
||||
|
||||
scrapymanager = ExecutionManager()
|
||||
|
|
|
|||
|
|
@ -0,0 +1,87 @@
|
|||
from twisted.internet import defer
|
||||
|
||||
from scrapy.http import Request
|
||||
from scrapy.utils.misc import arg_to_iter
|
||||
from scrapy import log
|
||||
from scrapy.spider import spiders
|
||||
|
||||
|
||||
class ExecutionQueue(object):
|
||||
|
||||
polling_delay = 5
|
||||
|
||||
def __init__(self, _spiders=spiders):
|
||||
self.spider_requests = []
|
||||
self._spiders = _spiders
|
||||
|
||||
def _append_next(self):
|
||||
"""Called when there are no more itemsl left in self.spider_requests.
|
||||
This method is meant to be overriden in subclasses to add new (spider,
|
||||
requests) tuples to self.spider_requests. It can return a Deferred.
|
||||
"""
|
||||
pass
|
||||
|
||||
def get_next(self):
|
||||
"""Return a tuple (spider, requests) containing a list of Requests and
|
||||
the Spider which will be used to crawl those Requests. If there aren't
|
||||
any more spiders to crawl it must return (None, []).
|
||||
|
||||
This method can return a deferred.
|
||||
"""
|
||||
if self.spider_requests:
|
||||
return self._get_next_now()
|
||||
d = defer.maybeDeferred(self._append_next)
|
||||
d.addCallback(lambda _: self._get_next_now())
|
||||
return d
|
||||
|
||||
def _get_next_now(self):
|
||||
try:
|
||||
return self.spider_requests.pop(0)
|
||||
except IndexError:
|
||||
return (None, [])
|
||||
|
||||
def is_finished(self):
|
||||
"""Return True if the queue is empty and there won't be any more
|
||||
spiders to crawl (this is for one-shot runs). If it returns ``False``
|
||||
Scrapy will keep polling this queue for new requests to scrape
|
||||
"""
|
||||
return not bool(self.spider_requests)
|
||||
|
||||
def append_spider(self, spider):
|
||||
"""Append a Spider to crawl"""
|
||||
requests = spider.start_requests()
|
||||
self.spider_requests.append((spider, requests))
|
||||
|
||||
def append_request(self, request, spider=None, **kwargs):
|
||||
if spider is None:
|
||||
spider = self._spiders.create_for_request(request, **kwargs)
|
||||
if spider:
|
||||
self.spider_requests.append((spider, [request]))
|
||||
|
||||
def append_url(self, url, spider=None, **kwargs):
|
||||
"""Append a URL to crawl with the given spider. If the spider is not
|
||||
given, a spider will be looked up based on the URL
|
||||
"""
|
||||
if spider is None:
|
||||
spider = self._spiders.create_for_request(Request(url), **kwargs)
|
||||
if spider:
|
||||
requests = arg_to_iter(spider.make_requests_from_url(url))
|
||||
self.spider_requests.append((spider, requests))
|
||||
|
||||
def append_spider_name(self, spider_name, **spider_kwargs):
|
||||
"""Append a spider to crawl given its name and optional arguments,
|
||||
which are used to instantiate it. The SpiderManager is used to lookup
|
||||
the spider
|
||||
"""
|
||||
try:
|
||||
spider = self._spiders.create(spider_name, **spider_kwargs)
|
||||
except KeyError:
|
||||
log.msg('Unable to find spider: %s' % spider_name, log.ERROR)
|
||||
else:
|
||||
self.append_spider(spider)
|
||||
|
||||
|
||||
class KeepAliveExecutionQueue(ExecutionQueue):
|
||||
|
||||
def is_finished(self):
|
||||
return False
|
||||
|
|
@ -18,6 +18,7 @@ from scrapy.utils.response import open_in_browser
|
|||
from scrapy.conf import settings
|
||||
from scrapy.core.manager import scrapymanager
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.core.queue import KeepAliveExecutionQueue
|
||||
from scrapy.http import Request, TextResponse
|
||||
|
||||
def relevant_var(varname):
|
||||
|
|
@ -54,10 +55,11 @@ class Shell(object):
|
|||
url = parse_url(request_or_url)
|
||||
request = Request(url)
|
||||
|
||||
spider = scrapymanager._create_spider_for_request(request, \
|
||||
BaseSpider('default'), log_multiple=True)
|
||||
spider = spiders.create_for_request(request, BaseSpider('default'), \
|
||||
log_multiple=True)
|
||||
|
||||
print "Fetching %s..." % request
|
||||
scrapyengine.open_spider(spider)
|
||||
response = threads.blockingCallFromThread(reactor, scrapyengine.schedule, \
|
||||
request, spider)
|
||||
if response:
|
||||
|
|
@ -108,7 +110,8 @@ class Shell(object):
|
|||
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
||||
|
||||
reactor.callInThread(self._console_thread, url)
|
||||
scrapymanager.start(keep_alive=True)
|
||||
scrapymanager.queue = KeepAliveExecutionQueue()
|
||||
scrapymanager.start()
|
||||
|
||||
def inspect_response(self, response):
|
||||
print
|
||||
|
|
|
|||
|
|
@ -0,0 +1,105 @@
|
|||
import unittest
|
||||
|
||||
from scrapy.core.queue import ExecutionQueue, KeepAliveExecutionQueue
|
||||
from scrapy.spider import BaseSpider
|
||||
from scrapy.http import Request
|
||||
|
||||
class TestSpider(BaseSpider):
|
||||
|
||||
name = "default"
|
||||
|
||||
def start_requests(self):
|
||||
return [Request("http://www.example.com/1"), \
|
||||
Request("http://www.example.com/2")]
|
||||
|
||||
def make_requests_from_url(self, url):
|
||||
return [Request(url + "/make1"), Request(url + "/make2")]
|
||||
|
||||
|
||||
class TestSpiderManager(object):
|
||||
|
||||
def create_for_request(self, request, **kwargs):
|
||||
return TestSpider('create_for_request', **kwargs)
|
||||
|
||||
def create(self, spider_name, **spider_kwargs):
|
||||
return TestSpider(spider_name, **spider_kwargs)
|
||||
|
||||
|
||||
class ExecutionQueueTest(unittest.TestCase):
|
||||
|
||||
queue_class = ExecutionQueue
|
||||
|
||||
def setUp(self):
|
||||
self.queue = self.queue_class(_spiders=TestSpiderManager())
|
||||
self.spider = TestSpider()
|
||||
self.request = Request('about:none')
|
||||
|
||||
def tearDown(self):
|
||||
del self.queue, self.spider, self.request
|
||||
|
||||
def test_is_finished(self):
|
||||
self.assert_(self.queue.is_finished())
|
||||
self.queue.append_request(self.request, self.spider)
|
||||
self.assert_(not self.queue.is_finished())
|
||||
|
||||
def test_append_spider(self):
|
||||
spider = TestSpider()
|
||||
self.queue.append_spider(spider)
|
||||
self.assert_(self.queue.spider_requests[0][0] is spider)
|
||||
self._assert_request_urls(self.queue.spider_requests[0][1],
|
||||
["http://www.example.com/1", "http://www.example.com/2"])
|
||||
|
||||
def test_append_request1(self):
|
||||
spider = TestSpider()
|
||||
request = Request('about:blank')
|
||||
self.queue.append_request(request, spider=spider)
|
||||
self.assert_(self.queue.spider_requests[0][0] is spider)
|
||||
self.assert_(self.queue.spider_requests[0][1][0] is request)
|
||||
|
||||
def test_append_request2(self):
|
||||
request = Request('about:blank')
|
||||
self.queue.append_request(request, arg='123')
|
||||
spider = self.queue.spider_requests[0][0]
|
||||
self.assert_(spider.name == 'create_for_request')
|
||||
self.assert_(spider.arg == '123')
|
||||
|
||||
def test_append_url(self):
|
||||
spider = TestSpider()
|
||||
url = 'http://www.example.com/asd'
|
||||
self.queue.append_url(url, spider=spider)
|
||||
self.assert_(self.queue.spider_requests[0][0] is spider)
|
||||
self._assert_request_urls(self.queue.spider_requests[0][1], \
|
||||
['http://www.example.com/asd/make1', 'http://www.example.com/asd/make2'])
|
||||
|
||||
def test_append_url2(self):
|
||||
url = 'http://www.example.com/asd'
|
||||
self.queue.append_url(url, arg='123')
|
||||
self._assert_request_urls(self.queue.spider_requests[0][1], \
|
||||
['http://www.example.com/asd/make1', 'http://www.example.com/asd/make2'])
|
||||
spider = self.queue.spider_requests[0][0]
|
||||
self.assert_(spider.name == 'create_for_request')
|
||||
self.assert_(spider.arg == '123')
|
||||
|
||||
def test_append_spider_name(self):
|
||||
self.queue.append_spider_name('test123', arg='123')
|
||||
spider = self.queue.spider_requests[0][0]
|
||||
self.assert_(spider.name == 'test123')
|
||||
self.assert_(spider.arg == '123')
|
||||
|
||||
def _assert_request_urls(self, requests, urls):
|
||||
assert all(isinstance(x, Request) for x in requests)
|
||||
self.assertEqual([x.url for x in requests], urls)
|
||||
|
||||
class KeepAliveExecutionQueueTest(ExecutionQueueTest):
|
||||
|
||||
queue_class = KeepAliveExecutionQueue
|
||||
|
||||
def test_is_finished(self):
|
||||
self.assert_(not self.queue.is_finished())
|
||||
self.queue.append_request(self.request, self.spider)
|
||||
self.assert_(not self.queue.is_finished())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
|
|
@ -97,7 +97,7 @@ class CrawlingSession(object):
|
|||
dispatcher.connect(self.response_downloaded, signals.response_downloaded)
|
||||
|
||||
scrapymanager.configure()
|
||||
scrapymanager.crawl_spider(self.spider)
|
||||
scrapymanager.queue.append_spider(self.spider)
|
||||
scrapymanager.start()
|
||||
self.port.stopListening()
|
||||
self.wasrun = True
|
||||
|
|
|
|||
|
|
@ -12,11 +12,11 @@ def get_engine_status(engine=None):
|
|||
global_tests = [
|
||||
"time()-engine.start_time",
|
||||
"engine.is_idle()",
|
||||
"engine.has_capacity()",
|
||||
"engine.scheduler.is_idle()",
|
||||
"len(engine.scheduler.pending_requests)",
|
||||
"engine.downloader.is_idle()",
|
||||
"len(engine.downloader.sites)",
|
||||
"engine.downloader.has_capacity()",
|
||||
"engine.scraper.is_idle()",
|
||||
"len(engine.scraper.sites)",
|
||||
]
|
||||
|
|
|
|||
Loading…
Reference in New Issue