diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py index b951209a8..92b5b7ced 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/manager.py @@ -24,6 +24,13 @@ priorities = { 20:'NORMAL', for val, attr in priorities.items(): setattr(sys.modules[__name__], "PRIORITY_%s" % attr, val ) +def my_import(name): + mod = __import__(name) + components = name.split('.') + for comp in components[1:]: + mod = getattr(mod, comp) + return mod + class Node: def __init__(self, remote, status, name, master): self.__remote = remote @@ -72,12 +79,12 @@ class Node: def _run_callback(status): if status['callresponse'][0] == 1: #slots are complete. Reschedule in master. This is a security issue because could happen that the slots were completed since last status update by another cluster (thinking at future with full-distributed worker-master clusters) - self.master.schedule(pending['domain'], pending['settings'], pending['priority'], pending["env"]) + self.master.schedule(pending['domain'], pending['settings'], pending['priority']) log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING) self._set_status(status) try: - deferred = self.__remote.callRemote("run", pending["domain"], pending["settings"], pending["env"]) + 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) @@ -87,13 +94,22 @@ class Node: class ClusterMaster(object): def __init__(self): - - if not settings.getbool('CLUSTER_MASTER_ENABLED'): + + 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("pending_cache_%s" % socket.gethostname(), "r") ) except IOError: self.pending = [] + self.nodes = {} dispatcher.connect(self._engine_started, signal=signals.engine_started) dispatcher.connect(self._engine_stopped, signal=signals.engine_stopped) @@ -138,7 +154,7 @@ class ClusterMaster(object): def remove_node(self, nodename): raise NotImplemented - def schedule(self, domains, spider_settings=None, env=None, priority=PRIORITY_NORMAL): + def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL): i = 0 for p in self.pending: if p['priority'] <= priority: @@ -146,7 +162,9 @@ class ClusterMaster(object): else: break for domain in domains: - self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'env': env, 'priority': priority}) + final_spider_settings = self.get_spider_groupsettings(domain) + final_spider_settings.update(spider_settings or {}) + self.pending.insert(i, {'domain': domain, 'settings': final_spider_settings, 'priority': priority}) self.update_nodes() def stop(self, domains): diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py index b99469520..e61d36b4c 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/web.py @@ -73,19 +73,8 @@ class ClusterMasterWeb(ClusterMaster): pass else: spider_settings[k] = v - - #other environment settings - envlist = args.get("env", [""])[0].split(sep) - env = {} - for e in envlist: - try: - k, v = e.strip().split("=") - except ValueError: - pass - else: - env[k] = v - - self.schedule(domains, spider_settings, env, priority) + + self.schedule(domains, spider_settings, priority) if ws: return self.ws_status(wc_request) @@ -173,18 +162,12 @@ class ClusterMasterWeb(ClusterMaster): s += "
\n" #spider settings - s += "Spider settings:
\n" - s += "\n" s += "
\n" - #other environment settings - s += "Other environment settings:
\n" - s += "\n" - s += "

\n" s += "\n" diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt b/scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt index 525525812..dc8042b63 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt +++ b/scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt @@ -18,8 +18,6 @@ Query parameters: 'settings': run settings for the specified domains. This is a comma separated list of = pairs. By default it is empty. - 'env': additional environment settings passed to the crawl process. 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) @@ -28,7 +26,7 @@ Examples: 1) Schedule argos.co.uk, diy.com, littlewoodsdirect.com spiders, with priority=PRIORITY_NOW, 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=PRIORITY_NOW&settings=UNAVAILABLES_NOTIFY=2,UNAVAILABLES_DAYS_BACK=3&env=PYTHONPATH=/opt/scraping/conf/feed + http://localhost:8080/cluster_master/ws/?format=pprint&schedule=argos.co.uk,diy.com,littlewoodsdirect.com&priority=PRIORITY_NOW&settings=UNAVAILABLES_NOTIFY=2,UNAVAILABLES_DAYS_BACK=3 2) Get status with pprint format: diff --git a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py index 1472832ed..46b6bbf94 100644 --- a/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py +++ b/scrapy/trunk/scrapy/contrib/pbcluster/worker/manager.py @@ -12,26 +12,25 @@ 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, env=None): + 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_original = env or {} - self.env = self.env_original.copy() + 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'}) for k in self.scrapy_settings: - self.env["SCRAPY_%s" % k] = self.scrapy_settings[k] + self.env["SCRAPY_%s" % k] = str(self.scrapy_settings[k]) 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, "env": self.env_original} + 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 @@ -81,16 +80,22 @@ class ClusterWorker(pb.Root): status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.") return status - def remote_run(self, domain, spider_settings=None, env=None): + def remote_run(self, domain, spider_settings=None): """Spawn process to run the given domain.""" if len(self.running) < self.maxproc: + try: + import pysvn + c=pysvn.Client() + r = c.update(settings["SVN_WORKDIR"] or ".") + log.msg("Updated to revision %s." %r[0].number ) + except: + pass 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, env) + scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings) args = [sys.executable, sys.argv[0], 'crawl', domain] proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env) self.running[domain] = scrapy_proc return self.status(0, "Started process %s." % scrapy_proc) return self.status(1, "No free slot to run another process.") -