mirror of https://github.com/scrapy/scrapy.git
pbcluster commit
--HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%4031
This commit is contained in:
parent
4caadf6b67
commit
53838d6a7d
|
|
@ -24,45 +24,67 @@ for val, attr in priorities.items():
|
|||
setattr(sys.modules[__name__], "PRIORITY_%s" % attr, val )
|
||||
|
||||
class Node:
|
||||
def __init__(self, remote, status, name):
|
||||
def __init__(self, remote, status, name, master):
|
||||
self.__remote = remote
|
||||
self._set_status(status)
|
||||
self.name = name
|
||||
self.master = master
|
||||
|
||||
def _set_status(self, status):
|
||||
if not status:
|
||||
self.available = False
|
||||
else:
|
||||
print ">----------<"
|
||||
self.available = True
|
||||
self.running = status['running']
|
||||
self.pending = status['pending']
|
||||
self.closing = status['closing']
|
||||
self.maxproc = status['maxproc']
|
||||
self.starttime = status['starttime']
|
||||
self.timestamp = status['timestamp']
|
||||
self.loadavg = status['loadavg']
|
||||
self.logdir = status['logdir']
|
||||
self.lastcallresponse = status['callresponse']
|
||||
print "Running: %s" % self.running
|
||||
free_slots = self.maxproc - len(self.running)
|
||||
while free_slots > 0 and self.master.pending:
|
||||
print "Free slots %s" % free_slots
|
||||
pending = self.master.pending.pop(0)
|
||||
self.run(pending)
|
||||
free_slots -= 1
|
||||
|
||||
def _remote_call(self, function, *args):
|
||||
def get_status(self):
|
||||
try:
|
||||
deferred = self.__remote.callRemote(function, *args)
|
||||
deferred = self.__remote.callRemote("status")
|
||||
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 get_status(self):
|
||||
self._remote_call("status")
|
||||
def stop(self, domain):
|
||||
try:
|
||||
deferred = self.__remote.callRemote("stop", domain)
|
||||
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 schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL):
|
||||
self._remote_call("schedule", domains, spider_settings, priority)
|
||||
def run(self, pending):
|
||||
|
||||
def stop(self, domains):
|
||||
self._remote_call("stop", domains)
|
||||
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'])
|
||||
log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING)
|
||||
self._set_status(status)
|
||||
|
||||
def remove(self, domains):
|
||||
self._remote_call("remove", domains)
|
||||
try:
|
||||
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)
|
||||
else:
|
||||
deferred.addCallbacks(callback=_run_callback, errback=lambda reason:log.msg(reason, log.ERROR))
|
||||
|
||||
class ClusterMaster(object):
|
||||
|
||||
|
|
@ -70,7 +92,7 @@ class ClusterMaster(object):
|
|||
if not settings.getbool('CLUSTER_MASTER_ENABLED'):
|
||||
raise NotConfigured
|
||||
self.nodes = {}
|
||||
self.queue = []
|
||||
self.pending = []
|
||||
dispatcher.connect(self._engine_started, signal=signals.engine_started)
|
||||
|
||||
def load_nodes(self):
|
||||
|
|
@ -84,7 +106,7 @@ class ClusterMaster(object):
|
|||
d.addCallbacks(callback=lambda obj: self.add_node(obj, _name), errback=_errback)
|
||||
|
||||
"""Loads nodes from the CLUSTER_MASTER_NODES setting"""
|
||||
|
||||
|
||||
for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
|
||||
if name not in self.nodes:
|
||||
server, port = url.split(":")
|
||||
|
|
@ -98,14 +120,14 @@ class ClusterMaster(object):
|
|||
log.msg("Could not connect to node %s in %s: %s." % (name, url, reason), log.ERROR)
|
||||
else:
|
||||
_make_callback(factory, name, url)
|
||||
|
||||
|
||||
def update_nodes(self):
|
||||
for node in self.nodes.itervalues():
|
||||
node.get_status()
|
||||
|
||||
def add_node(self, cworker, name):
|
||||
"""Add node given its node"""
|
||||
node = Node(cworker, None, name)
|
||||
node = Node(cworker, None, name, self)
|
||||
node.get_status()
|
||||
self.nodes[name] = node
|
||||
log.msg("Added cluster worker %s" % name)
|
||||
|
|
@ -113,11 +135,16 @@ class ClusterMaster(object):
|
|||
def remove_node(self, nodename):
|
||||
raise NotImplemented
|
||||
|
||||
def schedule(self, domains, spider_settings=None, nodename=None, priority=PRIORITY_NORMAL):
|
||||
if nodename:
|
||||
self.nodes[nodename].schedule(domains, spider_settings, priority)
|
||||
else:
|
||||
self._dispatch_domains(domains, spider_settings, priority)
|
||||
def schedule(self, domains, spider_settings=None, priority=PRIORITY_NORMAL):
|
||||
i = 0
|
||||
for p in self.pending:
|
||||
if p['priority'] <= priority:
|
||||
i += 1
|
||||
else:
|
||||
break
|
||||
for domain in domains:
|
||||
self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority})
|
||||
self.update_nodes()
|
||||
|
||||
def stop(self, domains):
|
||||
to_stop = {}
|
||||
|
|
@ -129,19 +156,21 @@ class ClusterMaster(object):
|
|||
to_stop[node.name].append(domain)
|
||||
|
||||
for nodename, domains in to_stop.iteritems():
|
||||
self.nodes[nodename].stop(domains)
|
||||
for domain in domains:
|
||||
self.nodes[nodename].stop(domain)
|
||||
|
||||
def remove(self, domains):
|
||||
to_remove = {}
|
||||
for domain in domains:
|
||||
node = self.pending.get(domain, None)
|
||||
if node:
|
||||
if node.name not in to_remove:
|
||||
to_remove[node.name] = []
|
||||
to_remove[node.name].append(domain)
|
||||
"""Remove all scheduled instances of the given domains (if it hasn't
|
||||
started yet). Otherwise use stop()"""
|
||||
|
||||
for nodename, domains in to_remove.iteritems():
|
||||
self.nodes[nodename].remove(domains)
|
||||
for domain in domains:
|
||||
to_remove = []
|
||||
for p in self.pending:
|
||||
if p['domain'] == domain:
|
||||
to_remove.append(p)
|
||||
|
||||
for p in to_remove:
|
||||
self.pending.remove(p)
|
||||
|
||||
def discard(self, domains):
|
||||
"""Stop and remove all running and pending instances of the given
|
||||
|
|
@ -158,63 +187,10 @@ class ClusterMaster(object):
|
|||
d[proc['domain']] = node
|
||||
return d
|
||||
|
||||
@property
|
||||
def pending(self):
|
||||
"""Return dict of pending domains as domain -> node"""
|
||||
d = {}
|
||||
for node in self.nodes.itervalues():
|
||||
for p in node.pending:
|
||||
d[p['domain']] = node
|
||||
return d
|
||||
|
||||
@property
|
||||
def available_nodes(self):
|
||||
return (node for node in self.nodes.itervalues() if node.available)
|
||||
|
||||
def _dispatch_domains(self, domains, spider_settings, priority):
|
||||
"""Schedule the given domains in the availables nodes as good as
|
||||
possible. The algorithm follows the next rules (in order):
|
||||
|
||||
1. search for nodes with available capacity(running < maxproc) and (if
|
||||
any) schedules the domains there
|
||||
|
||||
2. if there isn't any node with available capacity it schedules the
|
||||
domain in the node with the smallest number of pending spiders
|
||||
"""
|
||||
|
||||
to_schedule = {} # domains to schedule per node
|
||||
pending_node = [] # list of #pending, node
|
||||
|
||||
for node in self.available_nodes:
|
||||
capacity = node.maxproc - len(node.running)
|
||||
#order nodes in pending_node according to insertion position, calculated from priority comparison, for stage 2.
|
||||
i = 0
|
||||
for p in node.pending:
|
||||
if p['priority'] <= priority:
|
||||
i += 1
|
||||
else:
|
||||
break
|
||||
bisect.insort(pending_node, (i, node))
|
||||
|
||||
#stage 1: use available capacity
|
||||
to_schedule[node.name] = []
|
||||
while domains and capacity > 0:
|
||||
to_schedule[node.name].append(domains.pop(0))
|
||||
capacity -= 1
|
||||
if not domains:
|
||||
break
|
||||
|
||||
#stage 2: queue in pendings the remaining domains.
|
||||
# a) pops out minor insertion-point node b) schedules the domain c) reinserts the node in list with insertion-point incremented by one.
|
||||
for domain in domains:
|
||||
insert_point, node = pending_node.pop(0)
|
||||
to_schedule[node.name].append(domain)
|
||||
bisect.insort(pending_node, (insert_point+1, node))
|
||||
|
||||
for nodename, domains in to_schedule.iteritems():
|
||||
if domains:
|
||||
self.nodes[nodename].schedule(domains, spider_settings, priority)
|
||||
|
||||
def _engine_started(self):
|
||||
self.load_nodes()
|
||||
scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL'))
|
||||
|
|
|
|||
|
|
@ -29,14 +29,14 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
s += "<h2>Home</h2>\n"
|
||||
|
||||
s += "<table border='1'>\n"
|
||||
s += "<tr><th> </th><th>Name</th><th>Available</th><th>Running</th><th>Pending</th><th>Load.avg</th></tr>\n"
|
||||
s += "<tr><th> </th><th>Name</th><th>Available</th><th>Running</th><th>Load.avg</th></tr>\n"
|
||||
for node in self.nodes.itervalues():
|
||||
#chkbox = "<input type='checkbox' name='shutdown' value='%s' />" % domain if node.status in ["up", "idle"] else " "
|
||||
nodelink = "<a href='nodes/#%s'>%s</a>" % (node.name, node.name)
|
||||
chkbox = " "
|
||||
loadavg = "%.2f %.2f %.2f" % node.loadavg
|
||||
s += "<tr><td>%s</td><td>%s</td><td>%s</td><td>%d/%d</td><td>%d</td><td>%s</td></tr>\n" % \
|
||||
(chkbox, nodelink, node.available, len(node.running), node.maxproc, len(node.pending), loadavg)
|
||||
s += "<tr><td>%s</td><td>%s</td><td>%s</td><td>%d/%d</td><td>%s</td></tr>\n" % \
|
||||
(chkbox, nodelink, node.available, len(node.running), node.maxproc, loadavg)
|
||||
s += "</table>\n"
|
||||
|
||||
s += "</body>\n"
|
||||
|
|
@ -51,8 +51,7 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
self.update_nodes()
|
||||
|
||||
if "schedule" in args:
|
||||
node = args["node"][0] if "node" in args else None
|
||||
self.schedule(args["schedule"], nodename=node, priority=eval(args["priority"][0]))
|
||||
self.schedule(args["schedule"], priority=eval(args["priority"][0]))
|
||||
|
||||
if "stop" in args:
|
||||
self.stop(args["stop"])
|
||||
|
|
@ -90,20 +89,6 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
s += "</form>\n"
|
||||
else:
|
||||
s += "<p>No running domains on %s</p>\n" % node.name
|
||||
|
||||
# pending domains
|
||||
s += "<h3>Pending domains</h3>\n"
|
||||
if node.pending:
|
||||
s += "<form method='post' action='.'>\n"
|
||||
s += "<select name='remove' multiple='multiple' size='10'>\n"
|
||||
for p in node.pending:
|
||||
s += "<option value='%s'>%s (P:%s)</option>\n" % (p['domain'], p['domain'],p['priority'])
|
||||
s += "</select>\n"
|
||||
s += "<input type='hidden' name='node' value='%s'>\n" % node.name
|
||||
s += "<p><input type='submit' value='Remove selected pending domains on %s'></p>\n" % node.name
|
||||
s += "</form>\n"
|
||||
else:
|
||||
s += "<p>No pending domains on %s</p>\n" % node.name
|
||||
|
||||
return str(s)
|
||||
|
||||
|
|
@ -113,7 +98,7 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
|
||||
enabled_domains = set(spiders.asdict(include_disabled=False).keys())
|
||||
print "Enabled domains: %s" % len(enabled_domains)
|
||||
inactive_domains = enabled_domains - set(self.running.keys() + self.pending.keys())
|
||||
inactive_domains = enabled_domains - set(self.running.keys() + [p['domain'] for p in self.pending])
|
||||
|
||||
s = self.render_header()
|
||||
|
||||
|
|
@ -126,14 +111,6 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
s += "<option>%s</option>\n" % domain
|
||||
s += "</select>\n"
|
||||
s += "</br>\n"
|
||||
s += "Node (only available nodes shown):<br />\n"
|
||||
s += "<select name='node'>\n"
|
||||
s += "<option value='' selected='selected'>any</option>"
|
||||
for node in self.available_nodes:
|
||||
domcount = "%d/%d/%d" % (len(node.running), node.maxproc, len(node.pending))
|
||||
loadavg = "%.2f %.2f %.2f" % node.loadavg
|
||||
s += "<option value='%s'>%s [D: %s | LA: %s]</option>" % (node.name, node.name, domcount, loadavg)
|
||||
s += "</select><br />\n"
|
||||
s += "Priority:<br />\n"
|
||||
s += "<select name='priority'>\n"
|
||||
for p, pname in priorities.items():
|
||||
|
|
@ -150,9 +127,21 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
s += "<table border='1'>\n"
|
||||
s += "<tr><th>Domain</th><th>Status</th><th>Node</th></tr>\n"
|
||||
s += self._domains_table(self.running, '<b>running</b>')
|
||||
s += self._domains_table(self.pending, 'pending')
|
||||
s += "</table>\n"
|
||||
|
||||
# pending domains
|
||||
s += "<h3>Pending domains</h3>\n"
|
||||
if self.pending:
|
||||
s += "<form method='post' action='.'>\n"
|
||||
s += "<select name='remove' multiple='multiple' size='10'>\n"
|
||||
for p in self.pending:
|
||||
s += "<option value='%s'>%s (P:%s)</option>\n" % (p['domain'], p['domain'],p['priority'])
|
||||
s += "</select>\n"
|
||||
s += "<p><input type='submit' value='Remove selected pending domains'></p>\n"
|
||||
s += "</form>\n"
|
||||
else:
|
||||
s += "<p>No pending domains</p>\n"
|
||||
|
||||
return str(s)
|
||||
|
||||
def render_header(self):
|
||||
|
|
|
|||
|
|
@ -42,7 +42,6 @@ class ScrapyProcessProtocol(protocol.ProcessProtocol):
|
|||
def processEnded(self, status_object):
|
||||
log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile))
|
||||
del self.procman.running[self.domain]
|
||||
self.procman.next_pending()
|
||||
|
||||
class ClusterWorker(pb.Root):
|
||||
|
||||
|
|
@ -53,99 +52,46 @@ class ClusterWorker(pb.Root):
|
|||
self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC')
|
||||
self.logdir = settings['CLUSTER_WORKER_LOGDIR']
|
||||
self.running = {}
|
||||
self.pending = []
|
||||
self.starttime = time.time()
|
||||
port = settings.getint('CLUSTER_WORKER_PORT')
|
||||
scrapyengine.listenTCP(port, pb.PBServerFactory(self))
|
||||
|
||||
def remote_schedule(self, domains, spider_settings=None, priority=20):
|
||||
"""Schedule new domains to be crawled in a separate processes"""
|
||||
|
||||
responses = []
|
||||
for domain in domains:
|
||||
if len(self.running) < self.maxproc and domain not in self.running:
|
||||
self._run(domain, spider_settings)
|
||||
responses.append("Started %s" % self.running[domain])
|
||||
else:
|
||||
i = 0
|
||||
for p in self.pending:
|
||||
if p['priority'] <= priority:
|
||||
i += 1
|
||||
else:
|
||||
break
|
||||
self.pending.insert(i, {'domain': domain, 'settings': spider_settings, 'priority': priority})
|
||||
responses.append("Scheduled domain %s at position %s in queue" % (domain, i))
|
||||
return self.status(responses)
|
||||
|
||||
def remote_stop(self, domains):
|
||||
"""Stop running domains. For removing pending (not yet started) domains
|
||||
use remove() instead"""
|
||||
|
||||
responses = []
|
||||
for domain in domains:
|
||||
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"
|
||||
responses.append("Stopped process %s" % proc)
|
||||
else:
|
||||
responses.append("%s: domain not running." % domain)
|
||||
return self.status(responses)
|
||||
|
||||
def remote_remove(self, domains):
|
||||
"""Remove all scheduled instances of the given domains (if it hasn't
|
||||
started yet). Otherwise use stop()"""
|
||||
|
||||
responses = []
|
||||
for domain in domains:
|
||||
to_remove = []
|
||||
for p in self.pending:
|
||||
if p['domain'] == domain:
|
||||
to_remove.append(p)
|
||||
|
||||
for p in to_remove:
|
||||
self.pending.remove(p)
|
||||
responses.append("Unscheduled domain %s" % domain)
|
||||
return self.status(responses)
|
||||
def remote_stop(self, domain):
|
||||
"""Stop running domain."""
|
||||
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"
|
||||
return self.status(0, "Stopped process %s" % proc)
|
||||
else:
|
||||
return self.status(1, "%s: domain not running." % domain)
|
||||
|
||||
def remote_status(self):
|
||||
return self.status()
|
||||
|
||||
def status(self, response="Status Response"):
|
||||
def status(self, rcode=0, rstring=None):
|
||||
status = {}
|
||||
status["pending"] = self.pending
|
||||
status["running"] = [ self.running[k].as_dict() for k in self.running.keys() ]
|
||||
status["starttime"] = self.starttime
|
||||
status["timestamp"] = time.time()
|
||||
status["maxproc"] = self.maxproc
|
||||
status["loadavg"] = os.getloadavg()
|
||||
status["logdir"] = self.logdir
|
||||
status["callresponse"] = response
|
||||
status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.")
|
||||
return status
|
||||
|
||||
def next_pending(self):
|
||||
"""Run the next domain in the pending list, which is not already running"""
|
||||
|
||||
if len(self.running) >= self.maxproc:
|
||||
return
|
||||
for p in self.pending:
|
||||
if p['domain'] not in self.running:
|
||||
self._run(p['domain'], p['settings'])
|
||||
self.pending.remove(p)
|
||||
return
|
||||
|
||||
def _run(self, domain, spider_settings=None):
|
||||
"""Spawn process to run the given domain. Don't call this method
|
||||
directly. Instead use schedule()."""
|
||||
|
||||
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)
|
||||
|
||||
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
|
||||
|
||||
def remote_run(self, domain, spider_settings=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)
|
||||
|
||||
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.")
|
||||
|
||||
|
|
|
|||
|
|
@ -13,16 +13,10 @@ sys.argv.pop(0)
|
|||
|
||||
if not sys.argv:
|
||||
d.addCallback(lambda object: object.callRemote("status"))
|
||||
elif sys.argv[0] == "-d":
|
||||
try:
|
||||
priority = eval(sys.argv[2])
|
||||
except:
|
||||
priority = 20
|
||||
d.addCallback(lambda object: object.callRemote("schedule", [sys.argv[1]], priority=priority))
|
||||
elif sys.argv[0] == "-s":
|
||||
d.addCallback(lambda object: object.callRemote("stop", [sys.argv[1]]))
|
||||
d.addCallback(lambda object: object.callRemote("stop", sys.argv[1]))
|
||||
elif sys.argv[0] == "-r":
|
||||
d.addCallback(lambda object: object.callRemote("remove", [sys.argv[1]]))
|
||||
d.addCallback(lambda object: object.callRemote("run", sys.argv[1]))
|
||||
|
||||
d.addCallbacks(callback = util.println, errback = lambda reason: 'error: '+str(reason.value))
|
||||
d.addCallback(lambda _: reactor.stop())
|
||||
|
|
|
|||
Loading…
Reference in New Issue