forked from elastic/beats
-
Notifications
You must be signed in to change notification settings - Fork 0
/
consumergroup.go
101 lines (83 loc) · 2.25 KB
/
consumergroup.go
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
package consumergroup
import (
"crypto/tls"
"github.com/elastic/beats/libbeat/common"
"github.com/elastic/beats/libbeat/logp"
"github.com/elastic/beats/libbeat/outputs"
"github.com/elastic/beats/metricbeat/mb"
"github.com/elastic/beats/metricbeat/module/kafka"
)
// init registers the MetricSet with the central registry.
func init() {
if err := mb.Registry.AddMetricSet("kafka", "consumergroup", New); err != nil {
panic(err)
}
}
// MetricSet type defines all fields of the MetricSet
type MetricSet struct {
mb.BaseMetricSet
broker *kafka.Broker
topics nameSet
groups nameSet
}
type groupAssignment struct {
clientID string
memberID string
clientHost string
}
var debugf = logp.MakeDebug("kafka")
// New creates a new instance of the MetricSet.
func New(base mb.BaseMetricSet) (mb.MetricSet, error) {
logp.Warn("BETA: The kafka consumergroup metricset is beta")
config := defaultConfig
if err := base.Module().UnpackConfig(&config); err != nil {
return nil, err
}
var tls *tls.Config
tlsCfg, err := outputs.LoadTLSConfig(config.TLS)
if err != nil {
return nil, err
}
if tlsCfg != nil {
tls = tlsCfg.BuildModuleConfig("")
}
timeout := base.Module().Config().Timeout
cfg := kafka.BrokerSettings{
MatchID: true,
DialTimeout: timeout,
ReadTimeout: timeout,
ClientID: config.ClientID,
Retries: config.Retries,
Backoff: config.Backoff,
TLS: tls,
Username: config.Username,
Password: config.Password,
// consumer groups API requires at least 0.9.0.0
Version: kafka.Version{"0.9.0.0"},
}
return &MetricSet{
BaseMetricSet: base,
broker: kafka.NewBroker(base.Host(), cfg),
groups: makeNameSet(config.Groups...),
topics: makeNameSet(config.Topics...),
}, nil
}
func (m *MetricSet) Fetch() ([]common.MapStr, error) {
if err := m.broker.Connect(); err != nil {
logp.Err("broker connect failed: %v", err)
return nil, err
}
b := m.broker
defer b.Close()
brokerInfo := common.MapStr{
"id": b.ID(),
"address": b.AdvertisedAddr(),
}
var events []common.MapStr
emitEvent := func(event common.MapStr) {
event["broker"] = brokerInfo
events = append(events, event)
}
err := fetchGroupInfo(emitEvent, b, m.groups.pred(), m.topics.pred())
return events, err
}