diff --git a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/manager.py b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/manager.py index 74e9531c2..c8f625172 100644 --- a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/manager.py @@ -14,8 +14,6 @@ from scrapy.core.exceptions import NotConfigured from scrapy.contrib.pbcluster.worker.manager import ResponseCode from scrapy.conf import settings -DEFAULT_PRIORITY = settings.getint("DEFAULT_PRIORITY", 20) - def my_import(name): mod = __import__(name) components = name.split('.') @@ -35,9 +33,11 @@ class ClusterNodeBroker(pb.Referenceable): deferred = self._worker.callRemote("set_master", self) except pb.DeadReferenceError: self._set_status(None) - log.msg("ClusterMaster: lost connection to node %s." % (self.name), log.ERROR) + log.msg("ClusterMaster: 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 _eb(failure): + self._logfailure("Error while setting master to worker node", failure) + deferred.addCallbacks(callback=self._set_status, errback=_eb) def status_as_dict(self, verbosity=1): if verbosity == 0: @@ -61,6 +61,85 @@ class ClusterNodeBroker(pb.Referenceable): status["loadavg"] = self.loadavg return status + def update_status(self): + """Update status from this worker. This is called periodically.""" + try: + deferred = self._worker.callRemote("status") + except pb.DeadReferenceError: + self._set_status(None) + log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR) + else: + def _eb(failure): + self._logfailure("Error while updating status", failure) + deferred.addCallbacks(callback=self._set_status, errback=_eb) + + def stop(self, domain): + try: + deferred = self._worker.callRemote("stop", domain) + except pb.DeadReferenceError: + self._set_status(None) + log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR) + else: + def _eb(failure): + self._logfailure("Error while stopping domain=%s" % domain, failure) + deferred.addCallbacks(callback=self._set_status, errback=_eb) + + def run(self, domain_info): + """Run the given domain. + + domain_info is a dict of keys: + domain - the domain to run + settings - the settings to use + priority - the priority to use + """ + + domain = domain_info['domain'] + dsettings = domain_info['settings'] + priority = domain_info['priority'] + + def _run_errback(failure): + 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) + + 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) + 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) + + try: + log.msg("ClusterMaster: Running domain=%s at worker=%s" % (domain, self.name), log.DEBUG) + deferred = self._worker.callRemote("run", domain, dsettings) + except pb.DeadReferenceError: + self._set_status(None) + log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR) + else: + deferred.addCallbacks(callback=_run_callback, errback=_run_errback) + + def remote_update(self, worker_status, domain, domain_status): + """Called remotely form worker when domains finish to update status""" + self._set_status(worker_status) + if domain in self.master.loading and domain_status == "running": + self.master.loading.remove(domain) + self.master.statistics["domains"]["running"].add(domain) + elif domain_status in ("done", "terminated"): + self.master.statistics["domains"]["running"].remove(domain) + self.master.statistics["domains"]["scraped"][domain] = self.master.statistics["domains"]["scraped"].get(domain, 0) + 1 + self.master.statistics["scraped_count"] = self.master.statistics.get("scraped_count", 0) + 1 + if domain in self.master.statistics["domains"]["lost"]: + self.master.statistics["domains"]["lost"].remove(domain) + log.msg("ClusterMaster: Changed status to <%s> for domain=%s at worker=%s" % (domain_status, domain, self.name)) + + def _logfailure(self, msg, failure): + log.msg("ClusterMaster: %s (worker=%s)\n%s" % (msg, self.name, failure), log.ERROR) + def _set_status(self, status): if not status: self.alive = False @@ -85,72 +164,6 @@ class ClusterNodeBroker(pb.Referenceable): self.run(pending) self.master.loading.append(pending['domain']) - def update_status(self): - try: - deferred = self._worker.callRemote("status") - except pb.DeadReferenceError: - self._set_status(None) - log.msg("ClusterMaster: 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 stop(self, domain): - try: - deferred = self._worker.callRemote("stop", domain) - except pb.DeadReferenceError: - self._set_status(None) - log.msg("ClusterMaster: 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 run(self, domain_info): - """Run the given domain. - - domain_info keys: - domain - the domain to run - settings - the settings to use - priority - the priority to use - """ - - def _run_errback(reason): - log.msg(reason, log.ERROR) - self.master.loading.remove(domain_info['domain']) - self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority'] - 1) - log.msg("ClusterMaster: Domain %s rescheduled: lost connection to node." % domain_info['domain'], log.WARNING) - - def _run_callback(status): - if status['callresponse'][0] == ResponseCode.NO_FREE_SLOT: - # slots are complete. Reschedule in master with priority reduced by one. - # self.master.loading check should avoid this to happen - self.master.loading.remove(domain_info['domain']) - self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority'] - 1) - log.msg("ClusterMaster: Domain %s rescheduled: no availble processes in worker" % domain_info['domain'], log.WARNING) - elif status['callresponse'][0] == ResponseCode.DOMAIN_ALREADY_RUNNING: - # domain already running in node. Reschedule with same priority. - # self.master.loading check should avoid this to happen - self.master.loading.remove(domain_info['domain']) - self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority']) - log.msg("ClusterMaster: Domain %s rescheduled: already running in node." % domain_info['domain'], log.WARNING) - - try: - deferred = self._worker.callRemote("run", domain_info["domain"], domain_info["settings"]) - except pb.DeadReferenceError: - self._set_status(None) - log.msg("ClusterMaster: Lost connection to node %s." % (self.name), log.ERROR) - else: - deferred.addCallbacks(callback=_run_callback, errback=_run_errback) - - def remote_update(self, status, domain, domain_status): - self._set_status(status) - if domain in self.master.loading and domain_status == "running": - self.master.loading.remove(domain) - self.master.statistics["domains"]["running"].add(domain) - elif domain_status == "scraped": - self.master.statistics["domains"]["running"].remove(domain) - self.master.statistics["domains"]["scraped"][domain] = self.master.statistics["domains"]["scraped"].get(domain, 0) + 1 - self.master.statistics["scraped_count"] = self.master.statistics.get("scraped_count", 0) + 1 - if domain in self.master.statistics["domains"]["lost"]: - self.master.statistics["domains"]["lost"].remove(domain) class ScrapyPBClientFactory(pb.PBClientFactory): @@ -163,8 +176,9 @@ class ScrapyPBClientFactory(pb.PBClientFactory): def clientConnectionLost(self, *args, **kargs): pb.PBClientFactory.clientConnectionLost(self, *args, **kargs) - del self.master.nodes[self.nodename] - log.msg("ClusterMaster: Lost connection to %s. Node removed" % self.nodename ) + self.master.remove_node(self.nodename) + log.msg("ClusterMaster: Lost connection to worker=%s. Node removed" % self.nodename) + class ClusterMaster(object): @@ -172,7 +186,9 @@ class ClusterMaster(object): if not settings.getbool('CLUSTER_MASTER_ENABLED'): raise NotConfigured - if not settings['CLUSTER_MASTER_STATEFILE']: + + self.statefile = settings['CLUSTER_MASTER_STATEFILE'] + if not self.statefile: raise NotConfigured("ClusterMaster: Missing CLUSTER_MASTER_STATEFILE setting") # import groups settings @@ -182,12 +198,14 @@ class ClusterMaster(object): self.get_spider_groupsettings = lambda x: {} # load pending domains try: - statefile = open(settings["CLUSTER_MASTER_STATEFILE"], "r") + statefile = open(self.statefile, "r") self.pending = pickle.load(statefile) + log.msg("ClusterMaster: Loaded state from %s" % self.statefile) except IOError: self.pending = [] self.loading = [] self.nodes = {} + self.nodesconf = settings.get('CLUSTER_MASTER_NODES', {}) self.start_time = datetime.datetime.utcnow() # for more info about statistics see self.update_nodes() and ClusterNodeBroker.remote_update() self.statistics = {"domains": {"running": set(), "scraped": {}, "lost_count": {}, "lost": set()}, "scraped_count": 0 } @@ -201,34 +219,32 @@ class ClusterMaster(object): def load_nodes(self): """Loads nodes listed in CLUSTER_MASTER_NODES setting""" - for name, hostport in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): + for name, hostport in self.nodesconf.iteritems(): self.load_node(name, hostport) def load_node(self, name, hostport): """Creates the remote reference for a worker node""" server, port = hostport.split(":") port = int(port) - log.msg("ClusterMaster: Connecting to worker %s (%s)..." % (name, hostport)) + log.msg("ClusterMaster: Connecting to worker=%s (%s)..." % (name, hostport)) factory = ScrapyPBClientFactory(self, name) try: reactor.connectTCP(server, port, factory) except Exception, err: - log.msg("ClusterMaster: Could not connect to worker %s (%s): %s" % (name, hostport, err), log.ERROR) + log.msg("ClusterMaster: Could not connect to worker=%s (%s): %s" % (name, hostport, err), log.ERROR) else: - def _errback(_reason): - log.msg("ClusterMaster: Could not connect to worker %s (%s): %s" % (name, hostport, _reason), log.ERROR) + def _eb(failure): + self._logfailure("Error while loading worker node", failure) d = factory.getRootObject() - d.addCallbacks(callback=lambda obj: self.add_node(obj, name), errback=_errback) + d.addCallbacks(callback=lambda obj: self.add_node(obj, name), errback=_eb) def update_nodes(self): """Update worker nodes statistics""" - for name, hostport in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): + for name, hostport in self.nodesconf.iteritems(): if name in self.nodes and self.nodes[name].alive: - log.msg("ClusterMaster: Updating stats from worker node %s (%s)" % (name, hostport)) self.nodes[name].update_status() else: - log.msg("ClusterMaster: Reloading worker node %s (%s)" % (name, hostport)) self.load_node(name, hostport) real_running = set(self.running.keys()) @@ -241,7 +257,10 @@ class ClusterMaster(object): """Add node given its node""" node = ClusterNodeBroker(cworker, name, self) self.nodes[name] = node - log.msg("ClusterMaster: Added cluster worker %s" % name) + log.msg("ClusterMaster: Added worker=%s" % name) + + def remove_node(self, nodename): + del self.nodes[nodename] def disable_node(self, name): self.nodes[name].available = False @@ -249,10 +268,7 @@ class ClusterMaster(object): def enable_node(self, name): self.nodes[name].available = True - def remove_node(self, nodename): - raise NotImplemented - - def schedule(self, domains, spider_settings=None, priority=DEFAULT_PRIORITY): + def schedule(self, domains, spider_settings=None, priority=20): """Schedule the given domains, with the given priority""" insert_pos = len([p for p in self.pending if ['priority'] <= priority]) for domain in domains: @@ -324,9 +340,9 @@ class ClusterMaster(object): def _engine_started(self): self.load_nodes() - scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL')) + scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL', 60)) def _engine_stopped(self): - with open(settings["CLUSTER_MASTER_STATEFILE"], "w") as f: + with open(self.statefile, "w") as f: pickle.dump(self.pending, f) - log.msg("ClusterMaster: state saved in %s" % settings["CLUSTER_MASTER_STATEFILE"]) + log.msg("ClusterMaster: Saved state in %s" % self.statefile) diff --git a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/web.py b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/web.py index 0d8916496..42129ea0a 100644 --- a/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/web.py +++ b/scrapy/branches/cluster-refactor/scrapy/contrib/pbcluster/master/web.py @@ -4,7 +4,7 @@ from pydispatch import dispatcher from scrapy.spider import spiders from scrapy.management.web import banner, webconsole_discover_module -from scrapy.contrib.pbcluster.master.manager import ClusterMaster, DEFAULT_PRIORITY +from scrapy.contrib.pbcluster.master.manager import ClusterMaster from scrapy.utils.serialization import serialize class ClusterMasterWeb(ClusterMaster): @@ -61,7 +61,7 @@ class ClusterMasterWeb(ClusterMaster): else: sep = "\r" domains = args["schedule"] - priority = int(args.get("priority", [DEFAULT_PRIORITY])[0]) + priority = int(args.get("priority", [20])[0]) # spider settings slist = args.get("settings", [""])[0].split(sep) @@ -164,7 +164,7 @@ class ClusterMasterWeb(ClusterMaster): s += "
\n" s += "Priority:
\n" - s += "%s" % DEFAULT_PRIORITY + s += "%s" % 20 s += "
\n" # spider settings