mirror of https://github.com/scrapy/scrapy.git
changes to logging and DEFAULT_PRIORITY removed
--HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40334
This commit is contained in:
parent
52596f350c
commit
874ac0c256
|
|
@ -14,8 +14,6 @@ from scrapy.core.exceptions import NotConfigured
|
|||
from scrapy.contrib.pbcluster.worker.manager import ResponseCode
|
||||
from scrapy.conf import settings
|
||||
|
||||
DEFAULT_PRIORITY = settings.getint("DEFAULT_PRIORITY", 20)
|
||||
|
||||
def my_import(name):
|
||||
mod = __import__(name)
|
||||
components = name.split('.')
|
||||
|
|
@ -35,9 +33,11 @@ class ClusterNodeBroker(pb.Referenceable):
|
|||
deferred = self._worker.callRemote("set_master", self)
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: lost connection to node %s." % (self.name), log.ERROR)
|
||||
log.msg("ClusterMaster: 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 _eb(failure):
|
||||
self._logfailure("Error while setting master to worker node", failure)
|
||||
deferred.addCallbacks(callback=self._set_status, errback=_eb)
|
||||
|
||||
def status_as_dict(self, verbosity=1):
|
||||
if verbosity == 0:
|
||||
|
|
@ -61,6 +61,85 @@ class ClusterNodeBroker(pb.Referenceable):
|
|||
status["loadavg"] = self.loadavg
|
||||
return status
|
||||
|
||||
def update_status(self):
|
||||
"""Update status from this worker. This is called periodically."""
|
||||
try:
|
||||
deferred = self._worker.callRemote("status")
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR)
|
||||
else:
|
||||
def _eb(failure):
|
||||
self._logfailure("Error while updating status", failure)
|
||||
deferred.addCallbacks(callback=self._set_status, errback=_eb)
|
||||
|
||||
def stop(self, domain):
|
||||
try:
|
||||
deferred = self._worker.callRemote("stop", domain)
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR)
|
||||
else:
|
||||
def _eb(failure):
|
||||
self._logfailure("Error while stopping domain=%s" % domain, failure)
|
||||
deferred.addCallbacks(callback=self._set_status, errback=_eb)
|
||||
|
||||
def run(self, domain_info):
|
||||
"""Run the given domain.
|
||||
|
||||
domain_info is a dict of keys:
|
||||
domain - the domain to run
|
||||
settings - the settings to use
|
||||
priority - the priority to use
|
||||
"""
|
||||
|
||||
domain = domain_info['domain']
|
||||
dsettings = domain_info['settings']
|
||||
priority = domain_info['priority']
|
||||
|
||||
def _run_errback(failure):
|
||||
self._logfailure("Error while running domain=%s" % domain, failure)
|
||||
self.master.loading.remove(domain)
|
||||
newprio = priority - 1 # increase priority for reschedule
|
||||
self.master.schedule([domain], dsettings, newprio)
|
||||
|
||||
def _run_callback(status):
|
||||
if status['callresponse'][0] == ResponseCode.NO_FREE_SLOT:
|
||||
log.msg("ClusterMaster: No available slots at worker=%s when trying to run domain=%s" % (self.name, domain), log.WARNING)
|
||||
self.master.loading.remove(domain)
|
||||
newprio = priority - 1 # increase priority for rerunning asap
|
||||
self.master.schedule([domain], dsettings, newprio)
|
||||
elif status['callresponse'][0] == ResponseCode.DOMAIN_ALREADY_RUNNING:
|
||||
log.msg("ClusterMaster: Already running domain=%s at worker=%s" % (domain, self.name), log.WARNING)
|
||||
self.master.loading.remove(domain)
|
||||
self.master.schedule([domain], dsettings, priority)
|
||||
|
||||
try:
|
||||
log.msg("ClusterMaster: Running domain=%s at worker=%s" % (domain, self.name), log.DEBUG)
|
||||
deferred = self._worker.callRemote("run", domain, dsettings)
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: Lost connection to worker=%s." % self.name, log.ERROR)
|
||||
else:
|
||||
deferred.addCallbacks(callback=_run_callback, errback=_run_errback)
|
||||
|
||||
def remote_update(self, worker_status, domain, domain_status):
|
||||
"""Called remotely form worker when domains finish to update status"""
|
||||
self._set_status(worker_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 in ("done", "terminated"):
|
||||
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)
|
||||
log.msg("ClusterMaster: Changed status to <%s> for domain=%s at worker=%s" % (domain_status, domain, self.name))
|
||||
|
||||
def _logfailure(self, msg, failure):
|
||||
log.msg("ClusterMaster: %s (worker=%s)\n%s" % (msg, self.name, failure), log.ERROR)
|
||||
|
||||
def _set_status(self, status):
|
||||
if not status:
|
||||
self.alive = False
|
||||
|
|
@ -85,72 +164,6 @@ class ClusterNodeBroker(pb.Referenceable):
|
|||
self.run(pending)
|
||||
self.master.loading.append(pending['domain'])
|
||||
|
||||
def update_status(self):
|
||||
try:
|
||||
deferred = self._worker.callRemote("status")
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: 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._worker.callRemote("stop", domain)
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: 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, domain_info):
|
||||
"""Run the given domain.
|
||||
|
||||
domain_info keys:
|
||||
domain - the domain to run
|
||||
settings - the settings to use
|
||||
priority - the priority to use
|
||||
"""
|
||||
|
||||
def _run_errback(reason):
|
||||
log.msg(reason, log.ERROR)
|
||||
self.master.loading.remove(domain_info['domain'])
|
||||
self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority'] - 1)
|
||||
log.msg("ClusterMaster: Domain %s rescheduled: lost connection to node." % domain_info['domain'], log.WARNING)
|
||||
|
||||
def _run_callback(status):
|
||||
if status['callresponse'][0] == ResponseCode.NO_FREE_SLOT:
|
||||
# slots are complete. Reschedule in master with priority reduced by one.
|
||||
# self.master.loading check should avoid this to happen
|
||||
self.master.loading.remove(domain_info['domain'])
|
||||
self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority'] - 1)
|
||||
log.msg("ClusterMaster: Domain %s rescheduled: no availble processes in worker" % domain_info['domain'], log.WARNING)
|
||||
elif status['callresponse'][0] == ResponseCode.DOMAIN_ALREADY_RUNNING:
|
||||
# domain already running in node. Reschedule with same priority.
|
||||
# self.master.loading check should avoid this to happen
|
||||
self.master.loading.remove(domain_info['domain'])
|
||||
self.master.schedule([domain_info['domain']], domain_info['settings'], domain_info['priority'])
|
||||
log.msg("ClusterMaster: Domain %s rescheduled: already running in node." % domain_info['domain'], log.WARNING)
|
||||
|
||||
try:
|
||||
deferred = self._worker.callRemote("run", domain_info["domain"], domain_info["settings"])
|
||||
except pb.DeadReferenceError:
|
||||
self._set_status(None)
|
||||
log.msg("ClusterMaster: 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):
|
||||
|
||||
|
|
@ -163,8 +176,9 @@ class ScrapyPBClientFactory(pb.PBClientFactory):
|
|||
|
||||
def clientConnectionLost(self, *args, **kargs):
|
||||
pb.PBClientFactory.clientConnectionLost(self, *args, **kargs)
|
||||
del self.master.nodes[self.nodename]
|
||||
log.msg("ClusterMaster: Lost connection to %s. Node removed" % self.nodename )
|
||||
self.master.remove_node(self.nodename)
|
||||
log.msg("ClusterMaster: Lost connection to worker=%s. Node removed" % self.nodename)
|
||||
|
||||
|
||||
class ClusterMaster(object):
|
||||
|
||||
|
|
@ -172,7 +186,9 @@ class ClusterMaster(object):
|
|||
|
||||
if not settings.getbool('CLUSTER_MASTER_ENABLED'):
|
||||
raise NotConfigured
|
||||
if not settings['CLUSTER_MASTER_STATEFILE']:
|
||||
|
||||
self.statefile = settings['CLUSTER_MASTER_STATEFILE']
|
||||
if not self.statefile:
|
||||
raise NotConfigured("ClusterMaster: Missing CLUSTER_MASTER_STATEFILE setting")
|
||||
|
||||
# import groups settings
|
||||
|
|
@ -182,12 +198,14 @@ class ClusterMaster(object):
|
|||
self.get_spider_groupsettings = lambda x: {}
|
||||
# load pending domains
|
||||
try:
|
||||
statefile = open(settings["CLUSTER_MASTER_STATEFILE"], "r")
|
||||
statefile = open(self.statefile, "r")
|
||||
self.pending = pickle.load(statefile)
|
||||
log.msg("ClusterMaster: Loaded state from %s" % self.statefile)
|
||||
except IOError:
|
||||
self.pending = []
|
||||
self.loading = []
|
||||
self.nodes = {}
|
||||
self.nodesconf = settings.get('CLUSTER_MASTER_NODES', {})
|
||||
self.start_time = datetime.datetime.utcnow()
|
||||
# for more info about statistics see self.update_nodes() and ClusterNodeBroker.remote_update()
|
||||
self.statistics = {"domains": {"running": set(), "scraped": {}, "lost_count": {}, "lost": set()}, "scraped_count": 0 }
|
||||
|
|
@ -201,34 +219,32 @@ class ClusterMaster(object):
|
|||
|
||||
def load_nodes(self):
|
||||
"""Loads nodes listed in CLUSTER_MASTER_NODES setting"""
|
||||
for name, hostport in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
|
||||
for name, hostport in self.nodesconf.iteritems():
|
||||
self.load_node(name, hostport)
|
||||
|
||||
def load_node(self, name, hostport):
|
||||
"""Creates the remote reference for a worker node"""
|
||||
server, port = hostport.split(":")
|
||||
port = int(port)
|
||||
log.msg("ClusterMaster: Connecting to worker %s (%s)..." % (name, hostport))
|
||||
log.msg("ClusterMaster: Connecting to worker=%s (%s)..." % (name, hostport))
|
||||
factory = ScrapyPBClientFactory(self, name)
|
||||
try:
|
||||
reactor.connectTCP(server, port, factory)
|
||||
except Exception, err:
|
||||
log.msg("ClusterMaster: Could not connect to worker %s (%s): %s" % (name, hostport, err), log.ERROR)
|
||||
log.msg("ClusterMaster: Could not connect to worker=%s (%s): %s" % (name, hostport, err), log.ERROR)
|
||||
else:
|
||||
def _errback(_reason):
|
||||
log.msg("ClusterMaster: Could not connect to worker %s (%s): %s" % (name, hostport, _reason), log.ERROR)
|
||||
def _eb(failure):
|
||||
self._logfailure("Error while loading worker node", failure)
|
||||
|
||||
d = factory.getRootObject()
|
||||
d.addCallbacks(callback=lambda obj: self.add_node(obj, name), errback=_errback)
|
||||
d.addCallbacks(callback=lambda obj: self.add_node(obj, name), errback=_eb)
|
||||
|
||||
def update_nodes(self):
|
||||
"""Update worker nodes statistics"""
|
||||
for name, hostport in settings.get('CLUSTER_MASTER_NODES', {}).iteritems():
|
||||
for name, hostport in self.nodesconf.iteritems():
|
||||
if name in self.nodes and self.nodes[name].alive:
|
||||
log.msg("ClusterMaster: Updating stats from worker node %s (%s)" % (name, hostport))
|
||||
self.nodes[name].update_status()
|
||||
else:
|
||||
log.msg("ClusterMaster: Reloading worker node %s (%s)" % (name, hostport))
|
||||
self.load_node(name, hostport)
|
||||
|
||||
real_running = set(self.running.keys())
|
||||
|
|
@ -241,7 +257,10 @@ class ClusterMaster(object):
|
|||
"""Add node given its node"""
|
||||
node = ClusterNodeBroker(cworker, name, self)
|
||||
self.nodes[name] = node
|
||||
log.msg("ClusterMaster: Added cluster worker %s" % name)
|
||||
log.msg("ClusterMaster: Added worker=%s" % name)
|
||||
|
||||
def remove_node(self, nodename):
|
||||
del self.nodes[nodename]
|
||||
|
||||
def disable_node(self, name):
|
||||
self.nodes[name].available = False
|
||||
|
|
@ -249,10 +268,7 @@ class ClusterMaster(object):
|
|||
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):
|
||||
def schedule(self, domains, spider_settings=None, priority=20):
|
||||
"""Schedule the given domains, with the given priority"""
|
||||
insert_pos = len([p for p in self.pending if ['priority'] <= priority])
|
||||
for domain in domains:
|
||||
|
|
@ -324,9 +340,9 @@ class ClusterMaster(object):
|
|||
|
||||
def _engine_started(self):
|
||||
self.load_nodes()
|
||||
scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL'))
|
||||
scrapyengine.addtask(self.update_nodes, settings.getint('CLUSTER_MASTER_POLL_INTERVAL', 60))
|
||||
|
||||
def _engine_stopped(self):
|
||||
with open(settings["CLUSTER_MASTER_STATEFILE"], "w") as f:
|
||||
with open(self.statefile, "w") as f:
|
||||
pickle.dump(self.pending, f)
|
||||
log.msg("ClusterMaster: state saved in %s" % settings["CLUSTER_MASTER_STATEFILE"])
|
||||
log.msg("ClusterMaster: Saved state in %s" % self.statefile)
|
||||
|
|
|
|||
|
|
@ -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.pbcluster.master.manager import ClusterMaster, DEFAULT_PRIORITY
|
||||
from scrapy.contrib.pbcluster.master.manager import ClusterMaster
|
||||
from scrapy.utils.serialization import serialize
|
||||
|
||||
class ClusterMasterWeb(ClusterMaster):
|
||||
|
|
@ -61,7 +61,7 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
else:
|
||||
sep = "\r"
|
||||
domains = args["schedule"]
|
||||
priority = int(args.get("priority", [DEFAULT_PRIORITY])[0])
|
||||
priority = int(args.get("priority", [20])[0])
|
||||
|
||||
# spider settings
|
||||
slist = args.get("settings", [""])[0].split(sep)
|
||||
|
|
@ -164,7 +164,7 @@ class ClusterMasterWeb(ClusterMaster):
|
|||
s += "<br />\n"
|
||||
|
||||
s += "Priority:<br />\n"
|
||||
s += "<input type='text' name='priority'>%s</input>" % DEFAULT_PRIORITY
|
||||
s += "<input type='text' name='priority'>%s</input>" % 20
|
||||
s += "<br />\n"
|
||||
|
||||
# spider settings
|
||||
|
|
|
|||
Loading…
Reference in New Issue