From d6c52d51ed2b8118a2d68dd56f035793facd6895 Mon Sep 17 00:00:00 2001 From: Daniel Grana Date: Thu, 16 Apr 2009 18:58:56 +0000 Subject: [PATCH] cluster: promote new code as replacement of old pbcluster --HG-- rename : scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py => scrapy/trunk/scrapy/contrib/cluster/crawler/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/crawler/manager.py => scrapy/trunk/scrapy/contrib/cluster/crawler/manager.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/hooks/__init__.py => scrapy/trunk/scrapy/contrib/cluster/hooks/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/hooks/svn.py => scrapy/trunk/scrapy/contrib/cluster/hooks/svn.py rename : scrapy/trunk/scrapy/contrib/pbcluster/master/__init__.py => scrapy/trunk/scrapy/contrib/cluster/master/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/master/manager.py => scrapy/trunk/scrapy/contrib/cluster/master/manager.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/master/web.py => scrapy/trunk/scrapy/contrib/cluster/master/web.py rename : scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt => scrapy/trunk/scrapy/contrib/cluster/master/ws_api.txt rename : scrapy/trunk/scrapy/contrib/pbcluster/tools/scrapy-cluster-ctl.py => scrapy/trunk/scrapy/contrib/cluster/tools/scrapy-cluster-ctl.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/tools/test-worker.py => scrapy/trunk/scrapy/contrib/cluster/tools/test-worker.py rename : scrapy/trunk/scrapy/contrib/pbcluster/worker/__init__.py => scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/worker/manager.py => scrapy/trunk/scrapy/contrib/cluster/worker/manager.py extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%401064 --- .../trunk/scrapy/contrib/cluster/__init__.py | 3 + .../crawler/__init__.py | 0 .../cluster/crawler/manager.py | 0 .../cluster/hooks/__init__.py | 0 .../cluster/hooks/svn.py | 0 .../{pbcluster => cluster}/master/__init__.py | 0 .../cluster/master/manager.py | 2 +- .../cluster/master/web.py | 2 +- .../{pbcluster => cluster}/master/ws_api.txt | 0 .../tools/scrapy-cluster-ctl.py | 0 .../cluster/tools/test-worker.py | 0 .../{pbcluster => cluster}/worker/__init__.py | 0 .../cluster/worker/manager.py | 0 .../scrapy/contrib/pbcluster/__init__.py | 3 - .../contrib/pbcluster/crawler/manager.py | 38 -- .../contrib/pbcluster/master/manager.py | 330 ------------------ .../scrapy/contrib/pbcluster/master/web.py | 239 ------------- .../contrib/pbcluster/worker/manager.py | 138 -------- .../contrib/pbcluster/worker/testworker.py | 25 -- .../scrapy/contrib_exp/cluster/__init__.py | 3 - .../contrib_exp/cluster/crawler/__init__.py | 0 .../contrib_exp/cluster/master/__init__.py | 0 .../contrib_exp/cluster/master/ws_api.txt | 51 --- .../cluster/tools/scrapy-cluster-ctl.py | 80 ----- .../contrib_exp/cluster/worker/__init__.py | 0 25 files changed, 5 insertions(+), 909 deletions(-) create mode 100644 scrapy/trunk/scrapy/contrib/cluster/__init__.py rename scrapy/trunk/scrapy/contrib/{pbcluster => cluster}/crawler/__init__.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/crawler/manager.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/hooks/__init__.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/hooks/svn.py (100%) rename scrapy/trunk/scrapy/contrib/{pbcluster => cluster}/master/__init__.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/master/manager.py (99%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/master/web.py (99%) rename scrapy/trunk/scrapy/contrib/{pbcluster => cluster}/master/ws_api.txt (100%) rename scrapy/trunk/scrapy/contrib/{pbcluster => cluster}/tools/scrapy-cluster-ctl.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/tools/test-worker.py (100%) rename scrapy/trunk/scrapy/contrib/{pbcluster => cluster}/worker/__init__.py (100%) rename scrapy/trunk/scrapy/{contrib_exp => contrib}/cluster/worker/manager.py (100%) delete mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/__init__.py delete mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py delete mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py delete mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/master/web.py delete mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py delete mode 100755 scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py delete mode 100644 scrapy/trunk/scrapy/contrib_exp/cluster/__init__.py delete mode 100644 scrapy/trunk/scrapy/contrib_exp/cluster/crawler/__init__.py delete mode 100644 scrapy/trunk/scrapy/contrib_exp/cluster/master/__init__.py delete mode 100644 scrapy/trunk/scrapy/contrib_exp/cluster/master/ws_api.txt delete mode 100755 scrapy/trunk/scrapy/contrib_exp/cluster/tools/scrapy-cluster-ctl.py delete mode 100644 scrapy/trunk/scrapy/contrib_exp/cluster/worker/__init__.py diff --git a/scrapy/trunk/scrapy/contrib/cluster/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/__init__.py new file mode 100644 index 000000000..2c6267ad8 --- /dev/null +++ b/scrapy/trunk/scrapy/contrib/cluster/__init__.py @@ -0,0 +1,3 @@ +from scrapy.contrib.cluster.worker.manager import ClusterWorker +from scrapy.contrib.cluster.master.web import ClusterMasterWeb +from scrapy.contrib.cluster.crawler.manager import ClusterCrawler diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/crawler/__init__.py similarity index 100% rename from scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py rename to scrapy/trunk/scrapy/contrib/cluster/crawler/__init__.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/crawler/manager.py b/scrapy/trunk/scrapy/contrib/cluster/crawler/manager.py similarity index 100% rename from scrapy/trunk/scrapy/contrib_exp/cluster/crawler/manager.py rename to scrapy/trunk/scrapy/contrib/cluster/crawler/manager.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/hooks/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/hooks/__init__.py similarity index 100% rename from scrapy/trunk/scrapy/contrib_exp/cluster/hooks/__init__.py rename to scrapy/trunk/scrapy/contrib/cluster/hooks/__init__.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/hooks/svn.py b/scrapy/trunk/scrapy/contrib/cluster/hooks/svn.py similarity index 100% rename from scrapy/trunk/scrapy/contrib_exp/cluster/hooks/svn.py rename to scrapy/trunk/scrapy/contrib/cluster/hooks/svn.py diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/master/__init__.py similarity index 100% rename from scrapy/trunk/scrapy/contrib/pbcluster/master/__init__.py rename to scrapy/trunk/scrapy/contrib/cluster/master/__init__.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/master/manager.py b/scrapy/trunk/scrapy/contrib/cluster/master/manager.py similarity index 99% rename from scrapy/trunk/scrapy/contrib_exp/cluster/master/manager.py rename to scrapy/trunk/scrapy/contrib/cluster/master/manager.py index 202c6fb5f..3d22a64d4 100644 --- a/scrapy/trunk/scrapy/contrib_exp/cluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/cluster/master/manager.py @@ -11,7 +11,7 @@ from scrapy.core import signals from scrapy import log from scrapy.core.engine import scrapyengine from scrapy.core.exceptions import NotConfigured -from scrapy.contrib_exp.cluster.worker.manager import ResponseCode +from scrapy.contrib.cluster.worker.manager import ResponseCode from scrapy.conf import settings def my_import(name): diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/master/web.py b/scrapy/trunk/scrapy/contrib/cluster/master/web.py similarity index 99% rename from scrapy/trunk/scrapy/contrib_exp/cluster/master/web.py rename to scrapy/trunk/scrapy/contrib/cluster/master/web.py index 96e4b41ca..f1c1e9a42 100644 --- a/scrapy/trunk/scrapy/contrib_exp/cluster/master/web.py +++ b/scrapy/trunk/scrapy/contrib/cluster/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_exp.cluster.master.manager import ClusterMaster +from scrapy.contrib.cluster.master.manager import ClusterMaster from scrapy.utils.serialization import serialize class ClusterMasterWeb(ClusterMaster): diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt b/scrapy/trunk/scrapy/contrib/cluster/master/ws_api.txt similarity index 100% rename from scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt rename to scrapy/trunk/scrapy/contrib/cluster/master/ws_api.txt diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/tools/scrapy-cluster-ctl.py b/scrapy/trunk/scrapy/contrib/cluster/tools/scrapy-cluster-ctl.py similarity index 100% rename from scrapy/trunk/scrapy/contrib/pbcluster/tools/scrapy-cluster-ctl.py rename to scrapy/trunk/scrapy/contrib/cluster/tools/scrapy-cluster-ctl.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/tools/test-worker.py b/scrapy/trunk/scrapy/contrib/cluster/tools/test-worker.py similarity index 100% rename from scrapy/trunk/scrapy/contrib_exp/cluster/tools/test-worker.py rename to scrapy/trunk/scrapy/contrib/cluster/tools/test-worker.py diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/__init__.py b/scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py similarity index 100% rename from scrapy/trunk/scrapy/contrib/pbcluster/worker/__init__.py rename to scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/cluster/worker/manager.py similarity index 100% rename from scrapy/trunk/scrapy/contrib_exp/cluster/worker/manager.py rename to scrapy/trunk/scrapy/contrib/cluster/worker/manager.py diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py b/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py deleted file mode 100644 index 84b6df381..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/__init__.py +++ /dev/null @@ -1,3 +0,0 @@ -from scrapy.contrib.pbcluster.worker.manager import ClusterWorker -from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb -from scrapy.contrib.pbcluster.crawler.manager import ClusterCrawler \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py deleted file mode 100644 index f76eebe6e..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py +++ /dev/null @@ -1,38 +0,0 @@ -import os - -from twisted.spread import pb -from twisted.internet import reactor - -from scrapy.conf import settings -from scrapy import log -from scrapy.core.manager import scrapymanager -from scrapy.core.exceptions import NotConfigured - -class Broker(pb.Referenceable): - def __init__(self, crawler, remote): - self.__remote = remote - self.__crawler = crawler - try: - deferred = self.__remote.callRemote("register_crawler", os.getpid(), self) - except pb.DeadReferenceError: - self._set_status(None) - log.msg("Lost connection to node %s." % (self.name), log.ERROR) - else: - deferred.addCallbacks(callback=lambda x: None, errback=lambda reason: log.msg(reason, log.ERROR)) - def remote_stop(self): - scrapymanager.stop() - -class ClusterCrawler: - def __init__(self): - if not settings.getbool('CLUSTER_CRAWLER_ENABLED'): - raise NotConfigured - - self.worker = None - - factory = pb.PBClientFactory() - reactor.connectTCP("localhost", settings.getint('CLUSTER_WORKER_PORT'), factory) - d = factory.getRootObject() - def _set_worker(obj): - self.worker = Broker(self, obj) - d.addCallbacks(callback=_set_worker, errback=lambda reason: log.msg(reason, log.ERROR)) - diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py deleted file mode 100644 index b7ecb4c1c..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ /dev/null @@ -1,330 +0,0 @@ -import sys, datetime -import pickle - -from pydispatch import dispatcher - -from twisted.spread import pb -from twisted.internet import reactor - -from scrapy.core import signals -from scrapy import log -from scrapy.core.engine import scrapyengine -from scrapy.core.exceptions import NotConfigured -from scrapy.conf import settings - -DEFAULT_PRIORITY = settings.getint("DEFAULT_PRIORITY", 20) - -def my_import(name): - mod = __import__(name) - components = name.split('.') - for comp in components[1:]: - mod = getattr(mod, comp) - return mod - -class Broker(pb.Referenceable): - def __init__(self, remote, name, master): - self.__remote = remote - self.alive = False - self.name = name - self.master = master - self.available = True - try: - deferred = self.__remote.callRemote("set_master", self) - 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 status_as_dict(self, verbosity=1): - if verbosity == 0: - return - status = {"alive": self.alive} - if self.alive: - if verbosity == 1: - #dont show spider settings - status["running"] = [] - for proc in self.running: - proccopy = proc.copy() - del proccopy["settings"] - status["running"].append(proccopy) - elif verbosity == 2: - status["running"] = self.running - status["maxproc"] = self.maxproc - status["freeslots"] = self.maxproc - len(self.running) - status["available"] = self.available - status["starttime"] = self.starttime - status["timestamp"] = self.timestamp - status["loadavg"] = self.loadavg - return status - - def _set_status(self, status): - if not status: - self.alive = False - else: - self.alive = True - self.running = status['running'] - self.maxproc = status['maxproc'] - self.starttime = status['starttime'] - self.timestamp = status['timestamp'] - self.loadavg = status['loadavg'] - self.logdir = status['logdir'] - free_slots = self.maxproc - len(self.running) - - #load domains by one, so to mix up better the domain loading between nodes. The next one in the same node will be loaded - #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 pending['domain'] in self.master.running or pending['domain'] in self.master.loading: - self.master.schedule([pending['domain']], pending['settings'], pending['priority']) - else: - self.run(pending) - self.master.loading.append(pending['domain']) - - def update_status(self): - try: - 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 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 run(self, pending): - - def _run_errback(reason): - log.msg(reason, log.ERROR) - self.master.loading.remove(pending['domain']) - self.master.schedule([pending['domain']], pending['settings'], pending['priority'] - 1) - log.msg("Domain %s rescheduled: lost connection to node." % pending['domain'], log.WARNING) - - def _run_callback(status): - if status['callresponse'][0] == 1: - #slots are complete. Reschedule in master with priority reduced by one. - #self.master.loading check should avoid this to happen - self.master.loading.remove(pending['domain']) - self.master.schedule([pending['domain']], pending['settings'], pending['priority'] - 1) - log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING) - elif status['callresponse'][0] == 2: - #domain already running in node. Reschedule with same priority. - #self.master.loading check should avoid this to happen - self.master.loading.remove(pending['domain']) - self.master.schedule([pending['domain']], pending['settings'], pending['priority']) - log.msg("Domain %s rescheduled: already running in node." % pending['domain'], log.WARNING) - - 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=_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): - def __init__(self, master, nodename): - pb.PBClientFactory.__init__(self) - self.master = master - self.nodename = nodename - - def clientConnectionLost(self, *args, **kargs): - pb.PBClientFactory.clientConnectionLost(self, *args, **kargs) - del self.master.nodes[self.nodename] - log.msg("Removed node %s." % self.nodename ) - -class ClusterMaster: - - def __init__(self): - - if not (settings.getbool('CLUSTER_MASTER_ENABLED')): - raise NotConfigured - - #import groups settings - if settings.getbool('GROUPSETTINGS_ENABLED'): - self.get_spider_groupsettings = my_import(settings["GROUPSETTINGS_MODULE"]).get_spider_groupsettings - else: - self.get_spider_groupsettings = lambda x: {} - #load pending domains - try: - self.pending = pickle.load( open(settings["CLUSTER_MASTER_CACHEFILE"], "r") ) - except IOError: - self.pending = [] - self.loading = [] - self.nodes = {} - self.start_time = datetime.datetime.utcnow() - #on how statistics works, see self.update_nodes() and Broker.remote_update() - self.statistics = {"domains": {"running": set(), "scraped": {}, "lost_count": {}, "lost": set()}, "scraped_count": 0 } - 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) - - def load_nodes(self): - """Loads nodes listed in 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): - """Creates the remote reference for each worker node""" - def _make_callback(_factory, _name, _url): - - def _errback(_reason): - log.msg("Could not get remote node %s in %s: %s." % (_name, _url, _reason), log.ERROR) - - d = _factory.getRootObject() - d.addCallbacks(callback=lambda obj: self.add_node(obj, _name), errback=_errback) - - server, port = url.split(":") - port = int(port) - log.msg("Connecting to cluster worker %s..." % name) - log.msg("Server: %s, Port: %s" % (server, port)) - factory = ScrapyPBClientFactory(self, name) - 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 update_nodes(self): - for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems(): - if name in self.nodes and self.nodes[name].alive: - log.msg("Updating node. name: %s, url: %s" % (name, url) ) - self.nodes[name].update_status() - else: - log.msg("Reloading node. name: %s, url: %s" % (name, url) ) - self.load_node(name, url) - - real_running = set(self.running.keys()) - lost = self.statistics["domains"]["running"].difference(real_running) - for domain in lost: - self.statistics["domains"]["lost_count"][domain] = self.statistics["domains"]["lost_count"].get(domain, 0) + 1 - self.statistics["domains"]["lost"] = self.statistics["domains"]["lost"].union(lost) - - def add_node(self, cworker, name): - """Add node given its node""" - node = Broker(cworker, name, self) - self.nodes[name] = node - log.msg("Added cluster worker %s" % name) - - def disable_node(self, name): - self.nodes[name].available = False - - 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): - i = 0 - for p in self.pending: - if p['priority'] <= priority: - i += 1 - else: - break - for domain in domains: - pd = self.find_inpending(domain) - if pd: #domain already pending, so just change priority if new is higher - if priority < pd['priority']: - self.pending.remove(pd) - pd['priority'] = priority - self.pending.insert(i, pd) - else: - 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}) - - def stop(self, domains): - to_stop = {} - for domain in domains: - node = self.running.get(domain, None) - if node: - if node.name not in to_stop: - to_stop[node.name] = [] - to_stop[node.name].append(domain) - - for nodename, domains in to_stop.iteritems(): - for domain in domains: - self.nodes[nodename].stop(domain) - - def remove(self, domains): - """Remove all scheduled instances of the given domains (if it hasn't - started yet). Otherwise use stop()""" - - 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 - domains""" - self.remove(domains) - self.stop(domains) - - @property - def running(self): - """Return dict of running domains as domain -> node""" - d = {} - for node in self.nodes.itervalues(): - for proc in node.running: - d[proc['domain']] = node - return d - - @property - def available_nodes(self): - return (node for node in self.nodes.itervalues() if node.available) - - def find_inpending(self, domain): - for p in self.pending: - if domain == p['domain']: - return p - - def print_pending(self, verbosity=1): - if verbosity == 1: - pending = [] - for p in self.pending: - pp = p.copy() - del pp["settings"] - pending.append(pp) - return pending - elif verbosity == 2: - return self.pending - return - - def _engine_started(self): - self.load_nodes() - scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL')) - def _engine_stopped(self): - pickle.dump( self.pending, open(settings["CLUSTER_MASTER_CACHEFILE"], "w") ) - log.msg("Pending saved in %s" % settings["CLUSTER_MASTER_CACHEFILE"]) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py deleted file mode 100644 index 9fe14632e..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py +++ /dev/null @@ -1,239 +0,0 @@ -import datetime - -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.utils.serialization import serialize - -class ClusterMasterWeb(ClusterMaster): - webconsole_id = 'cluster_master' - webconsole_name = 'Cluster master' - - def __init__(self): - ClusterMaster.__init__(self) - - dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module) - - def webconsole_render(self, wc_request): - changes = "" - if wc_request.path == '/cluster_master/nodes/': - return self.render_nodes(wc_request) - elif wc_request.path == '/cluster_master/domains/': - return self.render_domains(wc_request) - elif wc_request.path == '/cluster_master/ws/': - return self.webconsole_control(wc_request, ws=True) - elif wc_request.args: - changes = self.webconsole_control(wc_request) - - s = self.render_header() - - s += "

Home

\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, loadavg) - s += "
 NameAvailableRunningLoad.avg
%s%s%s%d/%d%s
\n" - - s += "\n" - s += "\n" - - return str(s) - - def webconsole_control(self, wc_request, ws=False): - args = wc_request.args - if "updatenodes" in args: - self.update_nodes() - if ws: - return self.ws_status(wc_request) - - if "schedule" in args: - if ws: - sep = "," - domains = args["schedule"][0].split(sep) - else: - sep = "\r" - domains = args["schedule"] - priority = int(args.get("priority", [DEFAULT_PRIORITY])[0]) - - #spider settings - slist = args.get("settings", [""])[0].split(sep) - spider_settings = {} - for s in slist: - try: - k, v = s.strip().split("=") - except ValueError: - pass - else: - spider_settings[k] = v - - self.schedule(domains, spider_settings, priority) - if ws: - return self.ws_status(wc_request, verbosity=0) - - if "stop" in args: - if ws: - domains = args["stop"][0].split(",") - else: - domains=args["stop"] - self.stop(domains) - if ws: - return self.ws_status(wc_request) - - if "remove" in args: - if ws: - domains = args["remove"][0].split(",") - else: - domains=args["remove"] - self.remove(domains) - if ws: - return self.ws_status(wc_request) - if "disable_node" in args: - self.disable_node(args["disable_node"][0]) - if ws: - return self.ws_status(wc_request) - if "enable_node" in args: - self.enable_node(args["enable_node"][0]) - if ws: - return self.ws_status(wc_request) - if "statistics" in args: - if ws: - return self.ws_statistics(wc_request) - - if ws: - return self.ws_status(wc_request) - else: - return "" - - def render_nodes(self, wc_request): - if wc_request.args: - self.webconsole_control(wc_request) - - now = datetime.datetime.utcnow() - - s = self.render_header() - for node in self.nodes.itervalues(): - if node.available: - s += "

%s

\n" % (node.name, node.name) - - s += "

Running domains

\n" - if node.running: - s += "
\n" - s += "\n" - s += "\n" - for proc in node.running: - chkbox = "" % proc['domain'] if proc['status'] == "running" else " " - start_time = proc.get('starttime', None) - elapsed = now - start_time if start_time else None - s += "\n" % \ - (chkbox, proc['pid'], proc['domain'], proc['status'], elapsed, proc['logfile']) - s += "
 PIDDomainStatusRunning timeLog file
%s%s%s%s%s%s
\n" - s += "\n" % node.name - s += "

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

No running domains on %s

\n" % node.name - - return str(s) - - def render_domains(self, wc_request): - if wc_request.args: - self.webconsole_control(wc_request) - - enabled_domains = set(spiders.asdict(include_disabled=False).keys()) - print "Enabled domains: %s" % len(enabled_domains) - inactive_domains = enabled_domains - set(self.running.keys() + [p['domain'] for p in self.pending]) - - s = self.render_header() - - s += "

Schedule domains

\n" - - s += "Inactive domains (not running or pending)
" - s += "
\n" - s += "\n" - s += "
\n" - - s += "Priority:
\n" - s += "%s" % DEFAULT_PRIORITY - s += "
\n" - - #spider settings - s += "Overrided spider settings:
\n" - s += "\n" - s += "
\n" - - s += "

\n" - s += "
\n" - - s += "

Domains

\n" - - s += "\n" - s += "\n" - s += self._domains_table(self.running, 'running') - s += "
DomainStatusNode
\n" - - # pending domains - s += "

Pending domains

\n" - if self.pending: - s += "
\n" - s += "\n" - s += "

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

No pending domains

\n" - - return str(s) - - def render_header(self): - s = banner(self) - s += "

Nav: " - s += "Home | " - s += "Domains | " - s += "Nodes (update)" - s += "

" - return s - - def _domains_table(self, dict_, status): - s = "" - for domain, node in dict_.iteritems(): - s += "%s%s%s\n" % (domain, status, node.name) - return s - - def webconsole_discover_module(self): - return self - - def ws_status(self, wc_request, verbosity=1): - format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json' - verbosity = int(wc_request.args['verbosity'][0]) if 'verbosity' in wc_request.args else verbosity - wc_request.setHeader('content-type', 'text/plain') - status = {} - nodes_status = {} - if verbosity > 0: - for d, n in self.nodes.iteritems(): - nodes_status[d] = n.status_as_dict(verbosity) - status["nodes"] = nodes_status - status["pending"] = self.print_pending(verbosity) - status["loading"] = self.loading - content = serialize(status, format) - return content - return "" - - def ws_statistics(self, wc_request): - format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json' - content = serialize(self.statistics, format) - return content \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py deleted file mode 100644 index c45d7b032..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ /dev/null @@ -1,138 +0,0 @@ -import sys, os, time, datetime, pickle, gzip - -from twisted.internet import protocol, reactor -from twisted.spread import pb - -from scrapy import log -from scrapy.core.exceptions import NotConfigured -from scrapy.conf import settings -from scrapy.core.engine import scrapyengine - -class ScrapyProcessProtocol(protocol.ProcessProtocol): - def __init__(self, procman, domain, logfile=None, spider_settings=None): - self.procman = procman - self.domain = domain - self.logfile = logfile - self.start_time = datetime.datetime.utcnow() - self.status = "starting" - self.pid = -1 - self.env = {} - #We conserve original setting format for info purposes (avoid lots of unnecesary "SCRAPY_") - self.scrapy_settings = spider_settings or {} - self.scrapy_settings.update({'LOGFILE': self.logfile, 'CLUSTER_WORKER_ENABLED': 0, 'CLUSTER_CRAWLER_ENABLED': 1, 'WEBCONSOLE_ENABLED': 0}) - pickled_settings = pickle.dumps(self.scrapy_settings) - self.env["SCRAPY_PICKLED_SETTINGS_TO_OVERRIDE"] = pickled_settings - self.env["PYTHONPATH"] = ":".join(sys.path)#this is need so this crawl process knows where to locate local_scrapy_settings. - - def __str__(self): - return "" % (self.domain, self.pid, self.status) - - def as_dict(self): - return {"domain": self.domain, "pid": self.pid, "status": self.status, "settings": self.scrapy_settings, "logfile": self.logfile, "starttime": self.start_time} - - def connectionMade(self): - self.pid = self.transport.pid - log.msg("ClusterWorker: started domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) - self.transport.closeStdin() - self.status = "running" - self.procman.update_master(self.domain, "running") - - def processEnded(self, reason): - if settings.getbool('CLUSTER_WORKER_GZIP_LOGS'): - try: - f_in = open(self.logfile) - f_out = gzip.open("%s.gz" % self.logfile, "wb") - f_out.writelines(f_in) - f_in.close() - f_out.close() - os.remove(self.logfile) - self.logfile = "%s.gz" % self.logfile - except Exception, e: - log.msg("failed to compress %s exception=%s (domain=%s, pid=%s)" % (self.logfile, e, self.domain, self.pid)) - log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile)) - log.msg("Reason type: %s. value: %s" % (reason.type, reason.value) ) - del self.procman.running[self.domain] - del self.procman.crawlers[self.pid] - self.procman.update_master(self.domain, "scraped") - -class ClusterWorker(pb.Root): - - def __init__(self): - if not settings.getbool('CLUSTER_WORKER_ENABLED'): - raise NotConfigured - - self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') - self.logdir = settings['CLUSTER_LOGDIR'] - self.running = {}#a dict domain->ScrapyProcessControl - self.crawlers = {}#a dict pid->scrapy process remote pb connection - self.starttime = datetime.datetime.utcnow() - port = settings.getint('CLUSTER_WORKER_PORT') - scrapyengine.listenTCP(port, pb.PBServerFactory(self)) - log.msg("PYTHONPATH: %s" % repr(sys.path)) - - def status(self, rcode=0, rstring=None): - status = {} - status["running"] = [ self.running[k].as_dict() for k in self.running.keys() ] - status["starttime"] = self.starttime - status["timestamp"] = datetime.datetime.utcnow() - status["maxproc"] = self.maxproc - status["loadavg"] = os.getloadavg() - status["logdir"] = self.logdir - status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.") - return status - - def update_master(self, domain, domain_status): - try: - deferred = self.__master.callRemote("update", self.status(), domain, domain_status) - except pb.DeadReferenceError: - self.__master = None - log.msg("Lost connection to node %s." % (self.name), log.ERROR) - else: - deferred.addCallbacks(callback=lambda x: x, errback=lambda reason: log.msg(reason, log.ERROR)) - - def remote_set_master(self, master): - self.__master = master - return self.status() - - 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)) - d = self.crawlers[proc.pid].callRemote("stop") - def _close(): - proc.status = "closing" - d.addCallbacks(callback=_close, errback=lambda reason: log.msg(reason, log.ERROR)) - 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 remote_run(self, domain, spider_settings=None): - """Spawn process to run the given domain.""" - if len(self.running) < self.maxproc: - if not domain in self.running: - 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] - self.running[domain] = scrapy_proc - try: - import pysvn - c = pysvn.Client() - r = c.update(settings.get("CLUSTER_WORKER_SVNWORKDIR", ".")) - log.msg("Updated to revision %s." %r[0].number, level=log.DEBUG) - except pysvn.ClientError, e: - log.msg("Unable to svn update: %s" % e, level=log.WARNING) - except ImportError: - log.msg("pysvn module not available.", level=log.WARNING) - proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) - return self.status(0, "Started process %s." % scrapy_proc) - return self.status(2, "Domain %s already running." % domain ) - return self.status(1, "No free slot to run another process.") - - def remote_register_crawler(self, pid, crawler): - self.crawlers[pid] = crawler diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py deleted file mode 100755 index eee998910..000000000 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/testworker.py +++ /dev/null @@ -1,25 +0,0 @@ -#!/usr/bin/python2.5 - -from twisted.spread import pb -from twisted.internet import reactor -from twisted.python import util -import sys - -factory = pb.PBClientFactory() -reactor.connectTCP("localhost", 8789, factory) -d = factory.getRootObject() - -sys.argv.pop(0) - -if not sys.argv: - d.addCallback(lambda object: object.callRemote("status")) -elif sys.argv[0] == "-s": - d.addCallback(lambda object: object.callRemote("stop", sys.argv[1])) -elif sys.argv[0] == "-r": - d.addCallback(lambda object: object.callRemote("run", sys.argv[1])) -elif sys.argv[0] == "-t": - d.addCallback(lambda object: object.callRemote("statistics")) - -d.addCallbacks(callback = util.println, errback = lambda reason: 'error: '+str(reason.value)) -d.addCallback(lambda _: reactor.stop()) -reactor.run() diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/__init__.py b/scrapy/trunk/scrapy/contrib_exp/cluster/__init__.py deleted file mode 100644 index 5917df97c..000000000 --- a/scrapy/trunk/scrapy/contrib_exp/cluster/__init__.py +++ /dev/null @@ -1,3 +0,0 @@ -from scrapy.contrib_exp.cluster.worker.manager import ClusterWorker -from scrapy.contrib_exp.cluster.master.web import ClusterMasterWeb -from scrapy.contrib_exp.cluster.crawler.manager import ClusterCrawler diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/crawler/__init__.py b/scrapy/trunk/scrapy/contrib_exp/cluster/crawler/__init__.py deleted file mode 100644 index e69de29bb..000000000 diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/master/__init__.py b/scrapy/trunk/scrapy/contrib_exp/cluster/master/__init__.py deleted file mode 100644 index e69de29bb..000000000 diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/master/ws_api.txt b/scrapy/trunk/scrapy/contrib_exp/cluster/master/ws_api.txt deleted file mode 100644 index c47626a08..000000000 --- a/scrapy/trunk/scrapy/contrib_exp/cluster/master/ws_api.txt +++ /dev/null @@ -1,51 +0,0 @@ -Cluster Webservice API -====================== - -The webservice API is available at - -http://server:port/cluster_master/ws/ - -With no parameters, webservice returns the cluster status. - -Query parameters -================ - - - `format`: the answer format. By default, format=json. Other formats: pprint, pickle. - - - `schedule`: schedules a comma separated list of domains. Schedule function takes optional parameters: - - "priority": sets the queue priority for the specified domains (an integer). The default is setted by "DEFAULT_PRIORITY" - setting (20 if not given). A lower priority number implies more priority. - - "settings": run settings for the specified domains. This is a comma separated list of = pairs. By default it is empty. - - - `remove`: removes from pending list a comma separated list of domains. - - - `stop`: stops comma separated list of domains (they have to be running in some node) - - - `disable_node`: disables a node so no more domains will be loaded in it until enabled again (but it will finish to run the running domains) - - - `enable_node`: revert the state setted by 'disable_node' - - - `verbosity`: sets the output verbosity level (1 is the default minimal, 2 includes domain settings, 0 disables output) - - - `statistics`: shows the pending/running/scraped/lost statistics - -Examples: ---------- - -1) Schedule argos.co.uk, diy.com, littlewoodsdirect.com spiders, with priority=0, and settings UNAVAILABLES_NOTIFY=2 and UNAVAILABLES_DAYS_BACK=3. Answer with pprint format - - http://localhost:8080/cluster_master/ws/?format=pprint&schedule=argos.co.uk,diy.com,littlewoodsdirect.com&priority=0&settings=UNAVAILABLES_NOTIFY=2,UNAVAILABLES_DAYS_BACK=3 - -2) Get status with pprint format: - - http://localhost:8080/cluster_master/ws/?format=pprint - -3) Remove from pending lists domains argos.co.uk and diy.com. Answer with pprint format: - - http://localhost:8080/cluster_master/ws/?remove=argos.co.uk,diy.com - -4) Stop running domain littlewoodsdirect.com: - - http://localhost:8080/cluster_master/ws/?stop=littlewoodsdirect.com \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/tools/scrapy-cluster-ctl.py b/scrapy/trunk/scrapy/contrib_exp/cluster/tools/scrapy-cluster-ctl.py deleted file mode 100755 index f3ca78055..000000000 --- a/scrapy/trunk/scrapy/contrib_exp/cluster/tools/scrapy-cluster-ctl.py +++ /dev/null @@ -1,80 +0,0 @@ -#!/usr/bin/env python -""" -Cluster control script -""" - -from optparse import OptionParser -import urllib - -def main(): - parser = OptionParser(usage="Usage: scrapy-cluster-ctl.py [domain [domain [...]]] [options]" ) - parser.add_option("--disablenode", dest="disable_node", help="Disable given node (by name) so it will no accept more run requests.") - parser.add_option("--enablenode", dest="enable_node", help="Enable given node (by name) so it will accept again run requests") - parser.add_option("--format", dest="format", help="Output format. Default: pprint.", default="pprint") - parser.add_option("--list", metavar="FILE", dest="list", help="Specify a file from where to read domains, one per line.") - parser.add_option("--now", action="store_true", dest="now", help="Schedule domains to run with priority now.") - parser.add_option("--output", metavar="FILE", dest="output", help="Output file. If not given, output to stdout.") - parser.add_option("--port", dest="port", type="int", help="Cluster master port. Default: 8060.", default=8060) - parser.add_option("--remove", dest="remove", action="store_true", help="Remove from schedule domains given as args.") - parser.add_option("--schedule", dest="schedule", action="store_true", help="Schedule domains given as args.") - parser.add_option("--server", dest="server", help="Cluster master server name. Default: localhost.", default="localhost") - parser.add_option("--status", dest="status", action="store_true", help="Print cluster master status and quit.") - parser.add_option("--statistics", dest="statistics", action="store_true", help="Print cluster statistics") - parser.add_option("--stop", dest="stop", action="store_true", help="Stops a running domain.") - parser.add_option("--verbosity", dest="verbosity", type="int", help="Sets the report status verbosity.") - (opts, args) = parser.parse_args() - - output = "" - domains = [] - urlstring = "http://%s:%s/cluster_master/ws/" % (opts.server, opts.port) - post = {"format":opts.format} - if isinstance(opts.verbosity, int): - post["verbosity"] = opts.verbosity - - if args: - domains = ",".join(args) - elif opts.list: - try: - domainlist = [] - for d in open(opts.list, "r").readlines(): - domainlist.append(d.strip()) - domains = ",".join(domainlist) - except IOError: - print "Can't open file %s" % opts.list - - if opts.status: - pass - elif opts.statistics: - post["statistics"] = True - elif opts.schedule and domains: - post["schedule"] = domains - if opts.now: - post["priority"] = "0" - post["settings"] = "UNAVAILABLES_NOTIFY=2" - elif opts.remove and domains: - post["remove"] = domains - elif opts.stop and domains: - post["stop"] = domains - elif opts.disable_node: - post["disable_node"] = opts.disable_node - elif opts.enable_node: - post["enable_node"] = opts.enable_node - else: - parser.print_help() - return - - f = urllib.urlopen(urlstring, urllib.urlencode(post)) - output=f.read() - if not output: - return - if not opts.output: - print output - else: - try: - open(opts.output, "w").write(output) - except IOError: - open("/tmp/scrapy-cluster-schedule.tmp", "w").write(output) - print "Could not open file %s for writing. Output dumped to /tmp/scrapy-cluster-schedule.tmp instead." % opts.output - -if __name__ == '__main__': - main() diff --git a/scrapy/trunk/scrapy/contrib_exp/cluster/worker/__init__.py b/scrapy/trunk/scrapy/contrib_exp/cluster/worker/__init__.py deleted file mode 100644 index e69de29bb..000000000