mirror of https://github.com/scrapy/scrapy.git
Stats collectin: fixed race condition between stats persistance and population of stats on domain close
This commit is contained in:
parent
895c70e036
commit
884f0c878f
|
|
@ -1,8 +1,8 @@
|
|||
.. _topics-stats:
|
||||
|
||||
===============
|
||||
Stats Collector
|
||||
===============
|
||||
================
|
||||
Stats Collection
|
||||
================
|
||||
|
||||
Overview
|
||||
========
|
||||
|
|
@ -309,7 +309,8 @@ functionality:
|
|||
.. function:: stats_domain_closed(domain, reason, domain_stats)
|
||||
|
||||
Sent right after the stats domain is closed. You can use this signal to
|
||||
collect resources.
|
||||
collect resources, but not to add any more stats as the stats domain has
|
||||
already been close (use :signal:`stats_domain_closing` for that instead).
|
||||
|
||||
:param domain: the stats domain just closed
|
||||
:type domain: str
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ class StatsCollector(object):
|
|||
def __init__(self):
|
||||
self._dump = settings.getbool('STATS_DUMP')
|
||||
self._stats = {None: {}} # None is for global stats
|
||||
dispatcher.connect(self._engine_stopped, signal=signals.engine_stopped)
|
||||
|
||||
def get_value(self, key, default=None, domain=None):
|
||||
return self._stats[domain].get(key, default)
|
||||
|
|
@ -60,11 +61,16 @@ class StatsCollector(object):
|
|||
if self._dump:
|
||||
log.msg("Dumping domain stats:\n" + pprint.pformat(stats), \
|
||||
domain=domain)
|
||||
self._persist_stats(stats, domain)
|
||||
|
||||
def engine_stopped(self):
|
||||
def _engine_stopped(self):
|
||||
stats = self.get_stats()
|
||||
if self._dump:
|
||||
log.msg("Dumping global stats:\n" + pprint.pformat(self.get_stats()))
|
||||
log.msg("Dumping global stats:\n" + pprint.pformat(stats))
|
||||
self._persist_stats(stats, domain=None)
|
||||
|
||||
def _persist_stats(self, stats, domain=None):
|
||||
pass
|
||||
|
||||
class MemoryStatsCollector(StatsCollector):
|
||||
|
||||
|
|
@ -72,9 +78,8 @@ class MemoryStatsCollector(StatsCollector):
|
|||
super(MemoryStatsCollector, self).__init__()
|
||||
self.domain_stats = {}
|
||||
|
||||
def close_domain(self, domain, reason):
|
||||
self.domain_stats[domain] = self._stats[domain]
|
||||
super(MemoryStatsCollector, self).close_domain(domain, reason)
|
||||
def _persist_stats(self, stats, domain=None):
|
||||
self.domain_stats[domain] = stats
|
||||
|
||||
|
||||
class DummyStatsCollector(StatsCollector):
|
||||
|
|
|
|||
|
|
@ -16,14 +16,16 @@ class MysqlStatsCollector(StatsCollector):
|
|||
mysqluri = settings['STATS_MYSQL_URI']
|
||||
self._mysql_conn = mysql_connect(mysqluri, use_unicode=False) if mysqluri else None
|
||||
|
||||
def close_domain(self, domain, reason):
|
||||
if self._mysql_conn:
|
||||
stored = datetime.utcnow()
|
||||
datas = pickle.dumps(self._stats[domain])
|
||||
table = 'domain_data_history'
|
||||
def _persist_stats(self, stats, domain=None):
|
||||
if domain is None: # only store domain-specific stats
|
||||
return
|
||||
if self._mysql_conn is None:
|
||||
return
|
||||
stored = datetime.utcnow()
|
||||
datas = pickle.dumps(stats)
|
||||
table = 'domain_data_history'
|
||||
|
||||
c = self._mysql_conn.cursor()
|
||||
c.execute("INSERT INTO %s (domain,stored,data) VALUES (%%s,%%s,%%s)" % table, \
|
||||
(domain, stored, datas))
|
||||
self._mysql_conn.commit()
|
||||
super(MysqlStatsCollector, self).close_domain(domain, reason)
|
||||
c = self._mysql_conn.cursor()
|
||||
c.execute("INSERT INTO %s (domain,stored,data) VALUES (%%s,%%s,%%s)" % table, \
|
||||
(domain, stored, datas))
|
||||
self._mysql_conn.commit()
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ Requires the boto library: http://code.google.com/p/boto/
|
|||
from datetime import datetime
|
||||
|
||||
import boto
|
||||
from twisted.internet import threads
|
||||
|
||||
from scrapy.stats.collector import StatsCollector
|
||||
from scrapy import log
|
||||
|
|
@ -21,16 +22,17 @@ class SimpledbStatsCollector(StatsCollector):
|
|||
sdb = boto.connect_sdb()
|
||||
sdb.create_domain(self._sdbdomain)
|
||||
|
||||
def close_domain(self, domain, reason):
|
||||
if self._sdbdomain:
|
||||
if self._async:
|
||||
from twisted.internet import threads
|
||||
dfd = threads.deferToThread(self._persist_to_sdb, domain, \
|
||||
self._stats[domain].copy())
|
||||
dfd.addErrback(log.err, 'Error uploading stats to SimpleDB', domain=domain)
|
||||
else:
|
||||
self._persist_to_sdb(domain, self._stats[domain])
|
||||
super(SimpledbStatsCollector, self).close_domain(domain, reason)
|
||||
def _persist_stats(self, stats, domain=None):
|
||||
if domain is None: # only store domain-specific stats
|
||||
return
|
||||
if not self._sdbdomain:
|
||||
return
|
||||
if self._async:
|
||||
dfd = threads.deferToThread(self._persist_to_sdb, domain, stats.copy())
|
||||
dfd.addErrback(log.err, 'Error uploading stats to SimpleDB', \
|
||||
domain=domain)
|
||||
else:
|
||||
self._persist_to_sdb(domain, stats)
|
||||
|
||||
def _persist_to_sdb(self, domain, stats):
|
||||
ts = datetime.utcnow().isoformat()
|
||||
|
|
|
|||
Loading…
Reference in New Issue