From 53838d6a7dbdb41f8bf0ed7ff572d7db9344a378 Mon Sep 17 00:00:00 2001 From: olveyra Date: Mon, 30 Jun 2008 14:20:27 +0000 Subject: [PATCH] pbcluster commit --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%4031 --- .../contrib/pbcluster/master/manager.py | 146 ++++++++---------- .../scrapy/contrib/pbcluster/master/web.py | 47 +++--- .../contrib/pbcluster/worker/manager.py | 104 +++---------- .../contrib/pbcluster/worker/testworker.py | 10 +- 4 files changed, 106 insertions(+), 201 deletions(-) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index 102f35e10..12d26b945 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -24,45 +24,67 @@ for val, attr in priorities.items(): setattr(sys.modules[__name__], "PRIORITY_%s" % attr, val ) class Node: - def __init__(self, remote, status, name): + def __init__(self, remote, status, name, master): self.__remote = remote self._set_status(status) self.name = name + self.master = master def _set_status(self, status): if not status: self.available = False else: + print ">----------<" self.available = True self.running = status['running'] - self.pending = status['pending'] + self.closing = status['closing'] self.maxproc = status['maxproc'] self.starttime = status['starttime'] self.timestamp = status['timestamp'] self.loadavg = status['loadavg'] self.logdir = status['logdir'] - self.lastcallresponse = status['callresponse'] + print "Running: %s" % self.running + free_slots = self.maxproc - len(self.running) + while free_slots > 0 and self.master.pending: + print "Free slots %s" % free_slots + pending = self.master.pending.pop(0) + self.run(pending) + free_slots -= 1 - def _remote_call(self, function, *args): + def get_status(self): try: - deferred = self.__remote.callRemote(function, *args) + deferred = self.__remote.callRemote("status") except pb.DeadReferenceError: self._set_status(None) log.msg("Lost connection to node %s." % (self.name), log.ERROR) else: deferred.addCallbacks(callback=self._set_status, errback=lambda reason: log.msg(reason, log.ERROR)) - def get_status(self): - self._remote_call("status") + def stop(self, domain): + try: + deferred = self.__remote.callRemote("stop", domain) + except pb.DeadReferenceError: + self._set_status(None) + log.msg("Lost connection to node %s." % (self.name), log.ERROR) + else: + deferred.addCallbacks(callback=self._set_status, errback=lambda reason: log.msg(reason, log.ERROR)) - def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL): - self._remote_call("schedule", domains, spider_settings, priority) + def run(self, pending): - def stop(self, domains): - self._remote_call("stop", domains) + def _run_callback(status): + if status['callresponse'][0] == 1: + #slots are complete. Reschedule in master. This is a security issue because could happen that the slots were completed since last status update by another cluster (thinking at future with full-distributed worker-master clusters) + self.master.schedule(pending['domain'], pending['settings'], pending['priority']) + log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING) + self._set_status(status) - def remove(self, domains): - self._remote_call("remove", domains) + try: + deferred = self.__remote.callRemote("run", pending["domain"], pending["settings"]) + except pb.DeadReferenceError: + self._set_status(None) + log.msg("Lost connection to node %s." % (self.name), log.ERROR) + else: + deferred.addCallbacks(callback=_run_callback, errback=lambda reason:log.msg(reason, log.ERROR)) class ClusterMaster(object): @@ -70,7 +92,7 @@ class ClusterMaster(object): if not settings.getbool('CLUSTER_MASTER_ENABLED'): raise NotConfigured self.nodes = {} - self.queue = [] + self.pending = [] dispatcher.connect(self._engine_started, signal=signals.engine_started) def load_nodes(self): @@ -84,7 +106,7 @@ class ClusterMaster(object): d.addCallbacks(callback=lambda obj: self.add_node(obj, _name), errback=_errback) """Loads nodes from the CLUSTER_MASTER_NODES setting""" - + for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): if name not in self.nodes: server, port = url.split(":") @@ -98,14 +120,14 @@ class ClusterMaster(object): log.msg("Could not connect to node %s in %s: %s." % (name, url, reason), log.ERROR) else: _make_callback(factory, name, url) - + def update_nodes(self): for node in self.nodes.itervalues(): node.get_status() def add_node(self, cworker, name): """Add node given its node""" - node = Node(cworker, None, name) + node = Node(cworker, None, name, self) node.get_status() self.nodes[name] = node log.msg("Added cluster worker %s" % name) @@ -113,11 +135,16 @@ class ClusterMaster(object): def remove_node(self, nodename): raise NotImplemented - def schedule(self, domains, spider_settings=None, nodename=None, priority=PRIORITY_NORMAL): - if nodename: - self.nodes[nodename].schedule(domains, spider_settings, priority) - else: - self._dispatch_domains(domains, spider_settings, priority) + def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL): + i = 0 + for p in self.pending: + if p['priority'] <= priority: + i += 1 + else: + break + for domain in domains: + self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority}) + self.update_nodes() def stop(self, domains): to_stop = {} @@ -129,19 +156,21 @@ class ClusterMaster(object): to_stop[node.name].append(domain) for nodename, domains in to_stop.iteritems(): - self.nodes[nodename].stop(domains) + for domain in domains: + self.nodes[nodename].stop(domain) def remove(self, domains): - to_remove = {} - for domain in domains: - node = self.pending.get(domain, None) - if node: - if node.name not in to_remove: - to_remove[node.name] = [] - to_remove[node.name].append(domain) + """Remove all scheduled instances of the given domains (if it hasn't + started yet). Otherwise use stop()""" - for nodename, domains in to_remove.iteritems(): - self.nodes[nodename].remove(domains) + for domain in domains: + to_remove = [] + for p in self.pending: + if p['domain'] == domain: + to_remove.append(p) + + for p in to_remove: + self.pending.remove(p) def discard(self, domains): """Stop and remove all running and pending instances of the given @@ -158,63 +187,10 @@ class ClusterMaster(object): d[proc['domain']] = node return d - @property - def pending(self): - """Return dict of pending domains as domain -> node""" - d = {} - for node in self.nodes.itervalues(): - for p in node.pending: - d[p['domain']] = node - return d - @property def available_nodes(self): return (node for node in self.nodes.itervalues() if node.available) - def _dispatch_domains(self, domains, spider_settings, priority): - """Schedule the given domains in the availables nodes as good as - possible. The algorithm follows the next rules (in order): - - 1. search for nodes with available capacity(running < maxproc) and (if - any) schedules the domains there - - 2. if there isn't any node with available capacity it schedules the - domain in the node with the smallest number of pending spiders - """ - - to_schedule = {} # domains to schedule per node - pending_node = [] # list of #pending, node - - for node in self.available_nodes: - capacity = node.maxproc - len(node.running) - #order nodes in pending_node according to insertion position, calculated from priority comparison, for stage 2. - i = 0 - for p in node.pending: - if p['priority'] <= priority: - i += 1 - else: - break - bisect.insort(pending_node, (i, node)) - - #stage 1: use available capacity - to_schedule[node.name] = [] - while domains and capacity > 0: - to_schedule[node.name].append(domains.pop(0)) - capacity -= 1 - if not domains: - break - - #stage 2: queue in pendings the remaining domains. - # a) pops out minor insertion-point node b) schedules the domain c) reinserts the node in list with insertion-point incremented by one. - for domain in domains: - insert_point, node = pending_node.pop(0) - to_schedule[node.name].append(domain) - bisect.insort(pending_node, (insert_point+1, node)) - - for nodename, domains in to_schedule.iteritems(): - if domains: - self.nodes[nodename].schedule(domains, spider_settings, priority) - def _engine_started(self): self.load_nodes() scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL')) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py index 7306f7b03..6a54172f3 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py @@ -29,14 +29,14 @@ class ClusterMasterWeb(ClusterMaster): s += "

Home

\n" s += "\n" - s += "\n" + s += "\n" for node in self.nodes.itervalues(): #chkbox = "" % domain if node.status in ["up", "idle"] else " " nodelink = "%s" % (node.name, node.name) chkbox = " " loadavg = "%.2f %.2f %.2f" % node.loadavg - s += "\n" % \ - (chkbox, nodelink, node.available, len(node.running), node.maxproc, len(node.pending), loadavg) + s += "\n" % \ + (chkbox, nodelink, node.available, len(node.running), node.maxproc, loadavg) s += "
 NameAvailableRunningPendingLoad.avg
 NameAvailableRunningLoad.avg
%s%s%s%d/%d%d%s
%s%s%s%d/%d%s
\n" s += "\n" @@ -51,8 +51,7 @@ class ClusterMasterWeb(ClusterMaster): self.update_nodes() if "schedule" in args: - node = args["node"][0] if "node" in args else None - self.schedule(args["schedule"], nodename=node, priority=eval(args["priority"][0])) + self.schedule(args["schedule"], priority=eval(args["priority"][0])) if "stop" in args: self.stop(args["stop"]) @@ -90,20 +89,6 @@ class ClusterMasterWeb(ClusterMaster): s += "\n" else: s += "

No running domains on %s

\n" % node.name - - # pending domains - s += "

Pending domains

\n" - if node.pending: - s += "
\n" - s += "\n" - s += "\n" % node.name - s += "

\n" % node.name - s += "
\n" - else: - s += "

No pending domains on %s

\n" % node.name return str(s) @@ -113,7 +98,7 @@ class ClusterMasterWeb(ClusterMaster): enabled_domains = set(spiders.asdict(include_disabled=False).keys()) print "Enabled domains: %s" % len(enabled_domains) - inactive_domains = enabled_domains - set(self.running.keys() + self.pending.keys()) + inactive_domains = enabled_domains - set(self.running.keys() + [p['domain'] for p in self.pending]) s = self.render_header() @@ -126,14 +111,6 @@ class ClusterMasterWeb(ClusterMaster): s += "\n" % domain s += "\n" s += "
\n" - s += "Node (only available nodes shown):
\n" - s += "
\n" s += "Priority:
\n" s += "\n" + for p in self.pending: + s += "\n" % (p['domain'], p['domain'],p['priority']) + s += "\n" + s += "

\n" + s += "\n" + else: + s += "

No pending domains

\n" + return str(s) def render_header(self): diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py index 8aa208f71..7005aca60 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -42,7 +42,6 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): def processEnded(self, status_object): log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) del self.procman.running[self.domain] - self.procman.next_pending() class ClusterWorker(pb.Root): @@ -53,99 +52,46 @@ class ClusterWorker(pb.Root): self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') self.logdir = settings['CLUSTER_WORKER_LOGDIR'] self.running = {} - self.pending = [] self.starttime = time.time() port = settings.getint('CLUSTER_WORKER_PORT') scrapyengine.listenTCP(port, pb.PBServerFactory(self)) - def remote_schedule(self, domains, spider_settings=None, priority=20): - """Schedule new domains to be crawled in a separate processes""" - - responses = [] - for domain in domains: - if len(self.running) < self.maxproc and domain not in self.running: - self._run(domain, spider_settings) - responses.append("Started %s" % self.running[domain]) - else: - i = 0 - for p in self.pending: - if p['priority'] <= priority: - i += 1 - else: - break - self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority}) - responses.append("Scheduled domain %s at position %s in queue" % (domain, i)) - return self.status(responses) - - def remote_stop(self, domains): - """Stop running domains. For removing pending (not yet started) domains - use remove() instead""" - - responses = [] - for domain in domains: - if domain in self.running: - proc = self.running[domain] - log.msg("ClusterWorker: Sending shutdown signal to domain=%s, pid=%d" % (domain, proc.pid)) - proc.transport.signalProcess('INT') - proc.status = "closing" - responses.append("Stopped process %s" % proc) - else: - responses.append("%s: domain not running." % domain) - return self.status(responses) - - def remote_remove(self, domains): - """Remove all scheduled instances of the given domains (if it hasn't - started yet). Otherwise use stop()""" - - responses = [] - for domain in domains: - to_remove = [] - for p in self.pending: - if p['domain'] == domain: - to_remove.append(p) - - for p in to_remove: - self.pending.remove(p) - responses.append("Unscheduled domain %s" % domain) - return self.status(responses) + def remote_stop(self, domain): + """Stop running domain.""" + if domain in self.running: + proc = self.running[domain] + log.msg("ClusterWorker: Sending shutdown signal to domain=%s, pid=%d" % (domain, proc.pid)) + proc.transport.signalProcess('INT') + proc.status = "closing" + return self.status(0, "Stopped process %s" % proc) + else: + return self.status(1, "%s: domain not running." % domain) def remote_status(self): return self.status() - def status(self, response="Status Response"): + def status(self, rcode=0, rstring=None): status = {} - status["pending"] = self.pending status["running"] = [ self.running[k].as_dict() for k in self.running.keys() ] status["starttime"] = self.starttime status["timestamp"] = time.time() status["maxproc"] = self.maxproc status["loadavg"] = os.getloadavg() status["logdir"] = self.logdir - status["callresponse"] = response + status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.") return status - - def next_pending(self): - """Run the next domain in the pending list, which is not already running""" - - if len(self.running) >= self.maxproc: - return - for p in self.pending: - if p['domain'] not in self.running: - self._run(p['domain'], p['settings']) - self.pending.remove(p) - return - - def _run(self, domain, spider_settings=None): - """Spawn process to run the given domain. Don't call this method - directly. Instead use schedule().""" - - logfile = os.path.join(self.logdir, domain, time.strftime("%FT%T.log")) - if not os.path.exists(os.path.dirname(logfile)): - os.makedirs(os.path.dirname(logfile)) - scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings) - - args = [sys.executable, sys.argv[0], 'crawl', domain] - proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) - self.running[domain] = scrapy_proc + def remote_run(self, domain, spider_settings=None): + """Spawn process to run the given domain.""" + if len(self.running) < self.maxproc: + logfile = os.path.join(self.logdir, domain, time.strftime("%FT%T.log")) + if not os.path.exists(os.path.dirname(logfile)): + os.makedirs(os.path.dirname(logfile)) + scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings) + + args = [sys.executable, sys.argv[0], 'crawl', domain] + proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) + self.running[domain] = scrapy_proc + return self.status(0, "Started process %s." % scrapy_proc) + return self.status(1, "No free slot to run another process.") diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py index 23ef4c355..c4498e374 100755 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py @@ -13,16 +13,10 @@ sys.argv.pop(0) if not sys.argv: d.addCallback(lambda object: object.callRemote("status")) -elif sys.argv[0] == "-d": - try: - priority = eval(sys.argv[2]) - except: - priority = 20 - d.addCallback(lambda object: object.callRemote("schedule", [sys.argv[1]], priority=priority)) elif sys.argv[0] == "-s": - d.addCallback(lambda object: object.callRemote("stop", [sys.argv[1]])) + d.addCallback(lambda object: object.callRemote("stop", sys.argv[1])) elif sys.argv[0] == "-r": - d.addCallback(lambda object: object.callRemote("remove", [sys.argv[1]])) + d.addCallback(lambda object: object.callRemote("run", sys.argv[1])) d.addCallbacks(callback = util.println, errback = lambda reason: 'error: '+str(reason.value)) d.addCallback(lambda _: reactor.stop())