Several core changes:

Execution Manager:

* added control_reactor argument to delegate external twisted
  reactor control (for example by twistd)
* now it loads spiders (if not already loaded)
* now it stars the log (if not already started)
* removed *args from configure() method
* removed **opts from runonce and start methods

Execution engine:

* added control_reactor argument to to delegate external twisted
  reactor control (for example by twistd)
* changed some functions and method names for clarity
* improve handling of exceptions in st() method
* regrouped close_domain, closed_domain, and _close_domain method
  for legibilty

Scheduler:

* replaced pending_domains_count (dict) by pending_domains (set)
* simplified some doc
This commit is contained in:
Pablo Hoffman 2009-06-15 19:44:26 -03:00
parent 3c919f2562
commit e5b99a56c4
3 changed files with 94 additions and 134 deletions

View File

@ -56,6 +56,7 @@ class ExecutionEngine(object):
self.ports = []
self.running = False
self.paused = False
self.control_reactor = True
def configure(self, scheduler=None, downloader=None):
"""
@ -114,10 +115,11 @@ class ExecutionEngine(object):
p.stopListening()
self.ports = []
def start(self):
def start(self, control_reactor=True):
"""Start the execution engine"""
if not self.running:
reactor.callWhenRunning(self._mainloop)
self.control_reactor = control_reactor
reactor.callLater(0, self._mainloop)
self.start_time = datetime.now()
signals.send_catch_log(signal=signals.engine_started, sender=self.__class__)
self.addtask(self._mainloop, 5.0)
@ -126,7 +128,8 @@ class ExecutionEngine(object):
for args, kwargs in [t for t in self.ports if isinstance(t, tuple)]:
reactor.listenTCP(*args, **kwargs)
self.running = True
reactor.run() # blocking call
if control_reactor:
reactor.run() # blocking call
def stop(self):
"""Stop the execution engine"""
@ -141,7 +144,7 @@ class ExecutionEngine(object):
self.tasks = []
for p in [p for p in self.ports if not isinstance(p, tuple)]:
p.stopListening()
if reactor.running:
if self.control_reactor and reactor.running:
reactor.stop()
signals.send_catch_log(signal=signals.engine_stopped, sender=self.__class__)
@ -185,7 +188,7 @@ class ExecutionEngine(object):
return
# backout enqueing downloads if domain needs it
if self.downloader.needs_backout(domain):
if domain in self.cancelled or self.downloader.needs_backout(domain):
return
# Next pending request from scheduler
@ -206,7 +209,7 @@ class ExecutionEngine(object):
def domain_is_idle(self, domain):
scraping = self._scraping.get(domain)
pending = self.scheduler.domain_has_pending(domain)
pending = self.scheduler.domain_has_pending_requests(domain)
downloading = not self.downloader.domain_is_idle(domain)
haspipe = not self.pipeline.domain_is_idle(domain)
oninit = domain in self.initializing
@ -222,48 +225,49 @@ class ExecutionEngine(object):
def crawl(self, request, spider, domain_priority=0):
domain = spider.domain_name
def _process(response):
assert isinstance(response, (Response, Exception))
def _process_response(response):
assert isinstance(response, (Response, Exception)), "Expecting Response or Exception, got %s" % type(response).__name__
def _onpipelinefinish(pipe_result, item):
# _ can only be an item or a DropItem failure, since other
# failures are caught in ItemPipeine (item/pipeline.py)
if isinstance(pipe_result, Failure):
signals.send_catch_log(signal=signals.item_dropped, sender=self.__class__, item=item, spider=spider, response=response, exception=pipe_result.value)
else:
signals.send_catch_log(signal=signals.item_passed, sender=self.__class__, item=item, spider=spider, response=response, pipe_output=pipe_result)
self.next_request(spider)
def cb_spidermiddleware_output(spmw_result):
def cb_spider_output(output):
def cb_pipeline_output(pipe_result, item):
if isinstance(pipe_result, Failure):
# can only be a DropItem exception, since other exceptions are caught in the Item Pipeline (item/pipeline.py)
signals.send_catch_log(signal=signals.item_dropped, sender=self.__class__, item=item, spider=spider, response=response, exception=pipe_result.value)
else:
signals.send_catch_log(signal=signals.item_passed, sender=self.__class__, item=item, spider=spider, response=response, pipe_output=pipe_result)
self.next_request(spider)
def _onsuccess_per_item(item):
if isinstance(item, ScrapedItem):
log.msg("Scraped %s in <%s>" % (item, request.url), log.INFO, domain=domain)
signals.send_catch_log(signal=signals.item_scraped, sender=self.__class__, item=item, spider=spider, response=response)
piped = self.pipeline.pipe(item, spider)
piped.addBoth(_onpipelinefinish, item)
elif isinstance(item, Request):
signals.send_catch_log(signal=signals.request_received, sender=self.__class__, request=item, spider=spider, response=response)
self.crawl(request=item, spider=spider)
elif item is None:
pass # may be next time.
else:
log.msg("Spider must return Request, ScrapedItem or None, got '%s' while processing %s" % (type(item).__name__, request), log.WARNING, domain=domain)
if domain in self.cancelled:
return
elif isinstance(output, ScrapedItem):
log.msg("Scraped %s in <%s>" % (output, request.url), log.INFO, domain=domain)
signals.send_catch_log(signal=signals.item_scraped, sender=self.__class__, item=output, spider=spider, response=response)
piped = self.pipeline.pipe(output, spider)
piped.addBoth(cb_pipeline_output, output)
elif isinstance(output, Request):
signals.send_catch_log(signal=signals.request_received, sender=self.__class__, request=output, spider=spider, response=response)
self.crawl(request=output, spider=spider)
elif output is None:
pass # may be next time.
else:
log.msg("Spider must return Request, ScrapedItem or None, got '%s' while processing %s" % (type(output).__name__, request), log.WARNING, domain=domain)
def _onsuccess(result):
return deferred_imap(_onsuccess_per_item, result)
return deferred_imap(cb_spider_output, spmw_result)
def _onerror(_failure):
def eb_user(_failure):
if not isinstance(_failure.value, IgnoreRequest):
referer = None if not isinstance(response, Response) else response.request.headers.get('Referer', None)
log.msg("Error while spider was processing <%s> from <%s>: %s" % (request.url, referer, _failure), log.ERROR, domain=domain)
stats.incpath("%s/spider_exceptions/%s" % (domain, _failure.value.__class__.__name__))
def _bugtrap(_failure):
def eb_framework(_failure):
log.msg('FRAMEWORK BUG processing %s: %s' % (request, _failure), log.ERROR, domain=domain)
scd = self.scrape(request, response, spider)
scd.addCallbacks(_onsuccess, _onerror)
scd.addErrback(_bugtrap)
scd = self.spidermiddleware.scrape(request, response, spider)
scd.addCallbacks(cb_spidermiddleware_output, eb_user)
scd.addErrback(eb_framework)
self._scraping[domain].add(response)
scd.addBoth(lambda _: self._scraping[domain].remove(response))
@ -274,15 +278,12 @@ class ExecutionEngine(object):
if not isinstance(ex, IgnoreRequest):
log.msg("Unknown error propagated in %s: %s" % (request, _failure), log.ERROR, domain=domain)
request.deferred.addErrback(lambda _:None)
request.deferred.errback(_failure) #TODO: merge into spider middleware.
request.deferred.errback(_failure) # TODO: merge into spider middleware.
schd = self.schedule(request, spider, domain_priority)
schd.addCallbacks(_process, _cleanfailure)
schd.addCallbacks(_process_response, _cleanfailure)
return schd
def scrape(self, request, response, spider):
return self.spidermiddleware.scrape(request, response, spider)
def schedule(self, request, spider, domain_priority=0):
domain = spider.domain_name
if not self.scheduler.domain_is_open(domain):
@ -310,7 +311,7 @@ class ExecutionEngine(object):
def _add_starter(self, request, spider, domain_priority):
domain = spider.domain_name
if not self.scheduler.is_pending(domain):
if not self.scheduler.domain_is_pending(domain):
self.scheduler.add_domain(domain, priority=domain_priority)
self.starters[domain] = []
deferred = defer.Deferred()
@ -420,6 +421,7 @@ class ExecutionEngine(object):
log.msg("Domain opened", domain=domain)
spider = spider or spiders.fromdomain(domain)
self.cancelled.discard(domain)
self.scheduler.open_domain(domain)
self.downloader.open_domain(domain)
self.pipeline.open_domain(domain)
@ -436,27 +438,14 @@ class ExecutionEngine(object):
self._domain_idle(domain)
dfd.addCallback(_state)
def close_domain(self, domain):
"""Close (cancel) domain and clear all its outstanding requests"""
if domain not in self.cancelled:
self.cancelled.add(domain)
self._close_domain(domain)
def _close_domain(self, domain):
self.downloader.close_domain(domain)
def _domain_idle(self, domain):
"""Called when a domain gets idle. This function is called when there are no
remaining pages to download or schedule. It can be called multiple
times. If some extension raises a DontCloseDomain exception (in the
domain_idle signal handler) the domain is not closed until the next
loop and this function is guaranteed to be called (at least) once again
for this domain.
"""
Called when a domain gets idle. This function is called when there are no
remaining pages to download or scheduled. It can be called multiple
times. If the some extensions raises a DontCloseDomain exception the
domain won't be closed and this function is garanteed to be called
again (at least once) for this domain.
"""
# we get a callback from the downloader which completes closing
#log.msg("Finishing scraping %s" % domain, domain=domain)
#from traceback import print_stack; print_stack()
spider = spiders.fromdomain(domain)
try:
dispatcher.send(signal=signals.domain_idle, sender=self.__class__, domain=domain, spider=spider)
@ -474,6 +463,16 @@ class ExecutionEngine(object):
if self.is_idle() and not self.keep_alive:
self.stop()
def close_domain(self, domain):
"""Close (cancel) domain and clear all its outstanding requests"""
if domain not in self.cancelled:
log.msg("Closing domain", domain=domain)
self.cancelled.add(domain)
self._close_domain(domain)
def _close_domain(self, domain):
self.downloader.close_domain(domain)
def closed_domain(self, domain):
"""
This function is called after the domain has been closed, and throws
@ -489,6 +488,7 @@ class ExecutionEngine(object):
signals.send_catch_log(signal=signals.domain_closed, sender=self.__class__, domain=domain, spider=spider, status=status)
log.msg("Domain closed (%s)" % status, domain=domain)
self.cancelled.discard(domain)
self.starters.pop(domain, None)
self._mainloop()
def getstatus(self):
@ -501,22 +501,24 @@ class ExecutionEngine(object):
"datetime.now()-self.start_time",
"self.is_idle()",
"self.scheduler.is_idle()",
"len(self.scheduler.pending_domains_count)",
"len(self.scheduler.pending_domains)",
"self.downloader.is_idle()",
"len(self.downloader.sites)",
"self.downloader.has_capacity()",
"self.pipeline.is_idle()",
"len(self.pipeline.domaininfo)",
"len(self._scraping)",
]
domain_tests = [
"self.domain_is_idle(domain)",
"self.scheduler.domain_has_pending(domain)",
"self.scheduler.domain_has_pending_requests(domain)",
"len(self.scheduler.pending_requests[domain])",
"self.downloader.outstanding(domain)",
"len(self.downloader.request_queue(domain))",
"len(self.downloader.active_requests(domain))",
"self.pipeline.domain_is_idle(domain)",
"len(self.pipeline.domaininfo[domain])",
"len(self._scraping[domain])",
]
for test in global_tests:
@ -525,8 +527,10 @@ class ExecutionEngine(object):
for domain in self.downloader.sites:
s += "%s\n" % domain
for test in domain_tests:
s += " %-45s : %s\n" % (test, eval(test))
try:
s += " %-50s : %s\n" % (test, eval(test))
except Exception, e:
s += " %-50s : %s (exception)\n" % (test, type(e).__name__)
return s
def st(self): # shortcut for printing engine status (useful in telnet console)

View File

@ -28,8 +28,7 @@ class ExecutionManager(object):
self.control_reactor = control_reactor
if control_reactor:
self._install_signals()
else:
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
reactor.addSystemEventTrigger('before', 'shutdown', self.stop)
if not log.started:
log.start()

View File

@ -4,89 +4,53 @@ The Scrapy Scheduler
from twisted.internet import defer
from scrapy import log
from scrapy.utils.datatypes import PriorityQueue, PriorityStack
from scrapy.conf import settings
class Scheduler(object) :
class Scheduler(object):
"""The scheduler decides what to scrape next. In other words, it defines the
crawling order.
The scheduler schedules websites and requests to be scraped. Individual
web pages that are to be scraped are batched up into a "run" for a website.
As the domain is being scraped, pages that are discovered are added to the
scheduler.
Typical usage:
* next_domain() called each time a domain slot is freed, and return
next domain to be scraped.
* open_domain() called to commence scraping a website
* enqueue_request() called multiple times to enqueue new requests to be downloaded
* next_request() called multiple times when there is capacity to download requests
* close_domain() called when there are no more pages for a website
Notes:
1. The order in which you get back the list of pages to scrape is not
necesarily the order you put them in.
``pending_domains_count`` contains the names of all domains that are to be scheduled.
crawling order. The scheduler schedules websites and requests to be
scraped. Individual web pages that are to be scraped are batched up into a
"run" for a website. New pages discovered through the crawling process are
also added to the scheduler.
"""
def __init__(self):
self.pending_domains_count = {}
self.pending_domains = set()
self.domains_queue = PriorityQueue()
self.pending_requests = {}
self.dfo = settings.get('SCHEDULER_ORDER', '').upper() == 'DFO'
self.dfo = settings['SCHEDULER_ORDER'].upper() == 'DFO'
def domain_is_open(self, domain):
"""Check if scheduler's resources were allocated for a domain"""
return domain in self.pending_requests
def is_pending(self, domain):
def domain_is_pending(self, domain):
"""Check if a domain is waiting to be scraped in domain's queue."""
return domain in self.pending_domains_count
return domain in self.pending_domains
def domain_has_pending(self, domain):
def domain_has_pending_requests(self, domain):
"""Check if are there pending requests for a domain"""
if domain in self.pending_requests:
return bool(self.pending_requests[domain])
def next_domain(self) :
"""Return next domain available to scrape and remove it from available domains queue"""
if self.pending_domains_count:
if self.pending_domains:
domain = self.domains_queue.pop()[0]
if self.pending_domains_count[domain] == 1:
del self.pending_domains_count[domain]
else:
self.pending_domains_count[domain] -= 1
self.pending_domains.remove(domain)
return domain
return None
def add_domain(self, domain, priority=0):
"""This functions schedules a new domain to be scraped, with the given priority.
It doesn't check if the domain is already scheduled.
A domain can be scheduled twice, either with the same or with different
priority.
"""Add a new domain to be scraped, with the given priority. If the
domain is already scheduled, it does nothing.
"""
self.domains_queue.push(domain, priority)
if domain not in self.pending_domains_count:
self.pending_domains_count[domain] = 1
else:
self.pending_domains_count[domain] += 1
if domain not in self.pending_domains:
self.domains_queue.push(domain, priority)
self.pending_domains.add(domain)
def open_domain(self, domain):
"""Allocates resources for maintaining a schedule for domain."""
"""Allocates scheduling resources for the given domain"""
Priority = PriorityStack if self.dfo else PriorityQueue
self.pending_requests[domain] = Priority()
@ -104,24 +68,17 @@ class Scheduler(object) :
``(None, None)`` should be returned if there aren't requests pending
for the domain.
"""
try:
return self.pending_requests[domain].pop()[0] # [1] is priority
except (KeyError, IndexError):
return (None, None)
def close_domain(self, domain) :
"""Called once we are finished scraping a domain.
The scheduler will free any resources associated with the domain.
def close_domain(self, domain):
"""Called when a spider has finished scraping to free any resources
associated with the domain.
"""
try :
del self.pending_requests[domain]
except Exception, inst:
msg = "Could not clear pending pages for domain %s, %s" % (domain, inst)
log.msg(msg, level=log.WARNING)
self.pending_requests.pop(domain, None)
def remove_pending_domain(self, domain):
"""
@ -135,8 +92,8 @@ class Scheduler(object) :
need to call close_domain for open domains.
"""
if not self.domain_is_open(domain):
return self.pending_domains_count.pop(domain, 0)
if domain in self.pending_domains and not self.domain_is_open(domain):
return self.pending_domains.remove(domain)
def is_idle(self):
"""Checks if the schedulers has any request pendings"""