mirror of https://github.com/scrapy/scrapy.git
- added svn update support
- removed passing of env variables - added automatic group settings load --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%4040
This commit is contained in:
parent
8344819a92
commit
23b3408403
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -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 += "<br />\n"
|
||||
|
||||
#spider settings
|
||||
s += "Spider settings:<br />\n"
|
||||
s += "<textarea name='settings' rows='6'>\n"
|
||||
s += "Overrided spider settings:<br />\n"
|
||||
s += "<textarea name='settings' rows='4'>\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"
|
||||
|
||||
|
|
|
|||
|
|
@ -18,8 +18,6 @@ Query parameters:
|
|||
|
||||
'settings': run settings for the specified domains. This is a comma separated list of <setting_name>=<value> pairs. By default it is empty.
|
||||
|
||||
'env': additional environment settings passed to the crawl process. This is a comma separated list of <setting_name>=<value> 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:
|
||||
|
||||
|
|
|
|||
|
|
@ -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 "<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.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.")
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue