mirror of https://github.com/scrapy/scrapy.git
Cluster crawler fixes
--HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40155
This commit is contained in:
parent
5b4d9f8f85
commit
290702d988
|
|
@ -1,2 +1,3 @@
|
|||
from scrapy.contrib.pbcluster.worker.manager import ClusterWorker
|
||||
from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb
|
||||
from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb
|
||||
from scrapy.contrib.pbcluster.crawler.manager import ClusterCrawler
|
||||
|
|
@ -1,21 +1,24 @@
|
|||
import os
|
||||
|
||||
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
|
||||
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", domain, self)
|
||||
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=self._set_status, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
deferred.addCallbacks(callback=lambda x: None, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
def remote_stop(self):
|
||||
scrapymanager.stop()
|
||||
|
||||
|
|
@ -29,5 +32,7 @@ class ClusterCrawler:
|
|||
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))
|
||||
def _set_worker(obj):
|
||||
self.worker = Broker(self, obj)
|
||||
d.addCallbacks(callback=_set_worker, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
|
||||
|
|
@ -41,7 +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]
|
||||
del self.procman.crawlers[self.pid]
|
||||
self.procman.update_master(self.domain, "scraped")
|
||||
|
||||
class ClusterWorker(pb.Root):
|
||||
|
|
@ -53,7 +53,7 @@ class ClusterWorker(pb.Root):
|
|||
self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC')
|
||||
self.logdir = settings['CLUSTER_LOGDIR']
|
||||
self.running = {}#a dict domain->ScrapyProcessControl
|
||||
self.crawlers = {}#a dict domain->scrapy process remote pb connection
|
||||
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))
|
||||
|
|
@ -88,7 +88,7 @@ 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))
|
||||
d = self.crawler["domain"].callRemote("stop")
|
||||
d = self.crawlers[proc.pid].callRemote("stop")
|
||||
def _close():
|
||||
proc.status = "closing"
|
||||
d.addCallbacks(callback=_close, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
|
|
@ -123,5 +123,5 @@ class ClusterWorker(pb.Root):
|
|||
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
|
||||
def remote_register_crawler(self, pid, crawler):
|
||||
self.crawlers[pid] = crawler
|
||||
Loading…
Reference in New Issue