forked from grafana-cold-storage/metrictank
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathinput.go
More file actions
116 lines (98 loc) · 3.56 KB
/
Copy pathinput.go
File metadata and controls
116 lines (98 loc) · 3.56 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
// Package in provides interfaces, concrete implementations, and utilities
// to ingest data into metrictank
package input
import (
"fmt"
"time"
"gopkg.in/raintank/schema.v1"
"gopkg.in/raintank/schema.v1/msg"
"github.com/grafana/metrictank/idx"
"github.com/grafana/metrictank/mdata"
"github.com/grafana/metrictank/stats"
"github.com/raintank/worldping-api/pkg/log"
)
type Handler interface {
ProcessMetricData(md *schema.MetricData, partition int32)
ProcessMetricPoint(point schema.MetricPoint, format msg.Format, partition int32)
}
// TODO: clever way to document all metrics for all different inputs
// Default is a base handler for a metrics packet, aimed to be embedded by concrete implementations
type DefaultHandler struct {
receivedMD *stats.Counter32
receivedMP *stats.Counter32
receivedMPNO *stats.Counter32
invalidMD *stats.Counter32
invalidMP *stats.Counter32
unknownMP *stats.Counter32
pressureIdx *stats.Counter32
pressureTank *stats.Counter32
metrics mdata.Metrics
metricIndex idx.MetricIndex
}
func NewDefaultHandler(metrics mdata.Metrics, metricIndex idx.MetricIndex, input string) DefaultHandler {
return DefaultHandler{
receivedMD: stats.NewCounter32(fmt.Sprintf("input.%s.metricdata.received", input)),
receivedMP: stats.NewCounter32(fmt.Sprintf("input.%s.metricpoint.received", input)),
receivedMPNO: stats.NewCounter32(fmt.Sprintf("input.%s.metricpoint_no_org.received", input)),
invalidMD: stats.NewCounter32(fmt.Sprintf("input.%s.metricdata.invalid", input)),
invalidMP: stats.NewCounter32(fmt.Sprintf("input.%s.metricpoint.invalid", input)),
unknownMP: stats.NewCounter32(fmt.Sprintf("input.%s.metricpoint.unknown", input)),
pressureIdx: stats.NewCounter32(fmt.Sprintf("input.%s.pressure.idx", input)),
pressureTank: stats.NewCounter32(fmt.Sprintf("input.%s.pressure.tank", input)),
metrics: metrics,
metricIndex: metricIndex,
}
}
// ProcessMetricPoint updates the index if possible, and stores the data if we have an index entry
// concurrency-safe.
func (in DefaultHandler) ProcessMetricPoint(point schema.MetricPoint, format msg.Format, partition int32) {
if format == msg.FormatMetricPoint {
in.receivedMP.Inc()
} else {
in.receivedMPNO.Inc()
}
if !point.Valid() {
in.invalidMP.Inc()
log.Debug("in: Invalid metric %v", point)
return
}
pre := time.Now()
archive, _, ok := in.metricIndex.Update(point, partition)
in.pressureIdx.Add(int(time.Since(pre).Nanoseconds()))
if !ok {
in.unknownMP.Inc()
return
}
pre = time.Now()
m := in.metrics.GetOrCreate(point.MKey, archive.SchemaId, archive.AggId)
m.Add(point.Time, point.Value)
in.pressureTank.Add(int(time.Since(pre).Nanoseconds()))
}
// ProcessMetricData assures the data is stored and the metadata is in the index
// concurrency-safe.
func (in DefaultHandler) ProcessMetricData(md *schema.MetricData, partition int32) {
in.receivedMD.Inc()
err := md.Validate()
if err != nil {
in.invalidMD.Inc()
log.Debug("in: Invalid metric %v: %s", md, err)
return
}
if md.Time == 0 {
in.invalidMD.Inc()
log.Warn("in: invalid metric. metric.Time is 0. %s", md.Id)
return
}
mkey, err := schema.MKeyFromString(md.Id)
if err != nil {
log.Error(3, "in: Invalid metric %v: could not parse ID: %s", md, err)
return
}
pre := time.Now()
archive, _, _ := in.metricIndex.AddOrUpdate(mkey, md, partition)
in.pressureIdx.Add(int(time.Since(pre).Nanoseconds()))
pre = time.Now()
m := in.metrics.GetOrCreate(mkey, archive.SchemaId, archive.AggId)
m.Add(uint32(md.Time), md.Value)
in.pressureTank.Add(int(time.Since(pre).Nanoseconds()))
}