fix(cloudmon): qcloud k8s metric

This commit is contained in:
ioito
2022-08-30 17:20:38 +08:00
parent 47ff4eb2b7
commit 53557b09fb
7 changed files with 159 additions and 151 deletions
+17
View File
@@ -52,6 +52,23 @@ type KubeClusterDetails struct {
}
func (self KubeClusterDetails) GetMetricTags() map[string]string {
ret := map[string]string{
"res_type": "kube_cluster",
"cluster_id": self.ExternalClusterId,
"cluster_name": self.Name,
"status": self.Status,
"cloudregion": self.Cloudregion,
"cloudregion_id": self.CloudregionId,
"region_ext_id": self.RegionExtId,
"domain_id": self.DomainId,
"project_domain": self.ProjectDomain,
"account": self.Account,
"account_id": self.AccountId,
}
return ret
}
func (self KubeClusterDetails) GetMetricPairs() map[string]string {
ret := map[string]string{}
return ret
}
+55 -60
View File
@@ -20,7 +20,6 @@ import (
"sync"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
@@ -29,7 +28,6 @@ import (
"yunion.io/x/onecloud/pkg/cloudmon/options"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules/k8s"
"yunion.io/x/onecloud/pkg/util/influxdb"
)
@@ -536,62 +534,58 @@ func (self *SCollectByResourceIdDriver) CollectK8sMetrics(ctx context.Context, m
log.Infof("skip collect %s %s(%s) metric, because not with local kubeserver", vm.Name, manager.Name, manager.Id)
return
}
params := jsonutils.Marshal(map[string]interface{}{
"cluster": vm.ExternalClusterId,
"scope": "system",
"limit": 100,
})
resp, err := k8s.K8sNodes.List(s, params)
opts := &cloudprovider.MetricListOptions{
ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_K8S,
ResourceId: vm.ExternalId,
RegionExtId: vm.RegionExtId,
StartTime: start,
EndTime: end,
}
data, err := provider.GetMetrics(opts)
if err != nil {
log.Errorf("Pods.List error: %v", err)
if errors.Cause(err) != cloudprovider.ErrNotImplemented && errors.Cause(err) != cloudprovider.ErrNotSupported {
log.Errorf("get %s %s(%s) error: %v", opts.ResourceType, vm.Name, vm.Id, err)
return
}
return
}
nodes := []struct {
Id string
Name string
}{}
jsonutils.Update(&nodes, resp.Data)
for i := range nodes {
opts := &cloudprovider.MetricListOptions{
ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_K8S,
ResourceId: vm.ExternalId,
RegionExtId: vm.RegionExtId,
StartTime: start,
EndTime: end,
Node: nodes[i].Name,
}
data, err := provider.GetMetrics(opts)
if err != nil {
if errors.Cause(err) != cloudprovider.ErrNotImplemented && errors.Cause(err) != cloudprovider.ErrNotSupported {
log.Errorf("get %s %s(%s) error: %v", opts.ResourceType, vm.Name, vm.Id, err)
continue
}
continue
}
for _, values := range data {
for _, value := range values.Values {
metric := influxdb.SMetricData{
Name: values.MetricType.Name(),
Timestamp: value.Timestamp,
Tags: []influxdb.SKeyValue{},
Metrics: []influxdb.SKeyValue{
{
Key: values.MetricType.Key(),
Value: strconv.FormatFloat(value.Value, 'E', -1, 64),
},
tags := []influxdb.SKeyValue{}
for k, v := range vm.GetMetricTags() {
tags = append(tags, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
pairs := []influxdb.SKeyValue{}
for k, v := range vm.GetMetricPairs() {
pairs = append(pairs, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
for _, values := range data {
for _, value := range values.Values {
metric := influxdb.SMetricData{
Name: values.MetricType.Name(),
Timestamp: value.Timestamp,
Tags: []influxdb.SKeyValue{},
Metrics: []influxdb.SKeyValue{
{
Key: values.MetricType.Key(),
Value: strconv.FormatFloat(value.Value, 'E', -1, 64),
},
}
for k, v := range value.Tags {
metric.Tags = append(metric.Tags, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
mu.Lock()
metrics = append(metrics, metric)
mu.Unlock()
},
}
for k, v := range value.Tags {
metric.Tags = append(metric.Tags, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
metric.Metrics = append(metric.Metrics, pairs...)
mu.Lock()
metrics = append(metrics, metric)
mu.Unlock()
}
}
}(res[i])
@@ -1018,7 +1012,7 @@ func (self *SCollectByMetricTypeDriver) CollectK8sMetrics(ctx context.Context, m
metrics := []influxdb.SMetricData{}
var wg sync.WaitGroup
var mu sync.Mutex
for _, _metricType := range cloudprovider.ALL_RDS_METRIC_TYPES {
for _, _metricType := range cloudprovider.ALL_K8S_NODE_TYPES {
wg.Add(1)
go func(metricType cloudprovider.TMetricType) {
defer func() {
@@ -1054,6 +1048,13 @@ func (self *SCollectByMetricTypeDriver) CollectK8sMetrics(ctx context.Context, m
metric := influxdb.SMetricData{
Name: value.MetricType.Name(),
Timestamp: v.Timestamp,
Tags: []influxdb.SKeyValue{},
Metrics: []influxdb.SKeyValue{
{
Key: value.MetricType.Key(),
Value: strconv.FormatFloat(v.Value, 'E', -1, 64),
},
},
}
for k, v := range v.Tags {
metric.Tags = append([]influxdb.SKeyValue{
@@ -1063,12 +1064,6 @@ func (self *SCollectByMetricTypeDriver) CollectK8sMetrics(ctx context.Context, m
},
}, tags...)
}
metric.Metrics = []influxdb.SKeyValue{
{
Key: value.MetricType.Key(),
Value: strconv.FormatFloat(v.Value, 'E', -1, 64),
},
}
mu.Lock()
metrics = append(metrics, metric)
mu.Unlock()
+5
View File
@@ -352,3 +352,8 @@ func (self *QcloudCollect) CollectRedisMetrics(ctx context.Context, manager api.
log.Infof("send %d redis with %d metrics for %s(%s)", len(res), len(metrics), manager.Name, manager.Id)
return influxdb.BatchSendMetrics(urls, options.Options.InfluxDatabase, metrics, false)
}
func (self *QcloudCollect) CollectK8sMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.KubeClusterDetails, start, end time.Time) error {
base := &SCollectByResourceIdDriver{}
return base.CollectK8sMetrics(ctx, manager, provider, res, start, end)
}
+13 -22
View File
@@ -72,6 +72,8 @@ const (
RDS_METRIC_TYPE_CONN_USAGE TMetricType = "rds_conn.used_percent"
RDS_METRIC_TYPE_CONN_FAILED TMetricType = "rds_conn.failed_count"
METRIC_TAG_DATABASE = "database"
RDS_METRIC_TYPE_QPS TMetricType = "rds_qps.query_qps"
RDS_METRIC_TYPE_TPS TMetricType = "rds_tps.trans_qps"
RDS_METRIC_TYPE_INNODB_READ_BPS TMetricType = "rds_innodb.read_bps"
@@ -149,25 +151,13 @@ const (
// 磁盘利用率
METRIC_TAG_DEVICE = "device"
K8S_CLUSTER_METRIC_TYPE_CPU_USAGE TMetricType = "k8s_cluster.cpu_used_percent"
K8S_CLUSTER_METRIC_TYPE_MEM_USAGE TMetricType = "k8s_cluster.mem_used_percent"
K8S_CLUSTER_METRIC_TYPE_ALLOCATABLE_POD TMetricType = "k8s_cluster.allocatable_pod"
K8S_CLUSTER_METRIC_TYPE_TOTAL_CPUCORE TMetricType = "k8s_cluster.total_cpu"
K8S_CLUSTER_METRIC_TYPE_CPU_ALLOCATED TMetricType = "k8s_cluster.cpu_allocated_percent"
METRIC_TAG_NODE = "node"
K8S_NODE_METRIC_TYPE_CPU_USAGE TMetricType = "k8s_node.cpu_used_percent"
K8S_NODE_METRIC_TYPE_MEM_USAGE TMetricType = "k8s_node.mem_used_percent"
K8S_NODE_METRIC_TYPE_DISK_USAGE TMetricType = "k8s_node.disk_used_percent"
K8S_NODE_METRIC_TYPE_NET_BPS_RX TMetricType = "k8s_node.bps_recv"
K8S_NODE_METRIC_TYPE_NET_BPS_TX TMetricType = "k8s_node.bps_sent"
K8S_NODE_METRIC_TYPE_POD_RESTART_TOTAL TMetricType = "k8s_node.pod_restart_total"
K8S_POD_METRIC_TYPE_CPU_USAGE TMetricType = "k8s_pod.cpu_used_percent"
K8S_POD_METRIC_TYPE_MEM_USAGE TMetricType = "k8s_pod.mem_used_percent"
K8S_POD_METRIC_TYPE_RESTART_TOTAL TMetricType = "k8s_pod.restart_total"
K8S_POD_METRIC_TYPE_OOM_CONTAINER_COUNT TMetricType = "k8s_deploy.pod_oom_total"
K8S_POD_METRIC_TYPE_RESTARTING_COUNT TMetricType = "k8s_deploy.pod_restarting_total"
K8S_NODE_METRIC_TYPE_CPU_USAGE TMetricType = "k8s_node_cpu.usage_active"
K8S_NODE_METRIC_TYPE_MEM_USAGE TMetricType = "k8s_node_mem.used_percent"
K8S_NODE_METRIC_TYPE_DISK_USAGE TMetricType = "k8s_node_disk.used_percent"
K8S_NODE_METRIC_TYPE_NET_BPS_RX TMetricType = "k8s_node_netio.bps_recv"
K8S_NODE_METRIC_TYPE_NET_BPS_TX TMetricType = "k8s_node_netio.bps_sent"
)
var (
@@ -256,6 +246,11 @@ var (
BUCKET_METRIC_TYPE_LATECY,
BUCKET_METRYC_TYPE_REQ_COUNT,
}
ALL_K8S_NODE_TYPES = []TMetricType{
K8S_NODE_METRIC_TYPE_CPU_USAGE,
K8S_NODE_METRIC_TYPE_MEM_USAGE,
}
)
type MetricListOptions struct {
@@ -272,10 +267,6 @@ type MetricListOptions struct {
Interval int
// rds
Engine string
// k8s
Node string
Pod string
}
type MetricValue struct {
-16
View File
@@ -627,22 +627,6 @@ func (self *SAliyunClient) GetElbMetrics(opts *cloudprovider.MetricListOptions)
func (self *SAliyunClient) GetK8sMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
metricTags, tagKey := map[string]string{}, ""
switch opts.MetricType {
case cloudprovider.K8S_CLUSTER_METRIC_TYPE_CPU_USAGE:
metricTags = map[string]string{
"cluster.cpu.utilization": "",
}
case cloudprovider.K8S_CLUSTER_METRIC_TYPE_MEM_USAGE:
metricTags = map[string]string{
"cluster.memory.utilization": "",
}
case cloudprovider.K8S_POD_METRIC_TYPE_CPU_USAGE:
metricTags = map[string]string{
"pod.cpu.utilization": "",
}
case cloudprovider.K8S_POD_METRIC_TYPE_MEM_USAGE:
metricTags = map[string]string{
"pod.memory.utilization": "",
}
case cloudprovider.K8S_NODE_METRIC_TYPE_CPU_USAGE:
metricTags = map[string]string{
"node.cpu.utilization": "",
+17 -23
View File
@@ -170,7 +170,7 @@ func (self *SAzureClient) GetEcsMetrics(opts *cloudprovider.MetricListOptions) (
if strings.Contains(opts.ResourceId, "microsoft.classiccompute/virtualmachines") {
metricnamespace = "microsoft.classiccompute/virtualmachines"
}
ret, err := self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, "", "", opts.StartTime, opts.EndTime)
ret, err := self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, nil, "", opts.StartTime, opts.EndTime)
if err != nil {
return nil, err
}
@@ -237,27 +237,20 @@ func (self *SAzureClient) GetEcsMetrics(opts *cloudprovider.MetricListOptions) (
func (self *SAzureClient) GetRedisMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
metricnamespace := "Microsoft.Cache/redis"
metricnames := "percentProcessorTime,usedmemorypercentage,connectedclients,operationsPerSecond,alltotalkeys,expiredkeys,usedmemory,serverLoad,errors"
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, "", "", opts.StartTime, opts.EndTime)
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, nil, "", opts.StartTime, opts.EndTime)
}
func (self *SAzureClient) GetLbMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
metricnamespace := "Microsoft.Network/loadBalancers"
metricnames := "SnatConnectionCount,UsedSnatPorts"
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, "", "", opts.StartTime, opts.EndTime)
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, nil, "", opts.StartTime, opts.EndTime)
}
func (self *SAzureClient) GetK8sMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
filter := ""
metricnamespace := "Microsoft.ContainerService/managedClusters"
metricnames := "node_cpu_usage_percentage,node_memory_rss_percentage,node_disk_usage_percentage,node_network_in_bytes,node_network_out_bytes"
if len(opts.Node) > 0 {
filter = fmt.Sprintf("node eq '%s'", opts.Node)
} else if len(opts.Pod) > 0 {
metricnamespace = "insights.container/pods"
metricnames = "oomKilledContainerCount,restartingContainerCount"
}
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, "", filter, opts.StartTime, opts.EndTime)
filter := fmt.Sprintf("node eq '*'")
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, nil, filter, opts.StartTime, opts.EndTime)
}
func (self *SAzureClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
@@ -288,7 +281,7 @@ func (self *SAzureClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions) (
if result.Value[i].Name == "master" {
continue
}
metrics, err := self.getMetricValues(result.Value[i].ID, metricnamespace, metricnames, result.Value[i].Name, "", opts.StartTime, opts.EndTime)
metrics, err := self.getMetricValues(result.Value[i].ID, metricnamespace, metricnames, map[string]string{cloudprovider.METRIC_TAG_DATABASE: result.Value[i].Name}, "", opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("error: %v", err)
continue
@@ -305,7 +298,7 @@ func (self *SAzureClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions) (
default:
return nil, errors.Wrapf(cloudprovider.ErrNotSupported, opts.Engine)
}
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, "", "", opts.StartTime, opts.EndTime)
return self.getMetricValues(opts.ResourceId, metricnamespace, metricnames, nil, "", opts.StartTime, opts.EndTime)
}
type MetrifDefinitions struct {
@@ -356,7 +349,7 @@ func (self *SAzureClient) getMetricDefinitions(resourceId, filter string) (*Metr
return result, nil
}
func (self *SAzureClient) getMetricValues(resourceId, metricnamespace, metricnames string, database, filter string, startTime, endTime time.Time) ([]cloudprovider.MetricValues, error) {
func (self *SAzureClient) getMetricValues(resourceId, metricnamespace, metricnames string, metricTag map[string]string, filter string, startTime, endTime time.Time) ([]cloudprovider.MetricValues, error) {
ret := []cloudprovider.MetricValues{}
params := url.Values{}
params.Set("interval", "PT1M")
@@ -442,20 +435,21 @@ func (self *SAzureClient) getMetricValues(resourceId, metricnamespace, metricnam
metric.MetricType = cloudprovider.K8S_NODE_METRIC_TYPE_NET_BPS_RX
case "node_network_out_bytes":
metric.MetricType = cloudprovider.K8S_NODE_METRIC_TYPE_NET_BPS_TX
case "oomKilledContainerCount":
metric.MetricType = cloudprovider.K8S_POD_METRIC_TYPE_OOM_CONTAINER_COUNT
case "restartingContainerCount":
metric.MetricType = cloudprovider.K8S_POD_METRIC_TYPE_RESTARTING_COUNT
default:
log.Warningf("incognizance metric type %s", element.Name.Value)
continue
}
for _, timeserie := range element.Timeseries {
for _, data := range timeserie.Data {
tags := map[string]string{}
if len(database) > 0 {
tags["database"] = database
tags := map[string]string{}
for _, metadata := range timeserie.Metadatavalues {
if metadata.Name.Value == "node" { //k8s node
tags[cloudprovider.METRIC_TAG_NODE] = metadata.Value
}
}
for k, v := range metricTag {
tags[k] = v
}
for _, data := range timeserie.Data {
metric.Values = append(metric.Values, cloudprovider.MetricValue{
Timestamp: data.TimeStamp,
Value: data.GetValue(),
+52 -30
View File
@@ -107,43 +107,25 @@ func (self *SQcloudClient) GetMonitorData(ns string, name string, since time.Tim
return ret, nil
}
/*
func (r *SRegion) GetK8sMonitorData(metricNames []string, ns string, since time.Time, until time.Time,
demensions []SQcMetricDimension) ([]SK8SDataPoint, error) {
func (self *SQcloudClient) GetK8sMonitorData(ns, name, resourceId string, since time.Time, until time.Time, regionId string) ([]SK8SDataPoint, error) {
params := make(map[string]string)
params["Module"] = "monitor"
params["Region"] = r.Region
for index, name := range metricNames {
i := strconv.FormatInt(int64(index), 10)
params["MetricNames."+i] = name
}
params["Region"] = regionId
params["MetricNames.0"] = name
params["Namespace"] = ns
if !since.IsZero() {
params["StartTime"] = since.Format(timeutils.IsoTimeFormat)
}
if !until.IsZero() {
params["EndTime"] = until.Format(timeutils.IsoTimeFormat)
}
for index, metricDimension := range demensions {
i := strconv.FormatInt(int64(index), 10)
params["Conditions."+i+".Key"] = metricDimension.Name
params["Conditions."+i+".Operator"] = "="
params["Conditions."+i+".Value.0"] = metricDimension.Value
}
body, err := r.metricsRequest("DescribeStatisticData", params)
params["Period"] = "60"
params["StartTime"] = since.Format(timeutils.IsoTimeFormat)
params["EndTime"] = until.Format(timeutils.IsoTimeFormat)
params["Conditions.0.Key"] = "tke_cluster_instance_id"
params["Conditions.0.Operator"] = "in"
params["Conditions.0.Value.0"] = resourceId
body, err := self.metricsRequest("DescribeStatisticData", params)
if err != nil {
return nil, errors.Wrap(err, "region.MetricRequest")
}
dataArray := make([]SK8SDataPoint, 0)
err = body.Unmarshal(&dataArray, "Data")
if err != nil {
return nil, errors.Wrap(err, "resp.Unmarshal")
}
return dataArray, nil
ret := []SK8SDataPoint{}
return ret, body.Unmarshal(&ret, "Data")
}
*/
func (self *SQcloudClient) GetMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
switch opts.ResourceType {
@@ -153,6 +135,8 @@ func (self *SQcloudClient) GetMetrics(opts *cloudprovider.MetricListOptions) ([]
return self.GetRedisMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_RDS:
return self.GetRdsMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_K8S:
return self.GetK8sMetrics(opts)
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.ResourceType)
}
@@ -348,3 +332,41 @@ func (self *SQcloudClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions)
}
return ret, nil
}
func (self *SQcloudClient) GetK8sMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
ret := []cloudprovider.MetricValues{}
for metricType, metricName := range map[cloudprovider.TMetricType]string{
cloudprovider.K8S_NODE_METRIC_TYPE_CPU_USAGE: "K8sNodeCpuUsage",
cloudprovider.K8S_NODE_METRIC_TYPE_MEM_USAGE: "K8sNodeMemUsage",
} {
metrics, err := self.GetK8sMonitorData("QCE/TKE2", metricName, opts.ResourceId, opts.StartTime, opts.EndTime, opts.RegionExtId)
if err != nil {
log.Errorf("GetMonitorData error: %v", err)
continue
}
for i := range metrics {
for _, point := range metrics[i].Points {
tags := map[string]string{}
for _, dim := range point.Dimensions {
if dim.Name == "node" {
tags[cloudprovider.METRIC_TAG_NODE] = dim.Value
}
}
metric := cloudprovider.MetricValues{}
metric.Id = opts.ResourceId
metric.MetricType = metricType
for _, value := range point.Values {
metric.Values = append(metric.Values, cloudprovider.MetricValue{
Value: value.Value,
Timestamp: time.Unix(int64(value.Timestamp), 0),
Tags: tags,
})
}
if len(metric.Values) > 0 {
ret = append(ret, metric)
}
}
}
}
return ret, nil
}