mirror of https://github.com/scrapy/scrapy.git
cluster: promote new code as replacement of old pbcluster
--HG-- rename : scrapy/trunk/scrapy/contrib/pbcluster/crawler/__init__.py => scrapy/trunk/scrapy/contrib/cluster/crawler/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/crawler/manager.py => scrapy/trunk/scrapy/contrib/cluster/crawler/manager.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/hooks/__init__.py => scrapy/trunk/scrapy/contrib/cluster/hooks/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/hooks/svn.py => scrapy/trunk/scrapy/contrib/cluster/hooks/svn.py rename : scrapy/trunk/scrapy/contrib/pbcluster/master/__init__.py => scrapy/trunk/scrapy/contrib/cluster/master/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/master/manager.py => scrapy/trunk/scrapy/contrib/cluster/master/manager.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/master/web.py => scrapy/trunk/scrapy/contrib/cluster/master/web.py rename : scrapy/trunk/scrapy/contrib/pbcluster/master/ws_api.txt => scrapy/trunk/scrapy/contrib/cluster/master/ws_api.txt rename : scrapy/trunk/scrapy/contrib/pbcluster/tools/scrapy-cluster-ctl.py => scrapy/trunk/scrapy/contrib/cluster/tools/scrapy-cluster-ctl.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/tools/test-worker.py => scrapy/trunk/scrapy/contrib/cluster/tools/test-worker.py rename : scrapy/trunk/scrapy/contrib/pbcluster/worker/__init__.py => scrapy/trunk/scrapy/contrib/cluster/worker/__init__.py rename : scrapy/trunk/scrapy/contrib_exp/cluster/worker/manager.py => scrapy/trunk/scrapy/contrib/cluster/worker/manager.py extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%401064
This commit is contained in:
parent
821f6be3ce
commit
d6c52d51ed
|
|
@ -0,0 +1,3 @@
|
|||
from scrapy.contrib.cluster.worker.manager import ClusterWorker
|
||||
from scrapy.contrib.cluster.master.web import ClusterMasterWeb
|
||||
from scrapy.contrib.cluster.crawler.manager import ClusterCrawler
|
||||
|
|
@ -11,7 +11,7 @@ from scrapy.core import signals
|
|||
from scrapy import log
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.core.exceptions import NotConfigured
|
||||
from scrapy.contrib_exp.cluster.worker.manager import ResponseCode
|
||||
from scrapy.contrib.cluster.worker.manager import ResponseCode
|
||||
from scrapy.conf import settings
|
||||
|
||||
def my_import(name):
|
||||
|
|
@ -4,7 +4,7 @@ from pydispatch import dispatcher
|
|||
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.management.web import banner, webconsole_discover_module
|
||||
from scrapy.contrib_exp.cluster.master.manager import ClusterMaster
|
||||
from scrapy.contrib.cluster.master.manager import ClusterMaster
|
||||
from scrapy.utils.serialization import serialize
|
||||
|
||||
class ClusterMasterWeb(ClusterMaster):
|
||||
|
|
@ -1,3 +0,0 @@
|
|||
from scrapy.contrib.pbcluster.worker.manager import ClusterWorker
|
||||
from scrapy.contrib.pbcluster.master.web import ClusterMasterWeb
|
||||
from scrapy.contrib.pbcluster.crawler.manager import ClusterCrawler
|
||||
|
|
@ -1,38 +0,0 @@
|
|||
import os
|
||||
|
||||
from twisted.spread import pb
|
||||
from twisted.internet import reactor
|
||||
|
||||
from scrapy.conf import settings
|
||||
from scrapy import log
|
||||
from scrapy.core.manager import scrapymanager
|
||||
from scrapy.core.exceptions import NotConfigured
|
||||
|
||||
class Broker(pb.Referenceable):
|
||||
def __init__(self, crawler, remote):
|
||||
self.__remote = remote
|
||||
self.__crawler = crawler
|
||||
try:
|
||||
deferred = self.__remote.callRemote("register_crawler", os.getpid(), self)
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("Lost connection to node %s." % (self.name), log.ERROR)
|
||||
else:
|
||||
deferred.addCallbacks(callback=lambda x: None, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
def remote_stop(self):
|
||||
scrapymanager.stop()
|
||||
|
||||
class ClusterCrawler:
|
||||
def __init__(self):
|
||||
if not settings.getbool('CLUSTER_CRAWLER_ENABLED'):
|
||||
raise NotConfigured
|
||||
|
||||
self.worker = None
|
||||
|
||||
factory = pb.PBClientFactory()
|
||||
reactor.connectTCP("localhost", settings.getint('CLUSTER_WORKER_PORT'), factory)
|
||||
d = factory.getRootObject()
|
||||
def _set_worker(obj):
|
||||
self.worker = Broker(self, obj)
|
||||
d.addCallbacks(callback=_set_worker, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
|
||||
|
|
@ -1,330 +0,0 @@
|
|||
import sys, datetime
|
||||
import pickle
|
||||
|
||||
from pydispatch import dispatcher
|
||||
|
||||
from twisted.spread import pb
|
||||
from twisted.internet import reactor
|
||||
|
||||
from scrapy.core import signals
|
||||
from scrapy import log
|
||||
from scrapy.core.engine import scrapyengine
|
||||
from scrapy.core.exceptions import NotConfigured
|
||||
from scrapy.conf import settings
|
||||
|
||||
DEFAULT_PRIORITY = settings.getint("DEFAULT_PRIORITY", 20)
|
||||
|
||||
def my_import(name):
|
||||
mod = __import__(name)
|
||||
components = name.split('.')
|
||||
for comp in components[1:]:
|
||||
mod = getattr(mod, comp)
|
||||
return mod
|
||||
|
||||
class Broker(pb.Referenceable):
|
||||
def __init__(self, remote, name, master):
|
||||
self.__remote = remote
|
||||
self.alive = False
|
||||
self.name = name
|
||||
self.master = master
|
||||
self.available = True
|
||||
try:
|
||||
deferred = self.__remote.callRemote("set_master", self)
|
||||
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 status_as_dict(self, verbosity=1):
|
||||
if verbosity == 0:
|
||||
return
|
||||
status = {"alive": self.alive}
|
||||
if self.alive:
|
||||
if verbosity == 1:
|
||||
#dont show spider settings
|
||||
status["running"] = []
|
||||
for proc in self.running:
|
||||
proccopy = proc.copy()
|
||||
del proccopy["settings"]
|
||||
status["running"].append(proccopy)
|
||||
elif verbosity == 2:
|
||||
status["running"] = self.running
|
||||
status["maxproc"] = self.maxproc
|
||||
status["freeslots"] = self.maxproc - len(self.running)
|
||||
status["available"] = self.available
|
||||
status["starttime"] = self.starttime
|
||||
status["timestamp"] = self.timestamp
|
||||
status["loadavg"] = self.loadavg
|
||||
return status
|
||||
|
||||
def _set_status(self, status):
|
||||
if not status:
|
||||
self.alive = False
|
||||
else:
|
||||
self.alive = True
|
||||
self.running = status['running']
|
||||
self.maxproc = status['maxproc']
|
||||
self.starttime = status['starttime']
|
||||
self.timestamp = status['timestamp']
|
||||
self.loadavg = status['loadavg']
|
||||
self.logdir = status['logdir']
|
||||
free_slots = self.maxproc - len(self.running)
|
||||
|
||||
#load domains by one, so to mix up better the domain loading between nodes. The next one in the same node will be loaded
|
||||
#when there is no loading domain or in the next status update. This way also we load the nodes softly
|
||||
if self.available and free_slots > 0 and self.master.pending:
|
||||
pending = self.master.pending.pop(0)
|
||||
#if domain already running in some node, reschedule with same priority (so will be moved to run later)
|
||||
if pending['domain'] in self.master.running or pending['domain'] in self.master.loading:
|
||||
self.master.schedule([pending['domain']], pending['settings'], pending['priority'])
|
||||
else:
|
||||
self.run(pending)
|
||||
self.master.loading.append(pending['domain'])
|
||||
|
||||
def update_status(self):
|
||||
try:
|
||||
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 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 run(self, pending):
|
||||
|
||||
def _run_errback(reason):
|
||||
log.msg(reason, log.ERROR)
|
||||
self.master.loading.remove(pending['domain'])
|
||||
self.master.schedule([pending['domain']], pending['settings'], pending['priority'] - 1)
|
||||
log.msg("Domain %s rescheduled: lost connection to node." % pending['domain'], log.WARNING)
|
||||
|
||||
def _run_callback(status):
|
||||
if status['callresponse'][0] == 1:
|
||||
#slots are complete. Reschedule in master with priority reduced by one.
|
||||
#self.master.loading check should avoid this to happen
|
||||
self.master.loading.remove(pending['domain'])
|
||||
self.master.schedule([pending['domain']], pending['settings'], pending['priority'] - 1)
|
||||
log.msg("Domain %s rescheduled: no proc space in node." % pending['domain'], log.WARNING)
|
||||
elif status['callresponse'][0] == 2:
|
||||
#domain already running in node. Reschedule with same priority.
|
||||
#self.master.loading check should avoid this to happen
|
||||
self.master.loading.remove(pending['domain'])
|
||||
self.master.schedule([pending['domain']], pending['settings'], pending['priority'])
|
||||
log.msg("Domain %s rescheduled: already running in node." % pending['domain'], log.WARNING)
|
||||
|
||||
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=_run_errback)
|
||||
|
||||
def remote_update(self, status, domain, domain_status):
|
||||
self._set_status(status)
|
||||
if domain in self.master.loading and domain_status == "running":
|
||||
self.master.loading.remove(domain)
|
||||
self.master.statistics["domains"]["running"].add(domain)
|
||||
elif domain_status == "scraped":
|
||||
self.master.statistics["domains"]["running"].remove(domain)
|
||||
self.master.statistics["domains"]["scraped"][domain] = self.master.statistics["domains"]["scraped"].get(domain, 0) + 1
|
||||
self.master.statistics["scraped_count"] = self.master.statistics.get("scraped_count", 0) + 1
|
||||
if domain in self.master.statistics["domains"]["lost"]:
|
||||
self.master.statistics["domains"]["lost"].remove(domain)
|
||||
|
||||
class ScrapyPBClientFactory(pb.PBClientFactory):
|
||||
def __init__(self, master, nodename):
|
||||
pb.PBClientFactory.__init__(self)
|
||||
self.master = master
|
||||
self.nodename = nodename
|
||||
|
||||
def clientConnectionLost(self, *args, **kargs):
|
||||
pb.PBClientFactory.clientConnectionLost(self, *args, **kargs)
|
||||
del self.master.nodes[self.nodename]
|
||||
log.msg("Removed node %s." % self.nodename )
|
||||
|
||||
class ClusterMaster:
|
||||
|
||||
def __init__(self):
|
||||
|
||||
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(settings["CLUSTER_MASTER_CACHEFILE"], "r") )
|
||||
except IOError:
|
||||
self.pending = []
|
||||
self.loading = []
|
||||
self.nodes = {}
|
||||
self.start_time = datetime.datetime.utcnow()
|
||||
#on how statistics works, see self.update_nodes() and Broker.remote_update()
|
||||
self.statistics = {"domains": {"running": set(), "scraped": {}, "lost_count": {}, "lost": set()}, "scraped_count": 0 }
|
||||
self.global_settings = {}
|
||||
#load cluster global settings
|
||||
for sname in settings.getlist('GLOBAL_CLUSTER_SETTINGS'):
|
||||
self.global_settings[sname] = settings[sname]
|
||||
|
||||
dispatcher.connect(self._engine_started, signal=signals.engine_started)
|
||||
dispatcher.connect(self._engine_stopped, signal=signals.engine_stopped)
|
||||
|
||||
def load_nodes(self):
|
||||
"""Loads nodes listed in CLUSTER_MASTER_NODES setting"""
|
||||
for name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
|
||||
self.load_node(name, url)
|
||||
|
||||
def load_node(self, name, url):
|
||||
"""Creates the remote reference for each worker node"""
|
||||
def _make_callback(_factory, _name, _url):
|
||||
|
||||
def _errback(_reason):
|
||||
log.msg("Could not get remote node %s in %s: %s." % (_name, _url, _reason), log.ERROR)
|
||||
|
||||
d = _factory.getRootObject()
|
||||
d.addCallbacks(callback=lambda obj: self.add_node(obj, _name), errback=_errback)
|
||||
|
||||
server, port = url.split(":")
|
||||
port = int(port)
|
||||
log.msg("Connecting to cluster worker %s..." % name)
|
||||
log.msg("Server: %s, Port: %s" % (server, port))
|
||||
factory = ScrapyPBClientFactory(self, name)
|
||||
try:
|
||||
reactor.connectTCP(server, port, factory)
|
||||
except Exception, err:
|
||||
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 name, url in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
|
||||
if name in self.nodes and self.nodes[name].alive:
|
||||
log.msg("Updating node. name: %s, url: %s" % (name, url) )
|
||||
self.nodes[name].update_status()
|
||||
else:
|
||||
log.msg("Reloading node. name: %s, url: %s" % (name, url) )
|
||||
self.load_node(name, url)
|
||||
|
||||
real_running = set(self.running.keys())
|
||||
lost = self.statistics["domains"]["running"].difference(real_running)
|
||||
for domain in lost:
|
||||
self.statistics["domains"]["lost_count"][domain] = self.statistics["domains"]["lost_count"].get(domain, 0) + 1
|
||||
self.statistics["domains"]["lost"] = self.statistics["domains"]["lost"].union(lost)
|
||||
|
||||
def add_node(self, cworker, name):
|
||||
"""Add node given its node"""
|
||||
node = Broker(cworker, name, self)
|
||||
self.nodes[name] = node
|
||||
log.msg("Added cluster worker %s" % name)
|
||||
|
||||
def disable_node(self, name):
|
||||
self.nodes[name].available = False
|
||||
|
||||
def enable_node(self, name):
|
||||
self.nodes[name].available = True
|
||||
|
||||
def remove_node(self, nodename):
|
||||
raise NotImplemented
|
||||
|
||||
def schedule(self, domains, spider_settings=None, priority=DEFAULT_PRIORITY):
|
||||
i = 0
|
||||
for p in self.pending:
|
||||
if p['priority'] <= priority:
|
||||
i += 1
|
||||
else:
|
||||
break
|
||||
for domain in domains:
|
||||
pd = self.find_inpending(domain)
|
||||
if pd: #domain already pending, so just change priority if new is higher
|
||||
if priority < pd['priority']:
|
||||
self.pending.remove(pd)
|
||||
pd['priority'] = priority
|
||||
self.pending.insert(i, pd)
|
||||
else:
|
||||
final_spider_settings = self.get_spider_groupsettings(domain)
|
||||
final_spider_settings.update(self.global_settings)
|
||||
final_spider_settings.update(spider_settings or {})
|
||||
self.pending.insert(i, {'domain': domain, 'settings': final_spider_settings, 'priority': priority})
|
||||
|
||||
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():
|
||||
for domain in domains:
|
||||
self.nodes[nodename].stop(domain)
|
||||
|
||||
def remove(self, domains):
|
||||
"""Remove all scheduled instances of the given domains (if it hasn't
|
||||
started yet). Otherwise use stop()"""
|
||||
|
||||
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
|
||||
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 proc in node.running:
|
||||
d[proc['domain']] = node
|
||||
return d
|
||||
|
||||
@property
|
||||
def available_nodes(self):
|
||||
return (node for node in self.nodes.itervalues() if node.available)
|
||||
|
||||
def find_inpending(self, domain):
|
||||
for p in self.pending:
|
||||
if domain == p['domain']:
|
||||
return p
|
||||
|
||||
def print_pending(self, verbosity=1):
|
||||
if verbosity == 1:
|
||||
pending = []
|
||||
for p in self.pending:
|
||||
pp = p.copy()
|
||||
del pp["settings"]
|
||||
pending.append(pp)
|
||||
return pending
|
||||
elif verbosity == 2:
|
||||
return self.pending
|
||||
return
|
||||
|
||||
def _engine_started(self):
|
||||
self.load_nodes()
|
||||
scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL'))
|
||||
def _engine_stopped(self):
|
||||
pickle.dump( self.pending, open(settings["CLUSTER_MASTER_CACHEFILE"], "w") )
|
||||
log.msg("Pending saved in %s" % settings["CLUSTER_MASTER_CACHEFILE"])
|
||||
|
|
@ -1,239 +0,0 @@
|
|||
import datetime
|
||||
|
||||
from pydispatch import dispatcher
|
||||
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.management.web import banner, webconsole_discover_module
|
||||
from scrapy.contrib.pbcluster.master.manager import ClusterMaster, DEFAULT_PRIORITY
|
||||
from scrapy.utils.serialization import serialize
|
||||
|
||||
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_master/nodes/':
|
||||
return self.render_nodes(wc_request)
|
||||
elif wc_request.path == '/cluster_master/domains/':
|
||||
return self.render_domains(wc_request)
|
||||
elif wc_request.path == '/cluster_master/ws/':
|
||||
return self.webconsole_control(wc_request, ws=True)
|
||||
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> </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>%s</td></tr>\n" % \
|
||||
(chkbox, nodelink, node.available, len(node.running), node.maxproc, loadavg)
|
||||
s += "</table>\n"
|
||||
|
||||
s += "</body>\n"
|
||||
s += "</html>\n"
|
||||
|
||||
return str(s)
|
||||
|
||||
def webconsole_control(self, wc_request, ws=False):
|
||||
args = wc_request.args
|
||||
if "updatenodes" in args:
|
||||
self.update_nodes()
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
|
||||
if "schedule" in args:
|
||||
if ws:
|
||||
sep = ","
|
||||
domains = args["schedule"][0].split(sep)
|
||||
else:
|
||||
sep = "\r"
|
||||
domains = args["schedule"]
|
||||
priority = int(args.get("priority", [DEFAULT_PRIORITY])[0])
|
||||
|
||||
#spider settings
|
||||
slist = args.get("settings", [""])[0].split(sep)
|
||||
spider_settings = {}
|
||||
for s in slist:
|
||||
try:
|
||||
k, v = s.strip().split("=")
|
||||
except ValueError:
|
||||
pass
|
||||
else:
|
||||
spider_settings[k] = v
|
||||
|
||||
self.schedule(domains, spider_settings, priority)
|
||||
if ws:
|
||||
return self.ws_status(wc_request, verbosity=0)
|
||||
|
||||
if "stop" in args:
|
||||
if ws:
|
||||
domains = args["stop"][0].split(",")
|
||||
else:
|
||||
domains=args["stop"]
|
||||
self.stop(domains)
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
|
||||
if "remove" in args:
|
||||
if ws:
|
||||
domains = args["remove"][0].split(",")
|
||||
else:
|
||||
domains=args["remove"]
|
||||
self.remove(domains)
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
if "disable_node" in args:
|
||||
self.disable_node(args["disable_node"][0])
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
if "enable_node" in args:
|
||||
self.enable_node(args["enable_node"][0])
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
if "statistics" in args:
|
||||
if ws:
|
||||
return self.ws_statistics(wc_request)
|
||||
|
||||
if ws:
|
||||
return self.ws_status(wc_request)
|
||||
else:
|
||||
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():
|
||||
if node.available:
|
||||
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> </th><th>PID</th><th>Domain</th><th>Status</th><th>Running time</th><th>Log file</th></tr>\n"
|
||||
for proc in node.running:
|
||||
chkbox = "<input type='checkbox' name='stop' value='%s' />" % proc['domain'] if proc['status'] == "running" else " "
|
||||
start_time = proc.get('starttime', 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'], proc['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
|
||||
|
||||
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())
|
||||
print "Enabled domains: %s" % len(enabled_domains)
|
||||
inactive_domains = enabled_domains - set(self.running.keys() + [p['domain'] for p in self.pending])
|
||||
|
||||
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 += "Priority:<br />\n"
|
||||
s += "<input type='text' name='priority'>%s</input>" % DEFAULT_PRIORITY
|
||||
s += "<br />\n"
|
||||
|
||||
#spider settings
|
||||
s += "Overrided spider settings:<br />\n"
|
||||
s += "<textarea name='settings' rows='4'>\n"
|
||||
s += "UNAVAILABLES_NOTIFY=2\n"
|
||||
s += "</textarea>\n"
|
||||
s += "<br />\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 += "</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):
|
||||
s = banner(self)
|
||||
s += "<p>Nav: "
|
||||
s += "<a href='/cluster_master/'>Home</a> | "
|
||||
s += "<a href='/cluster_master/domains/'>Domains</a> | "
|
||||
s += "<a href='/cluster_master/nodes/'>Nodes</a> (<a href='/cluster_master/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
|
||||
|
||||
def ws_status(self, wc_request, verbosity=1):
|
||||
format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json'
|
||||
verbosity = int(wc_request.args['verbosity'][0]) if 'verbosity' in wc_request.args else verbosity
|
||||
wc_request.setHeader('content-type', 'text/plain')
|
||||
status = {}
|
||||
nodes_status = {}
|
||||
if verbosity > 0:
|
||||
for d, n in self.nodes.iteritems():
|
||||
nodes_status[d] = n.status_as_dict(verbosity)
|
||||
status["nodes"] = nodes_status
|
||||
status["pending"] = self.print_pending(verbosity)
|
||||
status["loading"] = self.loading
|
||||
content = serialize(status, format)
|
||||
return content
|
||||
return ""
|
||||
|
||||
def ws_statistics(self, wc_request):
|
||||
format = wc_request.args['format'][0] if 'format' in wc_request.args else 'json'
|
||||
content = serialize(self.statistics, format)
|
||||
return content
|
||||
|
|
@ -1,138 +0,0 @@
|
|||
import sys, os, time, datetime, pickle, gzip
|
||||
|
||||
from twisted.internet import protocol, reactor
|
||||
from twisted.spread import pb
|
||||
|
||||
from scrapy import log
|
||||
from scrapy.core.exceptions import NotConfigured
|
||||
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):
|
||||
self.procman = procman
|
||||
self.domain = domain
|
||||
self.logfile = logfile
|
||||
self.start_time = datetime.datetime.utcnow()
|
||||
self.status = "starting"
|
||||
self.pid = -1
|
||||
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, 'CLUSTER_CRAWLER_ENABLED': 1, 'WEBCONSOLE_ENABLED': 0})
|
||||
pickled_settings = pickle.dumps(self.scrapy_settings)
|
||||
self.env["SCRAPY_PICKLED_SETTINGS_TO_OVERRIDE"] = pickled_settings
|
||||
self.env["PYTHONPATH"] = ":".join(sys.path)#this is need so this crawl process knows where to locate local_scrapy_settings.
|
||||
|
||||
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}
|
||||
|
||||
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"
|
||||
self.procman.update_master(self.domain, "running")
|
||||
|
||||
def processEnded(self, reason):
|
||||
if settings.getbool('CLUSTER_WORKER_GZIP_LOGS'):
|
||||
try:
|
||||
f_in = open(self.logfile)
|
||||
f_out = gzip.open("%s.gz" % self.logfile, "wb")
|
||||
f_out.writelines(f_in)
|
||||
f_in.close()
|
||||
f_out.close()
|
||||
os.remove(self.logfile)
|
||||
self.logfile = "%s.gz" % self.logfile
|
||||
except Exception, e:
|
||||
log.msg("failed to compress %s exception=%s (domain=%s, pid=%s)" % (self.logfile, e, self.domain, self.pid))
|
||||
log.msg("ClusterWorker: finished domain=%s, pid=%d, log=%s" % (self.domain, self.pid, self.logfile))
|
||||
log.msg("Reason type: %s. value: %s" % (reason.type, reason.value) )
|
||||
del self.procman.running[self.domain]
|
||||
del self.procman.crawlers[self.pid]
|
||||
self.procman.update_master(self.domain, "scraped")
|
||||
|
||||
class ClusterWorker(pb.Root):
|
||||
|
||||
def __init__(self):
|
||||
if not settings.getbool('CLUSTER_WORKER_ENABLED'):
|
||||
raise NotConfigured
|
||||
|
||||
self.maxproc = settings.getint('CLUSTER_WORKER_MAXPROC')
|
||||
self.logdir = settings['CLUSTER_LOGDIR']
|
||||
self.running = {}#a dict domain->ScrapyProcessControl
|
||||
self.crawlers = {}#a dict pid->scrapy process remote pb connection
|
||||
self.starttime = datetime.datetime.utcnow()
|
||||
port = settings.getint('CLUSTER_WORKER_PORT')
|
||||
scrapyengine.listenTCP(port, pb.PBServerFactory(self))
|
||||
log.msg("PYTHONPATH: %s" % repr(sys.path))
|
||||
|
||||
def status(self, rcode=0, rstring=None):
|
||||
status = {}
|
||||
status["running"] = [ self.running[k].as_dict() for k in self.running.keys() ]
|
||||
status["starttime"] = self.starttime
|
||||
status["timestamp"] = datetime.datetime.utcnow()
|
||||
status["maxproc"] = self.maxproc
|
||||
status["loadavg"] = os.getloadavg()
|
||||
status["logdir"] = self.logdir
|
||||
status["callresponse"] = (rcode, rstring) if rstring else (0, "Status Response.")
|
||||
return status
|
||||
|
||||
def update_master(self, domain, domain_status):
|
||||
try:
|
||||
deferred = self.__master.callRemote("update", self.status(), domain, domain_status)
|
||||
except pb.DeadReferenceError:
|
||||
self.__master = None
|
||||
log.msg("Lost connection to node %s." % (self.name), log.ERROR)
|
||||
else:
|
||||
deferred.addCallbacks(callback=lambda x: x, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
|
||||
def remote_set_master(self, master):
|
||||
self.__master = master
|
||||
return self.status()
|
||||
|
||||
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))
|
||||
d = self.crawlers[proc.pid].callRemote("stop")
|
||||
def _close():
|
||||
proc.status = "closing"
|
||||
d.addCallbacks(callback=_close, errback=lambda reason: log.msg(reason, log.ERROR))
|
||||
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 remote_run(self, domain, spider_settings=None):
|
||||
"""Spawn process to run the given domain."""
|
||||
if len(self.running) < self.maxproc:
|
||||
if not domain in self.running:
|
||||
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]
|
||||
self.running[domain] = scrapy_proc
|
||||
try:
|
||||
import pysvn
|
||||
c = pysvn.Client()
|
||||
r = c.update(settings.get("CLUSTER_WORKER_SVNWORKDIR", "."))
|
||||
log.msg("Updated to revision %s." %r[0].number, level=log.DEBUG)
|
||||
except pysvn.ClientError, e:
|
||||
log.msg("Unable to svn update: %s" % e, level=log.WARNING)
|
||||
except ImportError:
|
||||
log.msg("pysvn module not available.", level=log.WARNING)
|
||||
proc = reactor.spawnProcess(scrapy_proc, sys.executable, args=args, env=scrapy_proc.env)
|
||||
return self.status(0, "Started process %s." % scrapy_proc)
|
||||
return self.status(2, "Domain %s already running." % domain )
|
||||
return self.status(1, "No free slot to run another process.")
|
||||
|
||||
def remote_register_crawler(self, pid, crawler):
|
||||
self.crawlers[pid] = crawler
|
||||
|
|
@ -1,25 +0,0 @@
|
|||
#!/usr/bin/python2.5
|
||||
|
||||
from twisted.spread import pb
|
||||
from twisted.internet import reactor
|
||||
from twisted.python import util
|
||||
import sys
|
||||
|
||||
factory = pb.PBClientFactory()
|
||||
reactor.connectTCP("localhost", 8789, factory)
|
||||
d = factory.getRootObject()
|
||||
|
||||
sys.argv.pop(0)
|
||||
|
||||
if not sys.argv:
|
||||
d.addCallback(lambda object: object.callRemote("status"))
|
||||
elif sys.argv[0] == "-s":
|
||||
d.addCallback(lambda object: object.callRemote("stop", sys.argv[1]))
|
||||
elif sys.argv[0] == "-r":
|
||||
d.addCallback(lambda object: object.callRemote("run", sys.argv[1]))
|
||||
elif sys.argv[0] == "-t":
|
||||
d.addCallback(lambda object: object.callRemote("statistics"))
|
||||
|
||||
d.addCallbacks(callback = util.println, errback = lambda reason: 'error: '+str(reason.value))
|
||||
d.addCallback(lambda _: reactor.stop())
|
||||
reactor.run()
|
||||
|
|
@ -1,3 +0,0 @@
|
|||
from scrapy.contrib_exp.cluster.worker.manager import ClusterWorker
|
||||
from scrapy.contrib_exp.cluster.master.web import ClusterMasterWeb
|
||||
from scrapy.contrib_exp.cluster.crawler.manager import ClusterCrawler
|
||||
|
|
@ -1,51 +0,0 @@
|
|||
Cluster Webservice API
|
||||
======================
|
||||
|
||||
The webservice API is available at
|
||||
|
||||
http://server:port/cluster_master/ws/
|
||||
|
||||
With no parameters, webservice returns the cluster status.
|
||||
|
||||
Query parameters
|
||||
================
|
||||
|
||||
- `format`: the answer format. By default, format=json. Other formats: pprint, pickle.
|
||||
|
||||
- `schedule`: schedules a comma separated list of domains. Schedule function takes optional parameters:
|
||||
|
||||
"priority": sets the queue priority for the specified domains (an integer). The default is setted by "DEFAULT_PRIORITY"
|
||||
setting (20 if not given). A lower priority number implies more priority.
|
||||
|
||||
"settings": run settings for the specified domains. 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)
|
||||
|
||||
- `disable_node`: disables a node so no more domains will be loaded in it until enabled again (but it will finish to run the running domains)
|
||||
|
||||
- `enable_node`: revert the state setted by 'disable_node'
|
||||
|
||||
- `verbosity`: sets the output verbosity level (1 is the default minimal, 2 includes domain settings, 0 disables output)
|
||||
|
||||
- `statistics`: shows the pending/running/scraped/lost statistics
|
||||
|
||||
Examples:
|
||||
---------
|
||||
|
||||
1) Schedule argos.co.uk, diy.com, littlewoodsdirect.com spiders, with priority=0, 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=0&settings=UNAVAILABLES_NOTIFY=2,UNAVAILABLES_DAYS_BACK=3
|
||||
|
||||
2) Get status with pprint format:
|
||||
|
||||
http://localhost:8080/cluster_master/ws/?format=pprint
|
||||
|
||||
3) Remove from pending lists domains argos.co.uk and diy.com. Answer with pprint format:
|
||||
|
||||
http://localhost:8080/cluster_master/ws/?remove=argos.co.uk,diy.com
|
||||
|
||||
4) Stop running domain littlewoodsdirect.com:
|
||||
|
||||
http://localhost:8080/cluster_master/ws/?stop=littlewoodsdirect.com
|
||||
|
|
@ -1,80 +0,0 @@
|
|||
#!/usr/bin/env python
|
||||
"""
|
||||
Cluster control script
|
||||
"""
|
||||
|
||||
from optparse import OptionParser
|
||||
import urllib
|
||||
|
||||
def main():
|
||||
parser = OptionParser(usage="Usage: scrapy-cluster-ctl.py [domain [domain [...]]] [options]" )
|
||||
parser.add_option("--disablenode", dest="disable_node", help="Disable given node (by name) so it will no accept more run requests.")
|
||||
parser.add_option("--enablenode", dest="enable_node", help="Enable given node (by name) so it will accept again run requests")
|
||||
parser.add_option("--format", dest="format", help="Output format. Default: pprint.", default="pprint")
|
||||
parser.add_option("--list", metavar="FILE", dest="list", help="Specify a file from where to read domains, one per line.")
|
||||
parser.add_option("--now", action="store_true", dest="now", help="Schedule domains to run with priority now.")
|
||||
parser.add_option("--output", metavar="FILE", dest="output", help="Output file. If not given, output to stdout.")
|
||||
parser.add_option("--port", dest="port", type="int", help="Cluster master port. Default: 8060.", default=8060)
|
||||
parser.add_option("--remove", dest="remove", action="store_true", help="Remove from schedule domains given as args.")
|
||||
parser.add_option("--schedule", dest="schedule", action="store_true", help="Schedule domains given as args.")
|
||||
parser.add_option("--server", dest="server", help="Cluster master server name. Default: localhost.", default="localhost")
|
||||
parser.add_option("--status", dest="status", action="store_true", help="Print cluster master status and quit.")
|
||||
parser.add_option("--statistics", dest="statistics", action="store_true", help="Print cluster statistics")
|
||||
parser.add_option("--stop", dest="stop", action="store_true", help="Stops a running domain.")
|
||||
parser.add_option("--verbosity", dest="verbosity", type="int", help="Sets the report status verbosity.")
|
||||
(opts, args) = parser.parse_args()
|
||||
|
||||
output = ""
|
||||
domains = []
|
||||
urlstring = "http://%s:%s/cluster_master/ws/" % (opts.server, opts.port)
|
||||
post = {"format":opts.format}
|
||||
if isinstance(opts.verbosity, int):
|
||||
post["verbosity"] = opts.verbosity
|
||||
|
||||
if args:
|
||||
domains = ",".join(args)
|
||||
elif opts.list:
|
||||
try:
|
||||
domainlist = []
|
||||
for d in open(opts.list, "r").readlines():
|
||||
domainlist.append(d.strip())
|
||||
domains = ",".join(domainlist)
|
||||
except IOError:
|
||||
print "Can't open file %s" % opts.list
|
||||
|
||||
if opts.status:
|
||||
pass
|
||||
elif opts.statistics:
|
||||
post["statistics"] = True
|
||||
elif opts.schedule and domains:
|
||||
post["schedule"] = domains
|
||||
if opts.now:
|
||||
post["priority"] = "0"
|
||||
post["settings"] = "UNAVAILABLES_NOTIFY=2"
|
||||
elif opts.remove and domains:
|
||||
post["remove"] = domains
|
||||
elif opts.stop and domains:
|
||||
post["stop"] = domains
|
||||
elif opts.disable_node:
|
||||
post["disable_node"] = opts.disable_node
|
||||
elif opts.enable_node:
|
||||
post["enable_node"] = opts.enable_node
|
||||
else:
|
||||
parser.print_help()
|
||||
return
|
||||
|
||||
f = urllib.urlopen(urlstring, urllib.urlencode(post))
|
||||
output=f.read()
|
||||
if not output:
|
||||
return
|
||||
if not opts.output:
|
||||
print output
|
||||
else:
|
||||
try:
|
||||
open(opts.output, "w").write(output)
|
||||
except IOError:
|
||||
open("/tmp/scrapy-cluster-schedule.tmp", "w").write(output)
|
||||
print "Could not open file %s for writing. Output dumped to /tmp/scrapy-cluster-schedule.tmp instead." % opts.output
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
Loading…
Reference in New Issue