mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #12213 from ioito/hotfix/qx-apsara-monitor
fix(cloudmon): batch send data
This commit is contained in:
@@ -17,10 +17,12 @@ package apsaramon
|
||||
import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
|
||||
api "yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"yunion.io/x/onecloud/pkg/cloudmon/collectors/common"
|
||||
@@ -44,45 +46,65 @@ func (self *SApsaraCloudReport) collectRegionMetricOfHost(region cloudprovider.I
|
||||
aliReg := region.(*apsara.SRegion)
|
||||
since, until, err := common.TimeRangeFromArgs(self.Args)
|
||||
if err != nil {
|
||||
return err
|
||||
return errors.Wrapf(err, "common.TimeRangeFromArgs")
|
||||
}
|
||||
for metricName, influxDbSpecs := range apsaraMetricSpecs {
|
||||
metricArray, err := aliReg.FetchMetricData(metricName, "acs_ecs_dashboard", since, until)
|
||||
if err != nil {
|
||||
log.Errorln(err)
|
||||
|
||||
instances := []api.ServerDetails{}
|
||||
instanceMaps := map[string]api.ServerDetails{}
|
||||
jsonutils.Update(&instances, servers)
|
||||
for i := range instances {
|
||||
if len(instances[i].ExternalId) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
instances := []api.ServerDetails{}
|
||||
instanceMaps := map[string]api.ServerDetails{}
|
||||
jsonutils.Update(&instances, servers)
|
||||
for i := range instances {
|
||||
if len(instances[i].ExternalId) == 0 {
|
||||
continue
|
||||
}
|
||||
instanceMaps[instances[i].ExternalId] = instances[i]
|
||||
}
|
||||
|
||||
metrics := []SApsaraMetric{}
|
||||
jsonutils.Update(&metrics, metricArray)
|
||||
|
||||
for _, rtnMetric := range metrics {
|
||||
server, ok := instanceMaps[rtnMetric.InstanceId]
|
||||
if ok {
|
||||
metric, err := common.FillVMCapacity(jsonutils.Marshal(server).(*jsonutils.JSONDict))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
dataList = append(dataList, metric)
|
||||
serverMetric, err := self.collectMetricFromThisServer(jsonutils.Marshal(server), rtnMetric, influxDbSpecs)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
dataList = append(dataList, serverMetric)
|
||||
}
|
||||
instanceMaps[instances[i].ExternalId] = instances[i]
|
||||
metric, err := common.FillVMCapacity(jsonutils.Marshal(instances[i]).(*jsonutils.JSONDict))
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "common.FillVMCapacity")
|
||||
}
|
||||
dataList = append(dataList, metric)
|
||||
}
|
||||
return common.SendMetrics(self.Session, dataList, self.Args.Debug, "")
|
||||
|
||||
err = common.SendMetrics(self.Session, dataList, self.Args.Debug, "")
|
||||
if err != nil {
|
||||
log.Errorf("send server base metric error: %v", err)
|
||||
}
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(len(apsaraMetricSpecs))
|
||||
for _metricName, _influxDbSpecs := range apsaraMetricSpecs {
|
||||
go func(metricName string, influxDbSpecs []string) {
|
||||
defer wg.Done()
|
||||
dataList = []influxdb.SMetricData{}
|
||||
metricArray, err := aliReg.FetchMetricData(metricName, "acs_ecs_dashboard", since, until)
|
||||
if err != nil {
|
||||
log.Errorln(err)
|
||||
return
|
||||
}
|
||||
|
||||
metrics := []SApsaraMetric{}
|
||||
jsonutils.Update(&metrics, metricArray)
|
||||
|
||||
for _, rtnMetric := range metrics {
|
||||
server, ok := instanceMaps[rtnMetric.InstanceId]
|
||||
if ok {
|
||||
serverMetric, err := self.collectMetricFromThisServer(jsonutils.Marshal(server), rtnMetric, influxDbSpecs)
|
||||
if err != nil {
|
||||
log.Errorf("collect %s error: %v", metricName, err)
|
||||
continue
|
||||
}
|
||||
dataList = append(dataList, serverMetric)
|
||||
}
|
||||
}
|
||||
log.Infof("SendMetrics %s length: %d", metricName, len(dataList))
|
||||
err = common.SendMetrics(self.Session, dataList, self.Args.Debug, "")
|
||||
if err != nil {
|
||||
log.Errorf("SendMetrics %s length: %d error: %v", metricName, len(dataList), err)
|
||||
}
|
||||
|
||||
}(_metricName, _influxDbSpecs)
|
||||
}
|
||||
wg.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SApsaraCloudReport) collectRegionMetricOfRedis(region cloudprovider.ICloudRegion,
|
||||
@@ -242,10 +264,8 @@ func (self *SApsaraCloudReport) collectRegionMetricOfElb(region cloudprovider.IC
|
||||
return common.SendMetrics(self.Session, dataList, self.Args.Debug, "")
|
||||
}
|
||||
|
||||
func (self *SApsaraCloudReport) collectMetricFromThisServer(server jsonutils.JSONObject, rtnMetric SApsaraMetric,
|
||||
influxDbSpecs []string) (influxdb.SMetricData, error) {
|
||||
func (self *SApsaraCloudReport) collectMetricFromThisServer(server jsonutils.JSONObject, rtnMetric SApsaraMetric, influxDbSpecs []string) (influxdb.SMetricData, error) {
|
||||
metric, err := self.NewMetricFromJson(server)
|
||||
//metric, err := common.JsonToMetric(server.(*jsonutils.JSONDict), "", common.ServerTags, make([]string, 0))
|
||||
if err != nil {
|
||||
return influxdb.SMetricData{}, err
|
||||
}
|
||||
|
||||
@@ -77,7 +77,7 @@ var OtherHostTag = map[string]string{
|
||||
type ReportOptions struct {
|
||||
Batch int `help:"batch"`
|
||||
Count int `help:"count" json:"count"`
|
||||
Interval string `help:"interval""`
|
||||
Interval string `help:"interval"`
|
||||
Timeout int64 `help:"command timeout unit:second" default:"10"`
|
||||
SinceTime string `help:"sinceTime"`
|
||||
EndTime string `help:"endTime"`
|
||||
|
||||
@@ -144,6 +144,7 @@ func (r *SRegion) DescribeMetricList(department, name string, ns string, since t
|
||||
params := make(map[string]string)
|
||||
params["MetricName"] = name
|
||||
params["Namespace"] = ns
|
||||
params["Period"] = "60"
|
||||
params["Length"] = "2000"
|
||||
if len(nextToken) > 0 {
|
||||
params["NextToken"] = nextToken
|
||||
|
||||
Reference in New Issue
Block a user