enabled unsafeTracebacks to master for sending full tracebacks to workers, splitted master scheduled() method in 2 methods: schedule() and reschedule()

--HG--
extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40337
This commit is contained in:
Pablo Hoffman 2008-10-23 12:43:31 +00:00
parent 215151dd86
commit ce3bbd1a2b
1 changed files with 25 additions and 11 deletions

View File

@ -23,8 +23,9 @@ def my_import(name):
class ClusterNodeBroker(pb.Referenceable):
def __init__(self, remote, name, master):
self._worker = remote
def __init__(self, worker, name, master):
self.unsafeTracebacks = True
self._worker = worker
self.alive = False
self.name = name
self.master = master
@ -101,18 +102,18 @@ class ClusterNodeBroker(pb.Referenceable):
self._logfailure("Error while running domain=%s" % domain, failure)
self.master.loading.remove(domain)
newprio = priority - 1 # increase priority for reschedule
self.master.schedule([domain], dsettings, newprio)
self.master.reschedule([domain], dsettings, newprio, reason="error while try to run it")
def _run_callback(status):
if status['callresponse'][0] == ResponseCode.NO_FREE_SLOT:
log.msg("ClusterMaster: No available slots at worker=%s when trying to run domain=%s" % (self.name, domain), log.WARNING)
self.master.loading.remove(domain)
newprio = priority - 1 # increase priority for rerunning asap
self.master.schedule([domain], dsettings, newprio)
self.master.reschedule([domain], dsettings, newprio, reason="no available slots at worker=%s" % self.name)
elif status['callresponse'][0] == ResponseCode.DOMAIN_ALREADY_RUNNING:
log.msg("ClusterMaster: Already running domain=%s at worker=%s" % (domain, self.name), log.WARNING)
self.master.loading.remove(domain)
self.master.schedule([domain], dsettings, priority)
self.master.reschedule([domain], dsettings, priority, reason="domain already running at worker=%s" % self.name)
try:
log.msg("ClusterMaster: Running domain=%s at worker=%s" % (domain, self.name), log.DEBUG)
@ -126,6 +127,7 @@ class ClusterNodeBroker(pb.Referenceable):
def remote_update(self, worker_status, domain, domain_status):
"""Called remotely form worker when domains finish to update status"""
self._set_status(worker_status)
raise Exception
if domain in self.master.loading and domain_status == "running":
self.master.loading.remove(domain)
self.master.statistics["domains"]["running"].add(domain)
@ -157,9 +159,9 @@ class ClusterNodeBroker(pb.Referenceable):
# when there is no loading domain or in the next status update. This way also we load the nodes softly
if self.available and free_slots > 0 and self.master.pending:
pending = self.master.pending.pop(0)
# if domain already running in some node, reschedule with same priority (so will be moved to run later)
# if domain already running in some node, reschedule with same priority (so it will be run later)
if pending['domain'] in self.master.running or pending['domain'] in self.master.loading:
self.master.schedule([pending['domain']], pending['settings'], pending['priority'])
self.master.reschedule([pending['domain']], pending['settings'], pending['priority'], reason="domain already running in other worker")
else:
self.run(pending)
self.master.loading.append(pending['domain'])
@ -171,6 +173,7 @@ class ScrapyPBClientFactory(pb.PBClientFactory):
def __init__(self, master, nodename):
pb.PBClientFactory.__init__(self)
self.unsafeTracebacks = True
self.master = master
self.nodename = nodename
@ -234,7 +237,7 @@ class ClusterMaster(object):
log.msg("ClusterMaster: Could not connect to worker=%s (%s): %s" % (name, hostport, err), log.ERROR)
else:
def _eb(failure):
self._logfailure("Error while loading worker node", failure)
log.msg("ClusterMaster: Could not connect to worker=%s (%s): %s" % (name, hostport, failure.value), log.ERROR)
d = factory.getRootObject()
d.addCallbacks(callback=lambda obj: self.add_node(obj, name), errback=_eb)
@ -268,8 +271,9 @@ class ClusterMaster(object):
def enable_node(self, name):
self.nodes[name].available = True
def schedule(self, domains, spider_settings=None, priority=20):
"""Schedule the given domains, with the given priority"""
def _schedule(self, domains, spider_settings=None, priority=20):
"""Private method which performs the schedule of the given domains,
with the given priority. Used for both scheduling and rescheduling."""
insert_pos = len([p for p in self.pending if ['priority'] <= priority])
for domain in domains:
pd = self.get_first_pending(domain)
@ -283,7 +287,17 @@ class ClusterMaster(object):
final_spider_settings.update(self.global_settings)
final_spider_settings.update(spider_settings or {})
self.pending.insert(insert_pos, {'domain': domain, 'settings': final_spider_settings, 'priority': priority})
log.msg("ClusterMaster: Scheduled domain=%s priority=%s" % (domain, priority), log.DEBUG)
def schedule(self, domains, spider_settings=None, priority=20):
"""Schedule the given domains, with the given priority"""
self._schedule(domains, spider_settings, priority)
log.msg("clustermaster: Scheduled domains=%s with priority=%s" % (','.join(domains), priority), log.DEBUG)
def reschedule(self, domains, spider_settings=None, priority=20, reason=None):
"""Reschedule the given domains, with the given priority"""
self._schedule(domains, spider_settings, priority)
log.msg("clustermaster: Rescheduled domains=%s with priority=%s reason='%s'" % (','.join(domains), priority, reason), log.DEBUG)
def stop(self, domains):
"""Stop the given domains"""