deleting old cluster branch

--HG--
extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40151
This commit is contained in:
olveyra 2008-08-07 12:26:43 +00:00
parent 59d6e92582
commit 4606b7f9c5
7 changed files with 0 additions and 583 deletions

View File

@ -1,2 +0,0 @@
from scrapy.contrib.cluster.master import ClusterMaster, ClusterMasterWeb
from scrapy.contrib.cluster.worker import ClusterWorker, ClusterWorkerWeb

View File

@ -1,2 +0,0 @@
from scrapy.contrib.cluster.master.manager import ClusterMaster
from scrapy.contrib.cluster.master.web import ClusterMasterWeb

View File

@ -1,189 +0,0 @@
import urlparse
import urllib
import bisect
from pydispatch import dispatcher
from twisted.web.client import getPage
from scrapy.core import log, signals
from scrapy.core.engine import scrapyengine
from scrapy.core.exceptions import NotConfigured
from scrapy.utils.serialization import unserialize
from scrapy.conf import settings
class ClusterNode(object):
def __init__(self, name, url):
self.name = name
self.url = url
self.maxproc = 0
self.running = {}
self.pending = []
self.loadavg = (0.0, 0.0, 0.0)
self.status = "down" # down/crawling/idle/error
self.available = False
self.wsurl = urlparse.urljoin(self.url, "cluster_worker/ws/?format=json")
def update(self):
d = getPage(self.wsurl)
d.addCallbacks(self._cbUpdate, self._ebUpdate)
def schedule(self, domains):
args = [("schedule", domain) for domain in domains]
self._wsRequest(args)
def stop(self, domains):
args = [("stop", domain) for domain in domains]
self._wsRequest(args)
def remove(self, domains):
args = [("remove", domain) for domain in domains]
self._wsRequest(args)
def _wsRequest(self, args):
d = getPage("%s?%s" % (self.wsurl, urllib.urlencode(args)))
d.addCallbacks(self._cbUpdate, self._ebUpdate)
def _cbUpdate(self, jsonstatus):
try:
pmstatus = unserialize(jsonstatus, 'json')
self.maxproc = int(pmstatus.get('maxproc', None))
self.running = pmstatus.get('running') or {}
self.pending = pmstatus.get('pending') or []
self.loadavg = tuple(pmstatus.get('loadavg', (0.0, 0.0, 0.0)))
self.status = "crawling" if self.running else "idle"
self.available = True
except ValueError:
self.status = "error"
self.available = False
def _ebUpdate(self, err):
self.status = "down"
self.available = False
class ClusterMaster(object):
def __init__(self):
if not settings.getbool('CLUSTER_MASTER_ENABLED'):
raise NotConfigured
self.nodes = {}
dispatcher.connect(self._engine_started, signal=signals.engine_started)
def load_nodes(self):
"""Loads nodes from the CLUSTER_MASTER_NODES setting"""
for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
self.add_node(name, url)
def update_nodes(self):
for node in self.nodes.itervalues():
node.update()
def add_node(self, name, url):
"""Add node given its node"""
if name not in self.nodes:
node = ClusterNode(name, url)
self.nodes[name] = node
node.update()
def remove_node(self, nodename):
raise NotImplemented
def schedule(self, domains, nodename=None):
if nodename:
node = self.nodes[nodename]
node.schedule(domains)
else:
self._dispatch_domains(domains)
def stop(self, domains):
to_stop = {}
for domain in domains:
node = self.running.get(domain, None)
if node:
if node.name not in to_stop:
to_stop[node.name] = []
to_stop[node.name].append(domain)
for nodename, domains in to_stop.iteritems():
self.nodes[nodename].stop(domains)
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)
for nodename, domains in to_remove.iteritems():
self.nodes[nodename].remove(domains)
def discard(self, domains):
"""Stop and remove all running and pending instances of the given
domains"""
self.remove(domains)
self.stop(domains)
@property
def running(self):
"""Return dict of running domains as domain -> node"""
d = {}
for node in self.nodes.itervalues():
for domain in node.running.iterkeys():
d[domain] = node
return d
@property
def pending(self):
"""Return dict of pending domains as domain -> node"""
d = {}
for node in self.nodes.itervalues():
for domain in node.pending:
d[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):
"""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)
bisect.insort(pending_node, (len(node.pending), node))
to_schedule[node.name] = []
while domains and capacity > 0:
to_schedule[node.name].append(domains.pop(0))
capacity -= 1
if not domains:
break
for domain in domains:
pending, node = pending_node.pop(0)
to_schedule[node.name].append(domain)
bisect.insort(pending_node, (pending+1, node))
for nodename, domains in to_schedule.iteritems():
if domains:
self.nodes[nodename].schedule(domains)
def _engine_started(self):
self.load_nodes()
scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL'))

View File

@ -1,165 +0,0 @@
import datetime
from pydispatch import dispatcher
from scrapy.spider import spiders
from scrapy.management.web import banner, webconsole_discover_module
from scrapy.utils.serialization import parse_jsondatetime
from scrapy.contrib.cluster.master import ClusterMaster
class ClusterMasterWeb(ClusterMaster):
webconsole_id = 'cluster_master'
webconsole_name = 'Cluster master'
def __init__(self):
ClusterMaster.__init__(self)
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
def webconsole_render(self, wc_request):
changes = ""
if wc_request.path == '/cluster/nodes/':
return self.render_nodes(wc_request)
elif wc_request.path == '/cluster/domains/':
return self.render_domains(wc_request)
elif wc_request.args:
changes = self.webconsole_control(wc_request)
s = self.render_header()
s += "<h2>Home</h2>\n"
s += "<table border='1'>\n"
s += "<tr><th>&nbsp;</th><th>Name</th><th>Status</th><th>Running</th><th>Pending</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 "&nbsp;"
nodelink = "<a href='nodes/#%s'>%s</a>" % (node.name, node.name)
chkbox = "&nbsp;"
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.status, len(node.running), node.maxproc, len(node.pending), loadavg)
s += "</table>\n"
s += "</body>\n"
s += "</html>\n"
return str(s)
def webconsole_control(self, wc_request):
args = wc_request.args
if "updatenodes" in args:
self.update_nodes()
if "schedule" in args:
node = args["node"][0] if "node" in args else None
self.schedule(args["schedule"], nodename=node)
if "stop" in args:
self.stop(args["stop"])
if "remove" in args:
self.remove(args["remove"])
return ""
def render_nodes(self, wc_request):
if wc_request.args:
self.webconsole_control(wc_request)
now = datetime.datetime.utcnow()
s = self.render_header()
for node in self.nodes.itervalues():
s += "<h2><a name='%s'>%s</h2>\n" % (node.name, node.name)
s += "<h3>Running domains</h3>\n"
if node.running:
s += "<form method='post' action='.'>\n"
s += "<table border='1'>\n"
s += "<tr><th>&nbsp;</th><th>PID</th><th>Domain</th><th>Status</th><th>Running time</th><th>Log file</th></tr>\n"
for domain, proc in node.running.iteritems():
chkbox = "<input type='checkbox' name='stop' value='%s' />" % domain if proc['status'] == "running" else "&nbsp;"
start_time = parse_jsondatetime(proc.get('start_time', None))
elapsed = now - start_time if start_time else None
s += "<tr><td>%s</td><td>%s</td><td>%s</td><td>%s</td><td>%s</td><td>%s</td></tr>\n" % \
(chkbox, proc['pid'], domain, proc['status'], elapsed, proc['logfile'])
s += "</table>\n"
s += "<input type='hidden' name='node' value='%s'>\n" % node.name
s += "<p><input type='submit' value='Stop selected domains on %s'></p>\n" % node.name
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 domain in node.pending:
s += "<option>%s</option>\n" % domain
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)
def render_domains(self, wc_request):
if wc_request.args:
self.webconsole_control(wc_request)
enabled_domains = set(spiders.asdict(include_disabled=False).keys())
inactive_domains = enabled_domains - set(self.running.keys() + self.pending.keys())
s = self.render_header()
s += "<h2>Schedule domains</h2>\n"
s += "Inactive domains (not running or pending)<br />"
s += "<form method='post' action='.'>\n"
s += "<select name='schedule' multiple='multiple' size='10'>\n"
for domain in sorted(inactive_domains):
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>\n"
s += "<p><input type='submit' value='Schedule selected domains'></p>\n"
s += "</form>\n"
s += "<h2>Domains</h2>\n"
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"
return str(s)
def render_header(self):
s = banner(self)
s += "<p>Nav: "
s += "<a href='/cluster/'>Home</a> | "
s += "<a href='/cluster/domains/'>Domains</a> | "
s += "<a href='/cluster/nodes/'>Nodes</a> (<a href='/cluster/nodes/?updatenodes=1'>update</a>)"
s += "</p>"
return s
def _domains_table(self, dict_, status):
s = ""
for domain, node in dict_.iteritems():
s += "<tr><td>%s</td><td>%s</td><td>%s</td></tr>\n" % (domain, status, node.name)
return s
def webconsole_discover_module(self):
return self

View File

@ -1,2 +0,0 @@
from scrapy.contrib.cluster.worker.manager import ClusterWorker
from scrapy.contrib.cluster.worker.web import ClusterWorkerWeb

View File

@ -1,96 +0,0 @@
import sys
import os
import time
import datetime
from twisted.internet import protocol, reactor
from scrapy.core import log
from scrapy.core.exceptions import NotConfigured
from scrapy.conf import settings
class ScrapyProcessProtocol(protocol.ProcessProtocol):
def __init__(self, procman, domain, logfile=None):
self.procman = procman
self.domain = domain
self.logfile = logfile
self.start_time = datetime.datetime.utcnow()
self.status = "starting"
self.pid = -1
def __str__(self):
return "<ScrapyProcess domain=%s, pid=%s, status=%s>" % (self.domain, self.pid, self.status)
def connectionMade(self):
self.pid = self.transport.pid
log.msg("ClusterWorker: started domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile))
self.transport.closeStdin()
self.status = "running"
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(object):
def __init__(self):
if not settings.getbool('CLUSTER_WORKER_ENABLED'):
raise NotConfigured
self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC')
self.logdir = settings['CLUSTER_WORKER_LOGDIR']
self.running = {}
self.pending = []
def schedule(self, domain):
"""Schedule new domain to be crawled in a separate processes"""
if len(self.running) < self.maxproc and domain not in self.running:
self._run(domain)
else:
self.pending.append(domain)
def stop(self, domain):
"""Stop running domain. For removing pending (not yet started) domains
use remove() instead"""
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"
def remove(self, domain):
"""Remove all scheduled instances of the given domain (if it hasn't
started yet). Otherwise use stop()"""
while domain in self.pending:
self.pending.remove(domain)
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 domain in self.pending:
if domain not in self.running:
self._run(domain)
self.pending.remove(domain)
return
def _run(self, domain):
"""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)
env = {'SCRAPY_LOGFILE': logfile}
args = [sys.executable, sys.argv[0], 'crawl', domain]
proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=env)
self.running[domain] = scrapy_proc

View File

@ -1,127 +0,0 @@
import os
import datetime
from pydispatch import dispatcher
from scrapy.spider import spiders
from scrapy.management.web import banner, webconsole_discover_module
from scrapy.utils.serialization import serialize
from scrapy.contrib.cluster.worker import ClusterWorker
class ClusterWorkerWeb(ClusterWorker):
webconsole_id = 'cluster_worker'
webconsole_name = 'Cluster worker'
def __init__(self):
ClusterWorker.__init__(self)
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
def webconsole_render(self, wc_request):
changes = ""
if wc_request.path == '/cluster_worker/ws/':
return self.webconsole_control(wc_request, ws=True)
elif wc_request.args:
changes = self.webconsole_control(wc_request)
now = datetime.datetime.utcnow()
s = banner(self)
# running processes
s += "<h2>Running processes</h2>\n"
if self.running:
s += "<form method='post' action='.'>\n"
s += "<table border='1'>\n"
s += "<tr><th>&nbsp;</th><th>PID</th><th>Domain</th><th>Status</th><th>Log file</th><th>Running time</th></tr>\n"
for domain, proc in self.running.iteritems():
chkbox = "<input type='checkbox' name='stop' value='%s' />" % domain if proc.status == "running" else "&nbsp;"
elapsed = now - proc.start_time
s += "<tr><td>%s</td><td>%s</td><td>%s</td><td>%s</td><td>%s</td><td>%s</td></tr>\n" % \
(chkbox, proc.pid, domain, proc.status, proc.logfile, elapsed)
s += "</table>\n"
s += "<p><input type='submit' value='Stop selected domains'></p>\n"
s += "</form>\n"
else:
s += "<p>No running processes</p>\n"
# pending domains
s += "<h2>Pending domains</h2>\n"
if self.pending:
s += "<form method='post' action='.'>\n"
s += "<select name='remove' multiple='multiple' size='10'>\n"
for domain in self.pending:
s += "<option>%s</option>\n" % domain
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"
# schedule domains
enabled_domains = spiders.asdict(include_disabled=False).keys()
s += "<h2>Schedule domains</h2>\n"
s += "<form method='post' action='.'>\n"
s += "<select name='schedule' multiple='multiple' size='10'>\n"
for domain in sorted(enabled_domains):
s += "<option>%s</option>\n" % domain
s += "</select>\n"
s += "<p><input type='submit' value='Schedule selected domains'></p>\n"
s += "</form>\n"
s += changes
s += "</body>\n"
s += "</html>\n"
return s
def webconsole_control(self, wc_request, ws=False):
args = wc_request.args
if "schedule" in args:
for domain in args["schedule"]:
self.schedule(domain)
if ws:
return self.ws_status(wc_request)
else:
return "<p>Scheduled domains: <ul><li>%s</li></ul>" % "</li><li>".join(args["schedule"]) + "</p>\n"
if "stop" in args:
for domain in args["stop"]:
self.stop(domain)
if ws:
return self.ws_status(wc_request)
else:
return "<p>Stopped running processes: <ul><li>%s</li></ul>" % "</li><li>".join(args["stop"]) + "</p>\n"
if "remove" in args:
for domain in args["remove"]:
self.remove(domain)
if ws:
return self.ws_status(wc_request)
else:
return "<p>Removed pending domains: <ul><li>%s</li></ul>" % "</li><li>".join(args["remove"]) + "</p>\n"
if ws:
return self.ws_status(wc_request)
else:
return ""
def ws_status(self, wc_request):
format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json'
wc_request.setHeader('content-type', 'text/plain')
exported_proc_attrs = ['pid', 'status', 'start_time', 'logfile']
d = {'maxproc': self.maxproc, 'running': {}, 'pending': []}
for domain, proc in self.running.iteritems():
d2 = {}
for a in exported_proc_attrs:
d2[a] = getattr(proc, a)
d['running'][domain] = d2
d['pending'] = self.pending
d['loadavg'] = os.getloadavg()
content = serialize(d, format)
return content
def webconsole_discover_module(self):
return self