diff --git a/pkg/apis/compute/kube_cluster.go b/pkg/apis/compute/kube_cluster.go index 0ce497291a..8755b65522 100644 --- a/pkg/apis/compute/kube_cluster.go +++ b/pkg/apis/compute/kube_cluster.go @@ -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 } diff --git a/pkg/cloudmon/providerdriver/base.go b/pkg/cloudmon/providerdriver/base.go index f9ec678edc..dbed0c7650 100644 --- a/pkg/cloudmon/providerdriver/base.go +++ b/pkg/cloudmon/providerdriver/base.go @@ -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() diff --git a/pkg/cloudmon/providerdriver/qcloud.go b/pkg/cloudmon/providerdriver/qcloud.go index d9d350e97d..d63e179ada 100644 --- a/pkg/cloudmon/providerdriver/qcloud.go +++ b/pkg/cloudmon/providerdriver/qcloud.go @@ -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) +} diff --git a/pkg/cloudprovider/metrics.go b/pkg/cloudprovider/metrics.go index 11e37700d5..252ec5cb6c 100644 --- a/pkg/cloudprovider/metrics.go +++ b/pkg/cloudprovider/metrics.go @@ -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 { diff --git a/pkg/multicloud/aliyun/monitor.go b/pkg/multicloud/aliyun/monitor.go index 18c1127010..479f4c4aa7 100644 --- a/pkg/multicloud/aliyun/monitor.go +++ b/pkg/multicloud/aliyun/monitor.go @@ -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": "", diff --git a/pkg/multicloud/azure/monitor.go b/pkg/multicloud/azure/monitor.go index 64b7e72ddd..9c0a2bdb82 100644 --- a/pkg/multicloud/azure/monitor.go +++ b/pkg/multicloud/azure/monitor.go @@ -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(), diff --git a/pkg/multicloud/qcloud/monitor.go b/pkg/multicloud/qcloud/monitor.go index 43504bf4b8..c159d0ca05 100644 --- a/pkg/multicloud/qcloud/monitor.go +++ b/pkg/multicloud/qcloud/monitor.go @@ -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 +}