mirror of https://github.com/scrapy/scrapy.git
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
This commit is contained in:
parent
6cfbe78d63
commit
29cd1bc3cb
|
|
@ -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))
|
||||
|
||||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Reference in New Issue