From 29cd1bc3cb2db46e0a599af0b2de2daf06402b3e Mon Sep 17 00:00:00 2001 From: olveyra Date: Fri, 8 Aug 2008 16:57:33 +0000 Subject: [PATCH] Added pb-capable crawler. The idea is to improve cluster performance adding communication between crawler and master. At the momento, a remote stop method to the crawler was added to sustitute the previous stop based on kernel signal. Further will add monitoring functionality, because the processes are very silent, mainly when unavailable report is not issued, and offen happens lots of thing that nobody realize on if some fortuite events wouldn't happent --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40153 --- .../contrib/pbcluster/crawler/manager.py | 33 +++++++++++++++++++ .../contrib/pbcluster/master/manager.py | 6 ++-- .../contrib/pbcluster/worker/manager.py | 13 +++++--- 3 files changed, 45 insertions(+), 7 deletions(-) create mode 100644 scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py new file mode 100644 index 000000000..8d66b9799 --- /dev/null +++ b/scrapy/trunk/scrapy/contrib/pbcluster/crawler/manager.py @@ -0,0 +1,33 @@ +from twisted.spread import pb +from twisted.internet import reactor + +from scrapy.conf import settings +from scrapy.core import log +from scrapy.core.manager import scrapymanager + +class Broker(pb.Referenceable): + def __init__(self, crawler, remote): + self.__remote = remote + self.__crawler = crawler + try: + deferred = self.__remote.callRemote("register_crawler", domain, 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 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() + d.addCallbacks(callback=lambda obj: self.worker=Node(self, obj), errback=lambda reason: log.msg(reason, log.ERROR)) + \ No newline at end of file diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index 20737275c..a5aa97fe0 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -27,7 +27,7 @@ def my_import(name): mod = getattr(mod, comp) return mod -class Node(pb.Referenceable): +class Broker(pb.Referenceable): def __init__(self, remote, name, master): self.__remote = remote self.alive = False @@ -177,7 +177,7 @@ class ClusterMaster: self.loading = [] self.nodes = {} self.start_time = datetime.datetime.utcnow() - #on how statistics works, see self.update_nodes() and Nodes.remote_update() + #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 @@ -231,7 +231,7 @@ class ClusterMaster: def add_node(self, cworker, name): """Add node given its node""" - node = Node(cworker, name, self) + node = Broker(cworker, name, self) self.nodes[name] = node log.msg("Added cluster worker %s" % name) diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py index 981e1263e..c3d2f7502 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -19,7 +19,7 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): 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', 'WEBCONSOLE_ENABLED': '0'}) + 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. @@ -41,6 +41,7 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol): 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.domain] self.procman.update_master(self.domain, "scraped") class ClusterWorker(pb.Root): @@ -51,7 +52,8 @@ class ClusterWorker(pb.Root): self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC') self.logdir = settings['CLUSTER_LOGDIR'] - self.running = {} + self.running = {}#a dict domain->ScrapyProcessControl + self.crawlers = {}#a dict domain->scrapy process remote pb connection self.starttime = datetime.datetime.utcnow() port = settings.getint('CLUSTER_WORKER_PORT') scrapyengine.listenTCP(port, pb.PBServerFactory(self)) @@ -86,8 +88,8 @@ class ClusterWorker(pb.Root): 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" + d = self.crawler["domain"].callRemote("stop") + d.addCallbacks(callback=lambda x: proc.status="closing", 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) @@ -118,3 +120,6 @@ class ClusterWorker(pb.Root): 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, domain, crawler): + self.crawlers['domain'] = crawler \ No newline at end of file