mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 02:37:24 +08:00
fix: telegraf kakfa ordered configuration (#21548)
Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
@@ -2496,6 +2496,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") {
|
||||
@@ -2503,6 +2504,8 @@ func (h *SHostInfo) OnCatalogChanged(catalog mcclient.KeystoneServiceCatalogV3)
|
||||
} else {
|
||||
telegraf.BgReloadConf(conf)
|
||||
}
|
||||
} else {
|
||||
log.Infof("telegraf configuration no change")
|
||||
}
|
||||
|
||||
/*urls, _ = catalog.GetServiceURLs("elasticsearch",
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user