代码拉取完成,页面将自动刷新
package consumergroup
import (
"crypto/tls"
"github.com/elastic/beats/libbeat/common"
"github.com/elastic/beats/libbeat/common/cfgwarn"
"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) {
cfgwarn.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{String: "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
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。