added environment options

--HG--
extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%4037
This commit is contained in:
olveyra 2008-07-01 15:17:00 +00:00
parent d2684dc17b
commit 74ef3ab5b0
3 changed files with 38 additions and 20 deletions

View File

@ -71,12 +71,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'])
self.master.schedule(pending['domain'], pending['settings'], pending['priority'], pending["env"])
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"])
deferred = self.__remote.callRemote("run", pending["domain"], pending["settings"], pending["env"])
except pb.DeadReferenceError:
self._set_status(None)
log.msg("Lost connection to node %s." % (self.name), log.ERROR)
@ -132,7 +132,7 @@ class ClusterMaster(object):
def remove_node(self, nodename):
raise NotImplemented
def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL):
def schedule(self, domains, spider_settings=None, env=None, priority=PRIORITY_NORMAL):
i = 0
for p in self.pending:
if p['priority'] <= priority:
@ -140,7 +140,7 @@ class ClusterMaster(object):
else:
break
for domain in domains:
self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority})
self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'env': env, 'priority': priority})
self.update_nodes()
def stop(self, domains):

View File

@ -62,6 +62,8 @@ class ClusterMasterWeb(ClusterMaster):
sep = "\r"
domains = args["schedule"]
priority = eval(args.get("priority", ["PRIORITY_NORMAL"])[0])
#spider settings
slist = args.get("settings", [""])[0].split(sep)
spider_settings = {}
for s in slist:
@ -71,7 +73,19 @@ class ClusterMasterWeb(ClusterMaster):
pass
else:
spider_settings[k] = v
self.schedule(domains, spider_settings, priority)
#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)
if ws:
return self.ws_status(wc_request)
@ -157,11 +171,20 @@ class ClusterMasterWeb(ClusterMaster):
s += "<option value='%s'>%s</option>" % (p, pname)
s += "</select>\n"
s += "<br />\n"
#spider settings
s += "Spider settings:<br />\n"
s += "<textarea name='settings' rows='6'>\n"
s += "UNAVAILABLES_NOTIFY=2\n"
s += "UNAVAILABLES_DAYS_BACK=3\n"
s += "</textarea>\n"
s += "<br />\n"
#other environment settings
s += "Other environment settings:<br />\n"
s += "<textarea name='env' rows='3'>\n"
s += "</textarea>\n"
s += "<p><input type='submit' value='Schedule selected domains'></p>\n"
s += "</form>\n"
@ -210,12 +233,7 @@ class ClusterMasterWeb(ClusterMaster):
wc_request.setHeader('content-type', 'text/plain')
status = {}
nodes_status = {}
running = []
for d, n in self.nodes.iteritems():
nodes_status[d] = n.status_as_dict
running.extend(n.running)
status["nodes"] = nodes_status
status["pending"] = self.pending
status["running"] = running
content = serialize(status, format)
return content

View File

@ -12,26 +12,26 @@ 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):
def __init__(self, procman, domain, logfile=None, spider_settings=None, env=None):
self.procman = procman
self.domain = domain
self.logfile = logfile
self.start_time = datetime.datetime.utcnow()
self.status = "starting"
self.pid = -1
env = {'SCRAPY_LOGFILE': self.logfile, 'SCRAPY_CLUSTER_WORKER_ENABLED': '0', 'SCRAPY_WEBCONSOLE_ENABLED': '0'}
self.env_original = env or {}
self.env = self.env_original.copy()
#We conserve original setting format for info purposes (avoid lots of unnecesary "SCRAPY_")
self.settings = spider_settings or {}
for k in self.settings:
env["SCRAPY_%s" % k] = self.settings[k]
self.env = env
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]
def __str__(self):
return "<ScrapyProcess domain=%s, pid=%s, status=%s>" % (self.domain, self.pid, self.status)
def as_dict(self):
return {"domain": self.domain, "pid": self.pid, "status": self.status, "settings": self.settings, "logfile": self.logfile, "starttime": self.start_time}
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}
def connectionMade(self):
self.pid = self.transport.pid
@ -81,13 +81,13 @@ class ClusterWorker(pb.Root):
status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.")
return status
def remote_run(self, domain, spider_settings=None):
def remote_run(self, domain, spider_settings=None, env=None):
"""Spawn process to run the given domain."""
if len(self.running) < self.maxproc:
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)
scrapy_proc = ScrapyProcessProtocol(self, domain, logfile, spider_settings, env)
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