Add global statistics
This commit is contained in:
parent
a60be980c5
commit
cbd8048d31
|
@ -170,7 +170,8 @@ func (nodes *Nodes) worker() {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Returns global statistics for InfluxDB
|
// Returns global statistics for InfluxDB
|
||||||
func (nodes *Nodes) GlobalStats() (result GlobalStats) {
|
func (nodes *Nodes) GlobalStats() (result *GlobalStats) {
|
||||||
|
result = &GlobalStats{}
|
||||||
nodes.Lock()
|
nodes.Lock()
|
||||||
for _, node := range nodes.List {
|
for _, node := range nodes.List {
|
||||||
if node.Flags.Online {
|
if node.Flags.Online {
|
||||||
|
|
|
@ -20,7 +20,7 @@ type Collector struct {
|
||||||
connection *net.UDPConn // UDP socket
|
connection *net.UDPConn // UDP socket
|
||||||
queue chan *Response // received responses
|
queue chan *Response // received responses
|
||||||
msgType reflect.Type
|
msgType reflect.Type
|
||||||
iface string // interface name for the multicast binding
|
multicastAddr string
|
||||||
db *database.DB
|
db *database.DB
|
||||||
nodes *models.Nodes
|
nodes *models.Nodes
|
||||||
// Ticker and stopper
|
// Ticker and stopper
|
||||||
|
@ -44,10 +44,9 @@ func NewCollector(db *database.DB, nodes *models.Nodes, interval time.Duration,
|
||||||
conn.SetReadBuffer(maxDataGramSize)
|
conn.SetReadBuffer(maxDataGramSize)
|
||||||
|
|
||||||
collector := &Collector{
|
collector := &Collector{
|
||||||
CollectType: "nodeinfo statistics neighbours",
|
|
||||||
connection: conn,
|
connection: conn,
|
||||||
nodes: nodes,
|
nodes: nodes,
|
||||||
iface: iface,
|
multicastAddr: net.JoinHostPort(multiCastGroup+"%"+iface, port),
|
||||||
queue: make(chan *Response, 400),
|
queue: make(chan *Response, 400),
|
||||||
ticker: time.NewTicker(interval),
|
ticker: time.NewTicker(interval),
|
||||||
stop: make(chan interface{}, 1),
|
stop: make(chan interface{}, 1),
|
||||||
|
@ -69,15 +68,14 @@ func NewCollector(db *database.DB, nodes *models.Nodes, interval time.Duration,
|
||||||
func (coll *Collector) Close() {
|
func (coll *Collector) Close() {
|
||||||
// stop ticker
|
// stop ticker
|
||||||
coll.ticker.Stop()
|
coll.ticker.Stop()
|
||||||
coll.stop <- nil
|
close(coll.stop)
|
||||||
|
|
||||||
coll.connection.Close()
|
coll.connection.Close()
|
||||||
close(coll.queue)
|
close(coll.queue)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (coll *Collector) sendOnce() {
|
func (coll *Collector) sendOnce() {
|
||||||
coll.sendPacket(net.JoinHostPort(multiCastGroup+"%"+coll.iface, port))
|
coll.sendPacket(coll.multicastAddr)
|
||||||
log.Println("request", coll.CollectType)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (coll *Collector) sendPacket(address string) {
|
func (coll *Collector) sendPacket(address string) {
|
||||||
|
@ -86,7 +84,7 @@ func (coll *Collector) sendPacket(address string) {
|
||||||
log.Panic(err)
|
log.Panic(err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if _, err := coll.connection.WriteToUDP([]byte("GET "+coll.CollectType), addr); err != nil {
|
if _, err := coll.connection.WriteToUDP([]byte("GET nodeinfo statistics neighbours"), addr); err != nil {
|
||||||
log.Println("WriteToUDP failed:", err)
|
log.Println("WriteToUDP failed:", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
@ -98,6 +96,10 @@ func (coll *Collector) sender() {
|
||||||
case <-coll.stop:
|
case <-coll.stop:
|
||||||
return
|
return
|
||||||
case <-coll.ticker.C:
|
case <-coll.ticker.C:
|
||||||
|
// save global statistics
|
||||||
|
coll.db.AddPoint(database.MeasurementGlobal, nil, coll.nodes.GlobalStats().Fields(), time.Now())
|
||||||
|
|
||||||
|
// send the multicast packet to request per-node statistics
|
||||||
coll.sendOnce()
|
coll.sendOnce()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in New Issue