diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 2b67343870..9983214e28 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -2408,6 +2408,7 @@ func (h *SHostInfo) OnCatalogChanged(catalog mcclient.KeystoneServiceCatalogV3) } } if !reflect.DeepEqual(telegraf.GetConf(), conf) || (!strings.Contains(svcs, "telegraf") && !telegraf.IsActive()) { + log.Infof("telegraf configuration change, to reload ...") log.Debugf("telegraf config: %s", conf) telegraf.SetConf(conf) if !strings.Contains(svcs, "telegraf") { @@ -2415,6 +2416,8 @@ func (h *SHostInfo) OnCatalogChanged(catalog mcclient.KeystoneServiceCatalogV3) } else { telegraf.BgReloadConf(conf) } + } else { + log.Infof("telegraf configuration no change") } /*urls, _ = catalog.GetServiceURLs("elasticsearch", diff --git a/pkg/hostman/system_service/telegraf.go b/pkg/hostman/system_service/telegraf.go index 48c19250f9..961742d357 100644 --- a/pkg/hostman/system_service/telegraf.go +++ b/pkg/hostman/system_service/telegraf.go @@ -121,17 +121,31 @@ func (s *STelegraf) GetConfig(kwargs map[string]interface{}) string { if kafka, ok := kwargs["kafka"]; ok { kafkaConf, _ := kafka.(map[string]interface{}) conf += "[[outputs.kafka]]\n" - for k := range kafkaConf { - if k == "brokers" { - brokers, _ := kafkaConf["brokers"].([]string) - for i := range brokers { - brokers[i] = fmt.Sprintf("\"%s\"", brokers[i]) + for _, k := range []string{ + "brokers", + "topic", + "sasl_username", + "sasl_password", + "sasl_mechanism", + } { + if val, ok := kafkaConf[k]; ok { + if k == "brokers" { + brokers, _ := val.([]string) + for i := range brokers { + brokers[i] = fmt.Sprintf("\"%s\"", brokers[i]) + } + conf += fmt.Sprintf(" brokers = [%s]\n", strings.Join(brokers, ", ")) + } else { + conf += fmt.Sprintf(" %s = \"%s\"\n", k, val) } - conf += fmt.Sprintf(" brokers = [%s]\n", strings.Join(brokers, ", ")) - } else { - conf += fmt.Sprintf(" %s = \"%s\"\n", k, kafkaConf[k]) } } + conf += " compression_codec = 0\n" + conf += " required_acks = -1\n" + conf += " max_retry = 3\n" + conf += " data_format = \"json\"\n" + conf += " json_timestamp_units = \"1ms\"\n" + conf += " routing_tag = \"host\"\n" conf += "\n" }