mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 18:02:24 +00:00
metrics: fix meter
This commit is contained in:
parent
05d2af672e
commit
754c0c8211
21 changed files with 96 additions and 123 deletions
|
|
@ -114,8 +114,8 @@ type freezerTable struct {
|
||||||
tailId uint32 // number of the earliest file
|
tailId uint32 // number of the earliest file
|
||||||
|
|
||||||
headBytes int64 // Number of bytes written to the head file
|
headBytes int64 // Number of bytes written to the head file
|
||||||
readMeter metrics.Meter // Meter for measuring the effective amount of data read
|
readMeter *metrics.Meter // Meter for measuring the effective amount of data read
|
||||||
writeMeter metrics.Meter // Meter for measuring the effective amount of data written
|
writeMeter *metrics.Meter // Meter for measuring the effective amount of data written
|
||||||
sizeGauge *metrics.Gauge // Gauge for tracking the combined size of all freezer tables
|
sizeGauge *metrics.Gauge // Gauge for tracking the combined size of all freezer tables
|
||||||
|
|
||||||
logger log.Logger // Logger with database path and table name embedded
|
logger log.Logger // Logger with database path and table name embedded
|
||||||
|
|
@ -124,13 +124,13 @@ type freezerTable struct {
|
||||||
|
|
||||||
// newFreezerTable opens the given path as a freezer table.
|
// newFreezerTable opens the given path as a freezer table.
|
||||||
func newFreezerTable(path, name string, disableSnappy, readonly bool) (*freezerTable, error) {
|
func newFreezerTable(path, name string, disableSnappy, readonly bool) (*freezerTable, error) {
|
||||||
return newTable(path, name, metrics.NilMeter{}, metrics.NilMeter{}, metrics.NewGauge(), freezerTableSize, disableSnappy, readonly)
|
return newTable(path, name, metrics.NewInactiveMeter(), metrics.NewInactiveMeter(), metrics.NewGauge(), freezerTableSize, disableSnappy, readonly)
|
||||||
}
|
}
|
||||||
|
|
||||||
// newTable opens a freezer table, creating the data and index files if they are
|
// newTable opens a freezer table, creating the data and index files if they are
|
||||||
// non-existent. Both files are truncated to the shortest common length to ensure
|
// non-existent. Both files are truncated to the shortest common length to ensure
|
||||||
// they don't go out of sync.
|
// they don't go out of sync.
|
||||||
func newTable(path string, name string, readMeter metrics.Meter, writeMeter metrics.Meter, sizeGauge *metrics.Gauge, maxFilesize uint32, noCompression, readonly bool) (*freezerTable, error) {
|
func newTable(path string, name string, readMeter, writeMeter *metrics.Meter, sizeGauge *metrics.Gauge, maxFilesize uint32, noCompression, readonly bool) (*freezerTable, error) {
|
||||||
// Ensure the containing directory exists and open the indexEntry file
|
// Ensure the containing directory exists and open the indexEntry file
|
||||||
if err := os.MkdirAll(path, 0755); err != nil {
|
if err := os.MkdirAll(path, 0755); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|
|
||||||
|
|
@ -47,21 +47,21 @@ type triePrefetcher struct {
|
||||||
term chan struct{} // Channel to signal interruption
|
term chan struct{} // Channel to signal interruption
|
||||||
noreads bool // Whether to ignore state-read-only prefetch requests
|
noreads bool // Whether to ignore state-read-only prefetch requests
|
||||||
|
|
||||||
deliveryMissMeter metrics.Meter
|
deliveryMissMeter *metrics.Meter
|
||||||
|
|
||||||
accountLoadReadMeter metrics.Meter
|
accountLoadReadMeter *metrics.Meter
|
||||||
accountLoadWriteMeter metrics.Meter
|
accountLoadWriteMeter *metrics.Meter
|
||||||
accountDupReadMeter metrics.Meter
|
accountDupReadMeter *metrics.Meter
|
||||||
accountDupWriteMeter metrics.Meter
|
accountDupWriteMeter *metrics.Meter
|
||||||
accountDupCrossMeter metrics.Meter
|
accountDupCrossMeter *metrics.Meter
|
||||||
accountWasteMeter metrics.Meter
|
accountWasteMeter *metrics.Meter
|
||||||
|
|
||||||
storageLoadReadMeter metrics.Meter
|
storageLoadReadMeter *metrics.Meter
|
||||||
storageLoadWriteMeter metrics.Meter
|
storageLoadWriteMeter *metrics.Meter
|
||||||
storageDupReadMeter metrics.Meter
|
storageDupReadMeter *metrics.Meter
|
||||||
storageDupWriteMeter metrics.Meter
|
storageDupWriteMeter *metrics.Meter
|
||||||
storageDupCrossMeter metrics.Meter
|
storageDupCrossMeter *metrics.Meter
|
||||||
storageWasteMeter metrics.Meter
|
storageWasteMeter *metrics.Meter
|
||||||
}
|
}
|
||||||
|
|
||||||
func newTriePrefetcher(db Database, root common.Hash, namespace string, noreads bool) *triePrefetcher {
|
func newTriePrefetcher(db Database, root common.Hash, namespace string, noreads bool) *triePrefetcher {
|
||||||
|
|
|
||||||
|
|
@ -881,7 +881,7 @@ func (q *queue) DeliverReceipts(id string, receiptList [][]*types.Receipt, recei
|
||||||
// to access the queue, so they already need a lock anyway.
|
// to access the queue, so they already need a lock anyway.
|
||||||
func (q *queue) deliver(id string, taskPool map[common.Hash]*types.Header,
|
func (q *queue) deliver(id string, taskPool map[common.Hash]*types.Header,
|
||||||
taskQueue *prque.Prque[int64, *types.Header], pendPool map[string]*fetchRequest,
|
taskQueue *prque.Prque[int64, *types.Header], pendPool map[string]*fetchRequest,
|
||||||
reqTimer metrics.Timer, resInMeter metrics.Meter, resDropMeter metrics.Meter,
|
reqTimer metrics.Timer, resInMeter, resDropMeter *metrics.Meter,
|
||||||
results int, validate func(index int, header *types.Header) error,
|
results int, validate func(index int, header *types.Header) error,
|
||||||
reconstruct func(index int, result *fetchResult)) (int, error) {
|
reconstruct func(index int, result *fetchResult)) (int, error) {
|
||||||
// Short circuit if the data was never requested
|
// Short circuit if the data was never requested
|
||||||
|
|
|
||||||
|
|
@ -41,23 +41,23 @@ func (h *bidirectionalMeters) get(ingress bool) *hsMeters {
|
||||||
type hsMeters struct {
|
type hsMeters struct {
|
||||||
// peerError measures the number of errors related to incorrect peer
|
// peerError measures the number of errors related to incorrect peer
|
||||||
// behaviour, such as invalid message code, size, encoding, etc.
|
// behaviour, such as invalid message code, size, encoding, etc.
|
||||||
peerError metrics.Meter
|
peerError *metrics.Meter
|
||||||
|
|
||||||
// timeoutError measures the number of timeouts.
|
// timeoutError measures the number of timeouts.
|
||||||
timeoutError metrics.Meter
|
timeoutError *metrics.Meter
|
||||||
|
|
||||||
// networkIDMismatch measures the number of network id mismatch errors.
|
// networkIDMismatch measures the number of network id mismatch errors.
|
||||||
networkIDMismatch metrics.Meter
|
networkIDMismatch *metrics.Meter
|
||||||
|
|
||||||
// protocolVersionMismatch measures the number of differing protocol
|
// protocolVersionMismatch measures the number of differing protocol
|
||||||
// versions.
|
// versions.
|
||||||
protocolVersionMismatch metrics.Meter
|
protocolVersionMismatch *metrics.Meter
|
||||||
|
|
||||||
// genesisMismatch measures the number of differing genesises.
|
// genesisMismatch measures the number of differing genesises.
|
||||||
genesisMismatch metrics.Meter
|
genesisMismatch *metrics.Meter
|
||||||
|
|
||||||
// forkidRejected measures the number of differing forkids.
|
// forkidRejected measures the number of differing forkids.
|
||||||
forkidRejected metrics.Meter
|
forkidRejected *metrics.Meter
|
||||||
}
|
}
|
||||||
|
|
||||||
// newHandshakeMeters registers and returns handshake meters for the given
|
// newHandshakeMeters registers and returns handshake meters for the given
|
||||||
|
|
|
||||||
|
|
@ -62,14 +62,14 @@ type Database struct {
|
||||||
fn string // filename for reporting
|
fn string // filename for reporting
|
||||||
db *leveldb.DB // LevelDB instance
|
db *leveldb.DB // LevelDB instance
|
||||||
|
|
||||||
compTimeMeter metrics.Meter // Meter for measuring the total time spent in database compaction
|
compTimeMeter *metrics.Meter // Meter for measuring the total time spent in database compaction
|
||||||
compReadMeter metrics.Meter // Meter for measuring the data read during compaction
|
compReadMeter *metrics.Meter // Meter for measuring the data read during compaction
|
||||||
compWriteMeter metrics.Meter // Meter for measuring the data written during compaction
|
compWriteMeter *metrics.Meter // Meter for measuring the data written during compaction
|
||||||
writeDelayNMeter metrics.Meter // Meter for measuring the write delay number due to database compaction
|
writeDelayNMeter *metrics.Meter // Meter for measuring the write delay number due to database compaction
|
||||||
writeDelayMeter metrics.Meter // Meter for measuring the write delay duration due to database compaction
|
writeDelayMeter *metrics.Meter // Meter for measuring the write delay duration due to database compaction
|
||||||
diskSizeGauge *metrics.Gauge // Gauge for tracking the size of all the levels in the database
|
diskSizeGauge *metrics.Gauge // Gauge for tracking the size of all the levels in the database
|
||||||
diskReadMeter metrics.Meter // Meter for measuring the effective amount of data read
|
diskReadMeter *metrics.Meter // Meter for measuring the effective amount of data read
|
||||||
diskWriteMeter metrics.Meter // Meter for measuring the effective amount of data written
|
diskWriteMeter *metrics.Meter // Meter for measuring the effective amount of data written
|
||||||
memCompGauge *metrics.Gauge // Gauge for tracking the number of memory compaction
|
memCompGauge *metrics.Gauge // Gauge for tracking the number of memory compaction
|
||||||
level0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in level0
|
level0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in level0
|
||||||
nonlevel0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in non0 level
|
nonlevel0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in non0 level
|
||||||
|
|
|
||||||
|
|
@ -58,14 +58,14 @@ type Database struct {
|
||||||
fn string // filename for reporting
|
fn string // filename for reporting
|
||||||
db *pebble.DB // Underlying pebble storage engine
|
db *pebble.DB // Underlying pebble storage engine
|
||||||
|
|
||||||
compTimeMeter metrics.Meter // Meter for measuring the total time spent in database compaction
|
compTimeMeter *metrics.Meter // Meter for measuring the total time spent in database compaction
|
||||||
compReadMeter metrics.Meter // Meter for measuring the data read during compaction
|
compReadMeter *metrics.Meter // Meter for measuring the data read during compaction
|
||||||
compWriteMeter metrics.Meter // Meter for measuring the data written during compaction
|
compWriteMeter *metrics.Meter // Meter for measuring the data written during compaction
|
||||||
writeDelayNMeter metrics.Meter // Meter for measuring the write delay number due to database compaction
|
writeDelayNMeter *metrics.Meter // Meter for measuring the write delay number due to database compaction
|
||||||
writeDelayMeter metrics.Meter // Meter for measuring the write delay duration due to database compaction
|
writeDelayMeter *metrics.Meter // Meter for measuring the write delay duration due to database compaction
|
||||||
diskSizeGauge *metrics.Gauge // Gauge for tracking the size of all the levels in the database
|
diskSizeGauge *metrics.Gauge // Gauge for tracking the size of all the levels in the database
|
||||||
diskReadMeter metrics.Meter // Meter for measuring the effective amount of data read
|
diskReadMeter *metrics.Meter // Meter for measuring the effective amount of data read
|
||||||
diskWriteMeter metrics.Meter // Meter for measuring the effective amount of data written
|
diskWriteMeter *metrics.Meter // Meter for measuring the effective amount of data written
|
||||||
memCompGauge *metrics.Gauge // Gauge for tracking the number of memory compaction
|
memCompGauge *metrics.Gauge // Gauge for tracking the number of memory compaction
|
||||||
level0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in level0
|
level0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in level0
|
||||||
nonlevel0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in non0 level
|
nonlevel0CompGauge *metrics.Gauge // Gauge for tracking the number of table compaction in non0 level
|
||||||
|
|
|
||||||
|
|
@ -146,7 +146,7 @@ func (exp *exp) publishHistogram(name string, metric metrics.Histogram) {
|
||||||
exp.getFloat(name + ".999-percentile").Set(ps[4])
|
exp.getFloat(name + ".999-percentile").Set(ps[4])
|
||||||
}
|
}
|
||||||
|
|
||||||
func (exp *exp) publishMeter(name string, metric metrics.Meter) {
|
func (exp *exp) publishMeter(name string, metric *metrics.Meter) {
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
exp.getInt(name + ".count").Set(m.Count())
|
exp.getInt(name + ".count").Set(m.Count())
|
||||||
exp.getFloat(name + ".one-minute").Set(m.Rate1())
|
exp.getFloat(name + ".one-minute").Set(m.Rate1())
|
||||||
|
|
@ -200,7 +200,7 @@ func (exp *exp) syncToExpvar() {
|
||||||
exp.publishGaugeInfo(name, i.Snapshot())
|
exp.publishGaugeInfo(name, i.Snapshot())
|
||||||
case metrics.Histogram:
|
case metrics.Histogram:
|
||||||
exp.publishHistogram(name, i)
|
exp.publishHistogram(name, i)
|
||||||
case metrics.Meter:
|
case *metrics.Meter:
|
||||||
exp.publishMeter(name, i)
|
exp.publishMeter(name, i)
|
||||||
case metrics.Timer:
|
case metrics.Timer:
|
||||||
exp.publishTimer(name, i)
|
exp.publishTimer(name, i)
|
||||||
|
|
|
||||||
|
|
@ -87,7 +87,7 @@ func graphite(c *GraphiteConfig) error {
|
||||||
key := strings.Replace(strconv.FormatFloat(psKey*100.0, 'f', -1, 64), ".", "", 1)
|
key := strings.Replace(strconv.FormatFloat(psKey*100.0, 'f', -1, 64), ".", "", 1)
|
||||||
fmt.Fprintf(w, "%s.%s.%s-percentile %.2f %d\n", c.Prefix, name, key, ps[psIdx], now)
|
fmt.Fprintf(w, "%s.%s.%s-percentile %.2f %d\n", c.Prefix, name, key, ps[psIdx], now)
|
||||||
}
|
}
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, m.Count(), now)
|
fmt.Fprintf(w, "%s.%s.count %d %d\n", c.Prefix, name, m.Count(), now)
|
||||||
fmt.Fprintf(w, "%s.%s.one-minute %.2f %d\n", c.Prefix, name, m.Rate1(), now)
|
fmt.Fprintf(w, "%s.%s.one-minute %.2f %d\n", c.Prefix, name, m.Rate1(), now)
|
||||||
|
|
|
||||||
|
|
@ -20,7 +20,6 @@ package metrics
|
||||||
var (
|
var (
|
||||||
_ SampleSnapshot = (*emptySnapshot)(nil)
|
_ SampleSnapshot = (*emptySnapshot)(nil)
|
||||||
_ HistogramSnapshot = (*emptySnapshot)(nil)
|
_ HistogramSnapshot = (*emptySnapshot)(nil)
|
||||||
_ MeterSnapshot = (*emptySnapshot)(nil)
|
|
||||||
_ TimerSnapshot = (*emptySnapshot)(nil)
|
_ TimerSnapshot = (*emptySnapshot)(nil)
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -62,7 +62,7 @@ func readMeter(namespace, name string, i interface{}) (string, map[string]interf
|
||||||
"p9999": ps[6],
|
"p9999": ps[6],
|
||||||
}
|
}
|
||||||
return measurement, fields
|
return measurement, fields
|
||||||
case metrics.Meter:
|
case *metrics.Meter:
|
||||||
ms := metric.Snapshot()
|
ms := metric.Snapshot()
|
||||||
measurement := fmt.Sprintf("%s%s.meter", namespace, name)
|
measurement := fmt.Sprintf("%s%s.meter", namespace, name)
|
||||||
fields := map[string]interface{}{
|
fields := map[string]interface{}{
|
||||||
|
|
|
||||||
|
|
@ -50,7 +50,7 @@ func LogScaled(r Registry, freq time.Duration, scale time.Duration, l Logger) {
|
||||||
l.Printf(" 95%%: %12.2f\n", ps[2])
|
l.Printf(" 95%%: %12.2f\n", ps[2])
|
||||||
l.Printf(" 99%%: %12.2f\n", ps[3])
|
l.Printf(" 99%%: %12.2f\n", ps[3])
|
||||||
l.Printf(" 99.9%%: %12.2f\n", ps[4])
|
l.Printf(" 99.9%%: %12.2f\n", ps[4])
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
l.Printf("meter %s\n", name)
|
l.Printf("meter %s\n", name)
|
||||||
l.Printf(" count: %9d\n", m.Count())
|
l.Printf(" count: %9d\n", m.Count())
|
||||||
|
|
|
||||||
|
|
@ -7,40 +7,21 @@ import (
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
type MeterSnapshot interface {
|
|
||||||
Count() int64
|
|
||||||
Rate1() float64
|
|
||||||
Rate5() float64
|
|
||||||
Rate15() float64
|
|
||||||
RateMean() float64
|
|
||||||
}
|
|
||||||
|
|
||||||
// Meters count events to produce exponentially-weighted moving average rates
|
|
||||||
// at one-, five-, and fifteen-minutes and a mean rate.
|
|
||||||
type Meter interface {
|
|
||||||
Mark(int64)
|
|
||||||
Snapshot() MeterSnapshot
|
|
||||||
Stop()
|
|
||||||
}
|
|
||||||
|
|
||||||
// GetOrRegisterMeter returns an existing Meter or constructs and registers a
|
// GetOrRegisterMeter returns an existing Meter or constructs and registers a
|
||||||
// new StandardMeter.
|
// new Meter.
|
||||||
// Be sure to unregister the meter from the registry once it is of no use to
|
// Be sure to unregister the meter from the registry once it is of no use to
|
||||||
// allow for garbage collection.
|
// allow for garbage collection.
|
||||||
func GetOrRegisterMeter(name string, r Registry) Meter {
|
func GetOrRegisterMeter(name string, r Registry) *Meter {
|
||||||
if nil == r {
|
if r == nil {
|
||||||
r = DefaultRegistry
|
r = DefaultRegistry
|
||||||
}
|
}
|
||||||
return r.GetOrRegister(name, NewMeter).(Meter)
|
return r.GetOrRegister(name, NewMeter).(*Meter)
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewMeter constructs a new StandardMeter and launches a goroutine.
|
// NewMeter constructs a new Meter and launches a goroutine.
|
||||||
// Be sure to call Stop() once the meter is of no use to allow for garbage collection.
|
// Be sure to call Stop() once the meter is of no use to allow for garbage collection.
|
||||||
func NewMeter() Meter {
|
func NewMeter() *Meter {
|
||||||
if !Enabled {
|
m := newMeter()
|
||||||
return NilMeter{}
|
|
||||||
}
|
|
||||||
m := newStandardMeter()
|
|
||||||
arbiter.Lock()
|
arbiter.Lock()
|
||||||
defer arbiter.Unlock()
|
defer arbiter.Unlock()
|
||||||
arbiter.meters[m] = struct{}{}
|
arbiter.meters[m] = struct{}{}
|
||||||
|
|
@ -53,57 +34,46 @@ func NewMeter() Meter {
|
||||||
|
|
||||||
// NewInactiveMeter returns a meter but does not start any goroutines. This
|
// NewInactiveMeter returns a meter but does not start any goroutines. This
|
||||||
// method is mainly intended for testing.
|
// method is mainly intended for testing.
|
||||||
func NewInactiveMeter() Meter {
|
func NewInactiveMeter() *Meter {
|
||||||
if !Enabled {
|
return newMeter()
|
||||||
return NilMeter{}
|
|
||||||
}
|
|
||||||
m := newStandardMeter()
|
|
||||||
return m
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewRegisteredMeter constructs and registers a new StandardMeter
|
// NewRegisteredMeter constructs and registers a new Meter
|
||||||
// and launches a goroutine.
|
// and launches a goroutine.
|
||||||
// Be sure to unregister the meter from the registry once it is of no use to
|
// Be sure to unregister the meter from the registry once it is of no use to
|
||||||
// allow for garbage collection.
|
// allow for garbage collection.
|
||||||
func NewRegisteredMeter(name string, r Registry) Meter {
|
func NewRegisteredMeter(name string, r Registry) *Meter {
|
||||||
return GetOrRegisterMeter(name, r)
|
return GetOrRegisterMeter(name, r)
|
||||||
}
|
}
|
||||||
|
|
||||||
// meterSnapshot is a read-only copy of the meter's internal values.
|
// MeterSnapshot is a read-only copy of the meter's internal values.
|
||||||
type meterSnapshot struct {
|
type MeterSnapshot struct {
|
||||||
count int64
|
count int64
|
||||||
rate1, rate5, rate15, rateMean float64
|
rate1, rate5, rate15, rateMean float64
|
||||||
}
|
}
|
||||||
|
|
||||||
// Count returns the count of events at the time the snapshot was taken.
|
// Count returns the count of events at the time the snapshot was taken.
|
||||||
func (m *meterSnapshot) Count() int64 { return m.count }
|
func (m *MeterSnapshot) Count() int64 { return m.count }
|
||||||
|
|
||||||
// Rate1 returns the one-minute moving average rate of events per second at the
|
// Rate1 returns the one-minute moving average rate of events per second at the
|
||||||
// time the snapshot was taken.
|
// time the snapshot was taken.
|
||||||
func (m *meterSnapshot) Rate1() float64 { return m.rate1 }
|
func (m *MeterSnapshot) Rate1() float64 { return m.rate1 }
|
||||||
|
|
||||||
// Rate5 returns the five-minute moving average rate of events per second at
|
// Rate5 returns the five-minute moving average rate of events per second at
|
||||||
// the time the snapshot was taken.
|
// the time the snapshot was taken.
|
||||||
func (m *meterSnapshot) Rate5() float64 { return m.rate5 }
|
func (m *MeterSnapshot) Rate5() float64 { return m.rate5 }
|
||||||
|
|
||||||
// Rate15 returns the fifteen-minute moving average rate of events per second
|
// Rate15 returns the fifteen-minute moving average rate of events per second
|
||||||
// at the time the snapshot was taken.
|
// at the time the snapshot was taken.
|
||||||
func (m *meterSnapshot) Rate15() float64 { return m.rate15 }
|
func (m *MeterSnapshot) Rate15() float64 { return m.rate15 }
|
||||||
|
|
||||||
// RateMean returns the meter's mean rate of events per second at the time the
|
// RateMean returns the meter's mean rate of events per second at the time the
|
||||||
// snapshot was taken.
|
// snapshot was taken.
|
||||||
func (m *meterSnapshot) RateMean() float64 { return m.rateMean }
|
func (m *MeterSnapshot) RateMean() float64 { return m.rateMean }
|
||||||
|
|
||||||
// NilMeter is a no-op Meter.
|
// Meter count events to produce exponentially-weighted moving average rates
|
||||||
type NilMeter struct{}
|
// at one-, five-, and fifteen-minutes and a mean rate.
|
||||||
|
type Meter struct {
|
||||||
func (NilMeter) Count() int64 { return 0 }
|
|
||||||
func (NilMeter) Mark(n int64) {}
|
|
||||||
func (NilMeter) Snapshot() MeterSnapshot { return (*emptySnapshot)(nil) }
|
|
||||||
func (NilMeter) Stop() {}
|
|
||||||
|
|
||||||
// StandardMeter is the standard implementation of a Meter.
|
|
||||||
type StandardMeter struct {
|
|
||||||
count atomic.Int64
|
count atomic.Int64
|
||||||
uncounted atomic.Int64 // not yet added to the EWMAs
|
uncounted atomic.Int64 // not yet added to the EWMAs
|
||||||
rateMean atomic.Uint64
|
rateMean atomic.Uint64
|
||||||
|
|
@ -113,8 +83,8 @@ type StandardMeter struct {
|
||||||
stopped atomic.Bool
|
stopped atomic.Bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func newStandardMeter() *StandardMeter {
|
func newMeter() *Meter {
|
||||||
return &StandardMeter{
|
return &Meter{
|
||||||
a1: NewEWMA1(),
|
a1: NewEWMA1(),
|
||||||
a5: NewEWMA5(),
|
a5: NewEWMA5(),
|
||||||
a15: NewEWMA15(),
|
a15: NewEWMA15(),
|
||||||
|
|
@ -123,7 +93,7 @@ func newStandardMeter() *StandardMeter {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Stop stops the meter, Mark() will be a no-op if you use it after being stopped.
|
// Stop stops the meter, Mark() will be a no-op if you use it after being stopped.
|
||||||
func (m *StandardMeter) Stop() {
|
func (m *Meter) Stop() {
|
||||||
if stopped := m.stopped.Swap(true); !stopped {
|
if stopped := m.stopped.Swap(true); !stopped {
|
||||||
arbiter.Lock()
|
arbiter.Lock()
|
||||||
delete(arbiter.meters, m)
|
delete(arbiter.meters, m)
|
||||||
|
|
@ -132,13 +102,13 @@ func (m *StandardMeter) Stop() {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Mark records the occurrence of n events.
|
// Mark records the occurrence of n events.
|
||||||
func (m *StandardMeter) Mark(n int64) {
|
func (m *Meter) Mark(n int64) {
|
||||||
m.uncounted.Add(n)
|
m.uncounted.Add(n)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Snapshot returns a read-only copy of the meter.
|
// Snapshot returns a read-only copy of the meter.
|
||||||
func (m *StandardMeter) Snapshot() MeterSnapshot {
|
func (m *Meter) Snapshot() *MeterSnapshot {
|
||||||
return &meterSnapshot{
|
return &MeterSnapshot{
|
||||||
count: m.count.Load() + m.uncounted.Load(),
|
count: m.count.Load() + m.uncounted.Load(),
|
||||||
rate1: m.a1.Snapshot().Rate(),
|
rate1: m.a1.Snapshot().Rate(),
|
||||||
rate5: m.a5.Snapshot().Rate(),
|
rate5: m.a5.Snapshot().Rate(),
|
||||||
|
|
@ -147,7 +117,7 @@ func (m *StandardMeter) Snapshot() MeterSnapshot {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *StandardMeter) tick() {
|
func (m *Meter) tick() {
|
||||||
// Take the uncounted values, add to count
|
// Take the uncounted values, add to count
|
||||||
n := m.uncounted.Swap(0)
|
n := m.uncounted.Swap(0)
|
||||||
count := m.count.Add(n)
|
count := m.count.Add(n)
|
||||||
|
|
@ -167,11 +137,11 @@ func (m *StandardMeter) tick() {
|
||||||
type meterArbiter struct {
|
type meterArbiter struct {
|
||||||
sync.RWMutex
|
sync.RWMutex
|
||||||
started bool
|
started bool
|
||||||
meters map[*StandardMeter]struct{}
|
meters map[*Meter]struct{}
|
||||||
ticker *time.Ticker
|
ticker *time.Ticker
|
||||||
}
|
}
|
||||||
|
|
||||||
var arbiter = meterArbiter{ticker: time.NewTicker(5 * time.Second), meters: make(map[*StandardMeter]struct{})}
|
var arbiter = meterArbiter{ticker: time.NewTicker(5 * time.Second), meters: make(map[*Meter]struct{})}
|
||||||
|
|
||||||
// tick meters on the scheduled interval
|
// tick meters on the scheduled interval
|
||||||
func (ma *meterArbiter) tick() {
|
func (ma *meterArbiter) tick() {
|
||||||
|
|
|
||||||
|
|
@ -30,10 +30,10 @@ func TestGetOrRegisterMeter(t *testing.T) {
|
||||||
func TestMeterDecay(t *testing.T) {
|
func TestMeterDecay(t *testing.T) {
|
||||||
ma := meterArbiter{
|
ma := meterArbiter{
|
||||||
ticker: time.NewTicker(time.Millisecond),
|
ticker: time.NewTicker(time.Millisecond),
|
||||||
meters: make(map[*StandardMeter]struct{}),
|
meters: make(map[*Meter]struct{}),
|
||||||
}
|
}
|
||||||
defer ma.ticker.Stop()
|
defer ma.ticker.Stop()
|
||||||
m := newStandardMeter()
|
m := newMeter()
|
||||||
ma.meters[m] = struct{}{}
|
ma.meters[m] = struct{}{}
|
||||||
m.Mark(1)
|
m.Mark(1)
|
||||||
ma.tickMeters()
|
ma.tickMeters()
|
||||||
|
|
|
||||||
|
|
@ -87,7 +87,7 @@ func (c *OpenTSDBConfig) writeRegistry(w io.Writer, now int64, shortHostname str
|
||||||
fmt.Fprintf(w, "put %s.%s.95-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[2], shortHostname)
|
fmt.Fprintf(w, "put %s.%s.95-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[2], shortHostname)
|
||||||
fmt.Fprintf(w, "put %s.%s.99-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[3], shortHostname)
|
fmt.Fprintf(w, "put %s.%s.99-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[3], shortHostname)
|
||||||
fmt.Fprintf(w, "put %s.%s.999-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[4], shortHostname)
|
fmt.Fprintf(w, "put %s.%s.999-percentile %d %.2f host=%s\n", c.Prefix, name, now, ps[4], shortHostname)
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
fmt.Fprintf(w, "put %s.%s.count %d %d host=%s\n", c.Prefix, name, now, m.Count(), shortHostname)
|
fmt.Fprintf(w, "put %s.%s.count %d %d host=%s\n", c.Prefix, name, now, m.Count(), shortHostname)
|
||||||
fmt.Fprintf(w, "put %s.%s.one-minute %d %.2f host=%s\n", c.Prefix, name, now, m.Rate1(), shortHostname)
|
fmt.Fprintf(w, "put %s.%s.one-minute %d %.2f host=%s\n", c.Prefix, name, now, m.Rate1(), shortHostname)
|
||||||
|
|
|
||||||
|
|
@ -63,7 +63,7 @@ func (c *collector) Add(name string, i any) error {
|
||||||
c.addGaugeInfo(name, m.Snapshot())
|
c.addGaugeInfo(name, m.Snapshot())
|
||||||
case metrics.Histogram:
|
case metrics.Histogram:
|
||||||
c.addHistogram(name, m.Snapshot())
|
c.addHistogram(name, m.Snapshot())
|
||||||
case metrics.Meter:
|
case *metrics.Meter:
|
||||||
c.addMeter(name, m.Snapshot())
|
c.addMeter(name, m.Snapshot())
|
||||||
case metrics.Timer:
|
case metrics.Timer:
|
||||||
c.addTimer(name, m.Snapshot())
|
c.addTimer(name, m.Snapshot())
|
||||||
|
|
@ -106,7 +106,7 @@ func (c *collector) addHistogram(name string, m metrics.HistogramSnapshot) {
|
||||||
c.buff.WriteRune('\n')
|
c.buff.WriteRune('\n')
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *collector) addMeter(name string, m metrics.MeterSnapshot) {
|
func (c *collector) addMeter(name string, m *metrics.MeterSnapshot) {
|
||||||
c.writeGaugeCounter(name, m.Count())
|
c.writeGaugeCounter(name, m.Count())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -176,7 +176,7 @@ func (r *StandardRegistry) GetAll() map[string]map[string]interface{} {
|
||||||
values["95%"] = ps[2]
|
values["95%"] = ps[2]
|
||||||
values["99%"] = ps[3]
|
values["99%"] = ps[3]
|
||||||
values["99.9%"] = ps[4]
|
values["99.9%"] = ps[4]
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
values["count"] = m.Count()
|
values["count"] = m.Count()
|
||||||
values["1m.rate"] = m.Rate1()
|
values["1m.rate"] = m.Rate1()
|
||||||
|
|
@ -214,7 +214,7 @@ func (r *StandardRegistry) Unregister(name string) {
|
||||||
|
|
||||||
func (r *StandardRegistry) loadOrRegister(name string, i interface{}) (interface{}, bool, bool) {
|
func (r *StandardRegistry) loadOrRegister(name string, i interface{}) (interface{}, bool, bool) {
|
||||||
switch i.(type) {
|
switch i.(type) {
|
||||||
case *Counter, *CounterFloat64, *Gauge, *GaugeFloat64, *GaugeInfo, *Healthcheck, Histogram, Meter, Timer, ResettingTimer:
|
case *Counter, *CounterFloat64, *Gauge, *GaugeFloat64, *GaugeInfo, *Healthcheck, Histogram, *Meter, Timer, ResettingTimer:
|
||||||
default:
|
default:
|
||||||
return nil, false, false
|
return nil, false, false
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -45,7 +45,7 @@ func Syslog(r Registry, d time.Duration, w *syslog.Writer) {
|
||||||
ps[3],
|
ps[3],
|
||||||
ps[4],
|
ps[4],
|
||||||
))
|
))
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
w.Info(fmt.Sprintf(
|
w.Info(fmt.Sprintf(
|
||||||
"meter %s: count: %d 1-min: %.2f 5-min: %.2f 15-min: %.2f mean: %.2f",
|
"meter %s: count: %d 1-min: %.2f 5-min: %.2f 15-min: %.2f mean: %.2f",
|
||||||
|
|
|
||||||
|
|
@ -7,7 +7,11 @@ import (
|
||||||
|
|
||||||
type TimerSnapshot interface {
|
type TimerSnapshot interface {
|
||||||
HistogramSnapshot
|
HistogramSnapshot
|
||||||
MeterSnapshot
|
Count() int64
|
||||||
|
Rate1() float64
|
||||||
|
Rate5() float64
|
||||||
|
Rate15() float64
|
||||||
|
RateMean() float64
|
||||||
}
|
}
|
||||||
|
|
||||||
// Timer capture the duration and rate of events.
|
// Timer capture the duration and rate of events.
|
||||||
|
|
@ -32,7 +36,7 @@ func GetOrRegisterTimer(name string, r Registry) Timer {
|
||||||
|
|
||||||
// NewCustomTimer constructs a new StandardTimer from a Histogram and a Meter.
|
// NewCustomTimer constructs a new StandardTimer from a Histogram and a Meter.
|
||||||
// Be sure to call Stop() once the timer is of no use to allow for garbage collection.
|
// Be sure to call Stop() once the timer is of no use to allow for garbage collection.
|
||||||
func NewCustomTimer(h Histogram, m Meter) Timer {
|
func NewCustomTimer(h Histogram, m *Meter) Timer {
|
||||||
if !Enabled {
|
if !Enabled {
|
||||||
return NilTimer{}
|
return NilTimer{}
|
||||||
}
|
}
|
||||||
|
|
@ -80,7 +84,7 @@ func (NilTimer) UpdateSince(time.Time) {}
|
||||||
// and Meter.
|
// and Meter.
|
||||||
type StandardTimer struct {
|
type StandardTimer struct {
|
||||||
histogram Histogram
|
histogram Histogram
|
||||||
meter Meter
|
meter *Meter
|
||||||
mutex sync.Mutex
|
mutex sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -123,7 +127,7 @@ func (t *StandardTimer) UpdateSince(ts time.Time) {
|
||||||
// timerSnapshot is a read-only copy of another Timer.
|
// timerSnapshot is a read-only copy of another Timer.
|
||||||
type timerSnapshot struct {
|
type timerSnapshot struct {
|
||||||
histogram HistogramSnapshot
|
histogram HistogramSnapshot
|
||||||
meter MeterSnapshot
|
meter *MeterSnapshot
|
||||||
}
|
}
|
||||||
|
|
||||||
// Count returns the number of events recorded at the time the snapshot was
|
// Count returns the number of events recorded at the time the snapshot was
|
||||||
|
|
|
||||||
|
|
@ -59,7 +59,7 @@ func WriteOnce(r Registry, w io.Writer) {
|
||||||
fmt.Fprintf(w, " 95%%: %12.2f\n", ps[2])
|
fmt.Fprintf(w, " 95%%: %12.2f\n", ps[2])
|
||||||
fmt.Fprintf(w, " 99%%: %12.2f\n", ps[3])
|
fmt.Fprintf(w, " 99%%: %12.2f\n", ps[3])
|
||||||
fmt.Fprintf(w, " 99.9%%: %12.2f\n", ps[4])
|
fmt.Fprintf(w, " 99.9%%: %12.2f\n", ps[4])
|
||||||
case Meter:
|
case *Meter:
|
||||||
m := metric.Snapshot()
|
m := metric.Snapshot()
|
||||||
fmt.Fprintf(w, "meter %s\n", namedMetric.name)
|
fmt.Fprintf(w, "meter %s\n", namedMetric.name)
|
||||||
fmt.Fprintf(w, " count: %9d\n", m.Count())
|
fmt.Fprintf(w, " count: %9d\n", m.Count())
|
||||||
|
|
|
||||||
|
|
@ -45,11 +45,11 @@ var (
|
||||||
egressTrafficMeter = metrics.NewRegisteredMeter("p2p/egress", nil)
|
egressTrafficMeter = metrics.NewRegisteredMeter("p2p/egress", nil)
|
||||||
|
|
||||||
// general ingress/egress connection meters
|
// general ingress/egress connection meters
|
||||||
serveMeter metrics.Meter = metrics.NilMeter{}
|
serveMeter = metrics.NewInactiveMeter()
|
||||||
serveSuccessMeter metrics.Meter = metrics.NilMeter{}
|
serveSuccessMeter = metrics.NewInactiveMeter()
|
||||||
dialMeter metrics.Meter = metrics.NilMeter{}
|
dialMeter = metrics.NewInactiveMeter()
|
||||||
dialSuccessMeter metrics.Meter = metrics.NilMeter{}
|
dialSuccessMeter = metrics.NewInactiveMeter()
|
||||||
dialConnectionError metrics.Meter = metrics.NilMeter{}
|
dialConnectionError = metrics.NewInactiveMeter()
|
||||||
|
|
||||||
// handshake error meters
|
// handshake error meters
|
||||||
dialTooManyPeers = metrics.NewRegisteredMeter("p2p/dials/error/saturated", nil)
|
dialTooManyPeers = metrics.NewRegisteredMeter("p2p/dials/error/saturated", nil)
|
||||||
|
|
|
||||||
|
|
@ -38,7 +38,7 @@ func (c *counter) add(size int) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// report uploads the cached statistics to meters.
|
// report uploads the cached statistics to meters.
|
||||||
func (c *counter) report(count metrics.Meter, size metrics.Meter) {
|
func (c *counter) report(count, size *metrics.Meter) {
|
||||||
count.Mark(int64(c.n))
|
count.Mark(int64(c.n))
|
||||||
size.Mark(int64(c.size))
|
size.Mark(int64(c.size))
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue