From 585f35fb3ab9720f2886c183c16d9adb2a4e280f Mon Sep 17 00:00:00 2001 From: olveyra Date: Fri, 11 Jul 2008 16:21:23 +0000 Subject: [PATCH] Cluster Master improvements: - rescheduling now goes with original priority decreased by one - Added GLOBAL_CLUSTER_SETTINGS - Added PB remote method load_node so the worker also can initiate a connection --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%4054 --- .../contrib/pbcluster/master/manager.py | 50 ++++++++++++------- 1 file changed, 33 insertions(+), 17 deletions(-) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index ad01c4122..ac2d039ec 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -78,8 +78,8 @@ class Node: 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'], PRIORITY_NOW) + #slots are complete. Reschedule in master with priority reduced by one. This is a security issue because offen happens that the slots were completed and not yet notified because of the asynchronous response from worker. + self.master.schedule([pending['domain']], pending['settings'], pending['priority'] - 1) log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING) self._set_status(status) @@ -91,7 +91,7 @@ class Node: else: deferred.addCallbacks(callback=_run_callback, errback=lambda reason:log.msg(reason, log.ERROR)) -class ClusterMaster(object): +class ClusterMaster(pb.Root): def __init__(self): @@ -110,11 +110,26 @@ class ClusterMaster(object): self.pending = [] self.nodes = {} + + self.global_settings = {} + #load cluster global settings + for sname in settings.getlist('GLOBAL_CLUSTER_SETTINGS'): + self.global_settings[sname] = settings[sname] + dispatcher.connect(self._engine_started, signal=signals.engine_started) dispatcher.connect(self._engine_stopped, signal=signals.engine_stopped) + port = settings.getint('CLUSTER_MASTER_PORT') + scrapyengine.listenTCP(port, pb.PBServerFactory(self)) def load_nodes(self): + """Loads nodes from the CLUSTER_MASTER_NODES setting""" + + for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): + self.load_node(name, url) + + def load_node(self, name, url): + def _make_callback(_factory, _name, _url): def _errback(_reason): @@ -123,21 +138,21 @@ class ClusterMaster(object): d = _factory.getRootObject() d.addCallbacks(callback=lambda obj: self.add_node(obj, _name), errback=_errback) - """Loads nodes from the CLUSTER_MASTER_NODES setting""" + if name not in self.nodes: + server, port = url.split(":") + port = eval(port) + log.msg("Connecting to cluster worker %s..." % name) + log.msg("Server: %s, Port: %s" % (server, port)) + factory = pb.PBClientFactory() + try: + reactor.connectTCP(server, port, factory) + except Exception, err: + log.msg("Could not connect to node %s in %s: %s." % (name, url, reason), log.ERROR) + else: + _make_callback(factory, name, url) - for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): - if name not in self.nodes: - server, port = url.split(":") - port = eval(port) - log.msg("Connecting to cluster worker %s..." % name) - log.msg("Server: %s, Port: %s" % (server, port)) - factory = pb.PBClientFactory() - try: - reactor.connectTCP(server, port, factory) - except Exception, err: - log.msg("Could not connect to node %s in %s: %s." % (name, url, reason), log.ERROR) - else: - _make_callback(factory, name, url) + def remote_connect(self, name, url): + self.load_node(name, url) def update_nodes(self): for node in self.nodes.itervalues(): @@ -162,6 +177,7 @@ class ClusterMaster(object): break for domain in domains: final_spider_settings = self.get_spider_groupsettings(domain) + final_spider_settings.update(self.global_settings) final_spider_settings.update(spider_settings or {}) self.pending.insert(i, {'domain': domain, 'settings': final_spider_settings, 'priority': priority}) self.update_nodes()