Merge pull request #16277 from ioito/hotfix/qx-apsara-eip-metric

fix(cloudmon): add apsara eip metric
This commit is contained in:
Zexi Li
2023-03-23 15:04:27 +08:00
committed by GitHub
11 changed files with 497 additions and 36 deletions
+1 -1
View File
@@ -83,7 +83,7 @@ require (
k8s.io/client-go v0.19.3
k8s.io/cluster-bootstrap v0.19.3
moul.io/http2curl/v2 v2.3.0
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230322095023-a4470c57d2d2
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230323054147-6e2b259a89f2
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656
yunion.io/x/jsonutils v1.0.1-0.20220819091305-3bab322ab4fd
yunion.io/x/log v1.0.0
+2 -2
View File
@@ -1164,8 +1164,8 @@ sigs.k8s.io/structured-merge-diff/v4 v4.0.1/go.mod h1:bJZC9H9iH24zzfZ/41RGcq60oK
sigs.k8s.io/yaml v1.1.0/go.mod h1:UJmg0vDUVViEyp3mgSv9WPwZCDxu4rQW1olrI1uml+o=
sigs.k8s.io/yaml v1.2.0 h1:kr/MCeFWJWTwyaHoR9c8EjH9OumOmoF9YGiZd7lFm/Q=
sigs.k8s.io/yaml v1.2.0/go.mod h1:yfXDCHCao9+ENCvLSE62v9VSji2MKu5jeNfTrofGhJc=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230322095023-a4470c57d2d2 h1:hvIcldFUzhQVvyNd3mbZxNs/g1w2q6aWpEycXMeUSBI=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230322095023-a4470c57d2d2/go.mod h1:sWqblYRhQCO63xeKdxdAT3wczCCY8Mc+TDa3WD/a5Kg=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230323054147-6e2b259a89f2 h1:vqU3NZvfTPFKnj8ycGqH4CbC5q7y6GpCDMpVQjFjU8w=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230323054147-6e2b259a89f2/go.mod h1:sWqblYRhQCO63xeKdxdAT3wczCCY8Mc+TDa3WD/a5Kg=
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656 h1:0zlZD5uhZoIHgLVAWCz2aHaYk2ZrNsACCYD7R6EIBII=
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656/go.mod h1:Uxuou9WQIeJXNpy7t2fPLL0BYLvLiMvGQwY7Qc6aSws=
yunion.io/x/jsonutils v0.0.0-20190625054549-a964e1e8a051/go.mod h1:4N0/RVzsYL3kH3WE/H1BjUQdFiWu50JGCFQuuy+Z634=
+23 -1
View File
@@ -14,7 +14,9 @@
package compute
import "yunion.io/x/onecloud/pkg/apis"
import (
"yunion.io/x/onecloud/pkg/apis"
)
type SElasticipCreateInput struct {
apis.VirtualResourceCreateInput
@@ -87,6 +89,26 @@ type ElasticipDetails struct {
AssociateName string `json:"associate_name"`
}
func (self ElasticipDetails) GetMetricTags() map[string]string {
ret := map[string]string{
"id": self.Id,
"name": self.Name,
"status": self.Status,
"mode": self.Mode,
"cloudregion": self.Cloudregion,
"cloudregion_id": self.CloudregionId,
"region_ext_id": self.RegionExtId,
"tenant": self.Project,
"tenant_id": self.ProjectId,
"brand": self.Brand,
"domain_id": self.DomainId,
"project_domain": self.ProjectDomain,
"ip_addr": self.IpAddr,
"external_id": self.ExternalId,
}
return ret
}
type ElasticipSyncstatusInput struct {
}
+77
View File
@@ -15,7 +15,17 @@
package providerdriver
import (
"context"
"strconv"
"sync"
"time"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/util/influxdb"
)
type ApsaraCollect struct {
@@ -33,3 +43,70 @@ func (self *ApsaraCollect) IsSupportMetrics() bool {
func init() {
Register(&ApsaraCollect{})
}
func (self *ApsaraCollect) CollectEipMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ElasticipDetails, start, end time.Time) error {
metrics := []influxdb.SMetricData{}
var wg sync.WaitGroup
var mu sync.Mutex
for _, _metricType := range cloudprovider.ALL_EIP_TYPES {
wg.Add(1)
go func(metricType cloudprovider.TMetricType) {
defer func() {
wg.Done()
}()
opts := &cloudprovider.MetricListOptions{
ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_EIP,
MetricType: metricType,
StartTime: start,
EndTime: end,
}
data, err := provider.GetMetrics(opts)
if err != nil {
if errors.Cause(err) != cloudprovider.ErrNotImplemented && errors.Cause(err) != cloudprovider.ErrNotSupported {
log.Errorf("get eip %s(%s) %s error: %v", manager.Name, manager.Id, metricType, err)
return
}
return
}
for _, value := range data {
eip, ok := res[value.Id]
if !ok {
continue
}
tags := []influxdb.SKeyValue{}
for k, v := range eip.GetMetricTags() {
tags = append(tags, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
for _, v := range value.Values {
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(metric.Tags, influxdb.SKeyValue{
Key: k,
Value: v,
})
}
metric.Tags = append(metric.Tags, tags...)
mu.Lock()
metrics = append(metrics, metric)
mu.Unlock()
}
}
}(_metricType)
}
wg.Wait()
return self.sendMetrics(ctx, manager, "elasticip", len(res), metrics)
}
+4
View File
@@ -62,6 +62,10 @@ func (self *SBaseCollectDriver) CollectWireMetrics(ctx context.Context, manager
return cloudprovider.ErrNotImplemented
}
func (self *SBaseCollectDriver) CollectEipMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ElasticipDetails, start, end time.Time) error {
return cloudprovider.ErrNotImplemented
}
func (self *SBaseCollectDriver) CollectStorageMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.StorageDetails, start, end time.Time) error {
metrics := []influxdb.SMetricData{}
for _, storage := range res {
+1
View File
@@ -42,6 +42,7 @@ type ICollectDriver interface {
CollectK8sMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.KubeClusterDetails, start, end time.Time) error
CollectModelartsPoolMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ModelartsPoolDetails, start, end time.Time) error
CollectWireMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.WireDetails, start, end time.Time) error
CollectEipMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ElasticipDetails, start, end time.Time) error
}
func GetDriver(name string) (ICollectDriver, error) {
+26
View File
@@ -386,6 +386,7 @@ type SResources struct {
ModelartsPool TResource
Wires TResource
Projects TResource
ElasticIps TResource
}
func (self *SResources) IsInit() bool {
@@ -407,6 +408,7 @@ func NewResources() *SResources {
ModelartsPool: NewBaseResources(&compute.ModelartsPools),
Wires: NewBaseResources(&compute.Wires),
Projects: NewBaseResources(&identity.Projects),
ElasticIps: NewBaseResources(&compute.Elasticips),
}
}
@@ -462,6 +464,10 @@ func (self *SResources) Init(ctx context.Context, userCred mcclient.TokenCredent
if err != nil {
errs = append(errs, errors.Wrapf(err, "ModelartsPool.init"))
}
err = self.ElasticIps.init(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "ElasticIps.init"))
}
err = self.Wires.init(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "Wires.init"))
@@ -529,6 +535,10 @@ func (self *SResources) IncrementSync(ctx context.Context, userCred mcclient.Tok
if err != nil {
errs = append(errs, errors.Wrapf(err, "ModelartsPool.increment"))
}
err = self.ElasticIps.increment(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "Elasticips.increment"))
}
err = self.Wires.increment(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "Wires.increment"))
@@ -594,6 +604,10 @@ func (self *SResources) DecrementSync(ctx context.Context, userCred mcclient.Tok
if err != nil {
errs = append(errs, errors.Wrapf(err, "ModelartsPool.decrement"))
}
err = self.ElasticIps.decrement(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "ElasticIps.decrement"))
}
err = self.Projects.decrement(ctx)
if err != nil {
errs = append(errs, errors.Wrapf(err, "Projects.decrement"))
@@ -827,6 +841,18 @@ func (self *SResources) CollectMetrics(ctx context.Context, userCred mcclient.To
log.Errorf("CollectWireMetrics for %s(%s) error: %v", manager.Name, manager.Provider, err)
}
resources = self.ElasticIps.getResources(ctx, manager.Id)
eips := map[string]api.ElasticipDetails{}
err = jsonutils.Update(&eips, resources)
if err != nil {
log.Errorf("unmarsha eips resources error: %v", err)
}
err = driver.CollectEipMetrics(ctx, manager, provider, eips, startTime, endTime)
if err != nil && errors.Cause(err) != cloudprovider.ErrNotImplemented && errors.Cause(err) != cloudprovider.ErrNotSupported {
log.Errorf("CollectEipMetrics for %s(%s) error: %v", manager.Name, manager.Provider, err)
}
}(cloudproviders[i])
}
wg.Wait()
+1 -1
View File
@@ -1442,7 +1442,7 @@ sigs.k8s.io/structured-merge-diff/v4/value
# sigs.k8s.io/yaml v1.2.0
## explicit; go 1.12
sigs.k8s.io/yaml
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230322095023-a4470c57d2d2
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230323054147-6e2b259a89f2
## explicit; go 1.18
yunion.io/x/cloudmux/pkg/apis
yunion.io/x/cloudmux/pkg/apis/billing
+99 -5
View File
@@ -55,6 +55,7 @@ const (
METRIC_RESOURCE_TYPE_WIRE TResourceType = "wire"
METRIC_RESOURCE_TYPE_CLOUD_ACCOUNT TResourceType = "cloudaccount_balance"
METRIC_RESOURCE_TYPE_MODELARTS_POOL TResourceType = "modelarts"
METRIC_RESOURCE_TYPE_EIP TResourceType = "eip"
)
const (
@@ -106,8 +107,8 @@ const (
METRIC_TAG_DATABASE = "database"
// RDS QPS
// 支持平台: huawei, qcloud
// RDS QPS(每秒查询数)
// 支持平台: huawei, qcloud, aliyun, apsara
RDS_METRIC_TYPE_QPS TMetricType = "rds_qps.query_qps"
// RDS TPS
// 支持平台: huawei, qcloud
@@ -150,6 +151,13 @@ const (
// 支持平台: huawei, aliyun, apsara, azure, esxi, google, bingocloud, aws, jdcloud, ecloud, zstack, qcloud
VM_METRIC_TYPE_NET_BPS_TX TMetricType = "vm_netio.bps_sent"
// 虚拟机TCP连接数
// 支持平台: aliyun, apsara
VM_METRIC_TYPE_NET_TCP_CONNECTION TMetricType = "vm_netio.tcp_connections"
// 虚拟机进程监控
// 支持平台: aliyun, apsara
VM_METRIC_TYPE_PROCESS_NUMBER = "vm_process.number"
// 宿主机CPU使用率
// 支持平台: esxi
HOST_METRIC_TYPE_CPU_USAGE TMetricType = "cpu.usage_active"
@@ -201,15 +209,42 @@ const (
LB_METRIC_TYPE_SNAT_PORT TMetricType = "haproxy.used_snat_port"
// 支持平台: azure
LB_METRIC_TYPE_SNAT_CONN_COUNT TMetricType = "haproxy.snat_conn_count"
// 入速率
// 入带宽速率
// 支持平台: huawei, aliyun, apsara
LB_METRIC_TYPE_NET_BPS_RX TMetricType = "haproxy.bin"
// 出速率
// 出带宽速率
// 支持平台: huawei, aliyun, apsara
LB_METRIC_TYPE_NET_BPS_TX TMetricType = "haproxy.bout"
// 入包速率
// 支持平台: aliyun, apsara
LB_METRIC_TYPE_NET_PACKET_RX TMetricType = "haproxy.packet_rx"
// 出包速率
// 支持平台: aliyun, apsara
LB_METRIC_TYPE_NET_PACKET_TX TMetricType = "haproxy.packet_tx"
// 非活跃连接数
// 支持平台: apsara, aliyun
LB_METRIC_TYPE_NET_INACTIVE_CONNECTION = "haproxy.inactive_connection"
// 最大并发数
// 支持平台: apsara, aliyun
LB_METRIC_TYPE_MAX_CONNECTION = "haproxy.max_connection"
// 后端异常ECS实例个数
// 支持平台: apsara, aliyun
LB_METRIC_TYPE_UNHEALTHY_SERVER_COUNT = "haproxy.unhealthy_server_count"
// 状态码统计
// 支持平台: huawei, aliyun, apsara
LB_METRIC_TYPE_HRSP_COUNT TMetricType = "haproxy.hrsp_Nxx"
// 入方向丢弃流量
// 支持平台: aliyun
LB_METRIC_TYPE_DROP_TRAFFIC_TX = "haproxy.drop_traffic_tx"
// 出方向丢弃流量
// 支持平台: aliyun, apsara
LB_METRIC_TYPE_DROP_TRAFFIC_RX = "haproxy.drop_traffic_rx"
// 入方向丢弃包数
// 支持平台: aliyun
LB_METRIC_TYPE_DROP_PACKET_TX = "haproxy.drop_packet_tx"
// 出方向丢弃包数
// 支持平台: aliyun, apsara
LB_METRIC_TYPE_DROP_PACKET_RX = "haproxy.drop_packet_rx"
// 对象存储出速率
// 支持平台: huawei, aliyun, apsara
@@ -220,9 +255,24 @@ const (
// 请求延时
// 支持平台: huawei, aliyun, apsara
BUCKET_METRIC_TYPE_LATECY TMetricType = "oss_latency.req_late"
// 请求数量
// 总请求数量
// 支持平台: huawei, aliyun, apsara
BUCKET_METRYC_TYPE_REQ_COUNT TMetricType = "oss_req.req_count"
// 服务端请求错误数量
// 支持平台: aliyun, apsara
BUCKET_METRIC_TYPE_REQ_5XX_COUNT TMetricType = "oss_req.5xx_count"
// 服务端请求错误数量
// 支持平台: aliyun, apsara
BUCKET_METRIC_TYPE_REQ_4XX_COUNT TMetricType = "oss_req.4xx_count"
// 重定向数量
// 支持平台: aliyun, apsara
BUCKET_METRIC_TYPE_REQ_3XX_COUNT TMetricType = "oss_req.3xx_count"
// 正常请求数量
// 支持平台: aliyun, apsara
BUCKET_METRIC_TYPE_REQ_2XX_COUNT TMetricType = "oss_req.2xx_count"
// 存储总容量(bit)
// 支持平台: aliyun, apsara
BUCKET_METRIC_TYPE_STORAGE_SIZE = "oss_storage.size"
METRIC_TAG_REQUST = "request"
METRIC_TAG_REQUST_GET = "get"
@@ -242,6 +292,12 @@ const (
// 磁盘利用率
METRIC_TAG_DEVICE = "device"
// 进程名称
METRIC_TAG_PROCESS_NAME = "process_name"
// 状态
METRIC_TAG_STATE = "state"
METRIC_TAG_NODE = "node"
// k8s节点CPU使用率
@@ -272,6 +328,14 @@ const (
WIRE_METRIC_TYPE_MEM_USAGE TMetricType = "wire_mem.usage_percent"
WIRE_METRIC_TYPE_NET_RT TMetricType = "wire_net.rt" // 响应时间ms
WIRE_METRIC_TYPE_NET_UNREACHABLE_RATE TMetricType = "wire_net.unreachable_rate" // 不可达率
// EIP入带宽
EIP_METRIC_TYPE_NET_BPS_RX TMetricType = "eip_net.bps_recv"
// EIP出带宽
EIP_METRIC_TYPE_NET_BPS_TX TMetricType = "eip_net.bps_sent"
// EIP 出方向限速丢包率
EIP_METRIC_TYPE_NET_DROP_SPEED_TX TMetricType = "eip_net.drop_speed_rx"
)
var (
@@ -318,6 +382,9 @@ var (
VM_METRIC_TYPE_NET_BPS_RX,
VM_METRIC_TYPE_NET_BPS_TX,
VM_METRIC_TYPE_NET_TCP_CONNECTION,
VM_METRIC_TYPE_PROCESS_NUMBER,
}
ALL_REDIS_METRIC_TYPES = []TMetricType{
@@ -338,6 +405,19 @@ var (
LB_METRIC_TYPE_NET_BPS_RX,
LB_METRIC_TYPE_NET_BPS_TX,
LB_METRIC_TYPE_HRSP_COUNT,
LB_METRIC_TYPE_NET_PACKET_RX,
LB_METRIC_TYPE_NET_PACKET_TX,
LB_METRIC_TYPE_UNHEALTHY_SERVER_COUNT,
LB_METRIC_TYPE_NET_INACTIVE_CONNECTION,
LB_METRIC_TYPE_MAX_CONNECTION,
LB_METRIC_TYPE_DROP_PACKET_RX,
LB_METRIC_TYPE_DROP_PACKET_TX,
LB_METRIC_TYPE_DROP_TRAFFIC_RX,
LB_METRIC_TYPE_DROP_TRAFFIC_TX,
}
ALL_BUCKET_TYPES = []TMetricType{
@@ -345,12 +425,26 @@ var (
BUCKET_METRIC_TYPE_NET_BPS_RX,
BUCKET_METRIC_TYPE_LATECY,
BUCKET_METRYC_TYPE_REQ_COUNT,
BUCKET_METRIC_TYPE_REQ_5XX_COUNT,
BUCKET_METRIC_TYPE_REQ_4XX_COUNT,
BUCKET_METRIC_TYPE_REQ_3XX_COUNT,
BUCKET_METRIC_TYPE_REQ_2XX_COUNT,
BUCKET_METRIC_TYPE_STORAGE_SIZE,
}
ALL_K8S_NODE_TYPES = []TMetricType{
K8S_NODE_METRIC_TYPE_CPU_USAGE,
K8S_NODE_METRIC_TYPE_MEM_USAGE,
}
ALL_EIP_TYPES = []TMetricType{
EIP_METRIC_TYPE_NET_BPS_RX,
EIP_METRIC_TYPE_NET_BPS_TX,
EIP_METRIC_TYPE_NET_DROP_SPEED_TX,
}
)
type MetricListOptions struct {
+111 -4
View File
@@ -195,6 +195,9 @@ type MetricData struct {
Minimum float64
Maximum float64
State string
ProcessName string
Diskname string
Device string
@@ -225,6 +228,12 @@ func (d MetricData) GetTags() map[string]string {
if len(d.Device) > 0 {
ret[cloudprovider.METRIC_TAG_DEVICE] = fmt.Sprintf("%s(%s)", d.Device, d.Diskname)
}
if len(d.ProcessName) > 0 {
ret[cloudprovider.METRIC_TAG_PROCESS_NAME] = d.ProcessName
}
if len(d.State) > 0 {
ret[cloudprovider.METRIC_TAG_STATE] = d.State
}
return ret
}
@@ -310,6 +319,8 @@ func (self *SAliyunClient) GetMetrics(opts *cloudprovider.MetricListOptions) ([]
return self.GetElbMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_K8S:
return self.GetK8sMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_EIP:
return self.GetEipMetrics(opts)
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.ResourceType)
}
@@ -358,6 +369,16 @@ func (self *SAliyunClient) GetEcsMetrics(opts *cloudprovider.MetricListOptions)
metricTags = map[string]string{
"diskusage_utilization": "",
}
case cloudprovider.VM_METRIC_TYPE_PROCESS_NUMBER:
metricTags = map[string]string{
"process.count_processname": "",
}
tagKey = cloudprovider.METRIC_TAG_PROCESS_NAME
case cloudprovider.VM_METRIC_TYPE_NET_TCP_CONNECTION:
metricTags = map[string]string{
"network.tcp.connection_state": "",
}
tagKey = cloudprovider.METRIC_TAG_STATE
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
@@ -407,7 +428,7 @@ func (self *SAliyunClient) GetOssMetrics(opts *cloudprovider.MetricListOptions)
"InternetSend": cloudprovider.METRIC_TAG_NET_TYPE_INTERNET,
"IntranetSend": cloudprovider.METRIC_TAG_NET_TYPE_INTRANET,
}
tagKey = cloudprovider.METRIC_TAG_REQUST
tagKey = cloudprovider.METRIC_TAG_NET_TYPE
case cloudprovider.BUCKET_METRIC_TYPE_NET_BPS_RX:
metricTags = map[string]string{
"InternetRecv": cloudprovider.METRIC_TAG_NET_TYPE_INTERNET,
@@ -416,11 +437,31 @@ func (self *SAliyunClient) GetOssMetrics(opts *cloudprovider.MetricListOptions)
tagKey = cloudprovider.METRIC_TAG_NET_TYPE
case cloudprovider.BUCKET_METRYC_TYPE_REQ_COUNT:
metricTags = map[string]string{
"GetObjectCount": cloudprovider.METRIC_TAG_REQUST_GET,
"PostObjectCount": cloudprovider.METRIC_TAG_REQUST_POST,
"ServerErrorCount": cloudprovider.METRIC_TAG_REQUST_5XX,
"TotalRequestCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_5XX_COUNT:
metricTags = map[string]string{
"ServerErrorCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_4XX_COUNT:
metricTags = map[string]string{
"AuthorizationErrorCount": "authorization",
"ClientOtherErrorCount": "other",
"ClientTimeoutErrorCount": "timeout",
}
tagKey = cloudprovider.METRIC_TAG_REQUST
case cloudprovider.BUCKET_METRIC_TYPE_REQ_3XX_COUNT:
metricTags = map[string]string{
"RedirectCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_2XX_COUNT:
metricTags = map[string]string{
"SuccessCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_STORAGE_SIZE:
metricTags = map[string]string{
"MeteringStorageUtilization": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
@@ -522,6 +563,10 @@ func (self *SAliyunClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions)
metricTags = map[string]string{
"ConnectionUsage": "",
}
case cloudprovider.RDS_METRIC_TYPE_QPS:
metricTags = map[string]string{
"MySQL_QPS": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
@@ -567,6 +612,22 @@ func (self *SAliyunClient) GetElbMetrics(opts *cloudprovider.MetricListOptions)
"InstanceStatusCode5xx": cloudprovider.METRIC_TAG_REQUST_5XX,
}
tagKey = cloudprovider.METRIC_TAG_REQUST
case cloudprovider.LB_METRIC_TYPE_DROP_PACKET_RX:
metricTags = map[string]string{
"InstanceDropPacketRX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_PACKET_TX:
metricTags = map[string]string{
"InstanceDropPacketTX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_TRAFFIC_RX:
metricTags = map[string]string{
"InstanceDropTrafficRX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_TRAFFIC_TX:
metricTags = map[string]string{
"InstanceDropTrafficTX": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
@@ -632,3 +693,49 @@ func (self *SAliyunClient) GetK8sMetrics(opts *cloudprovider.MetricListOptions)
}
return ret, nil
}
func (self *SAliyunClient) GetEipMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
metricTags, tagKey := map[string]string{}, ""
switch opts.MetricType {
case cloudprovider.EIP_METRIC_TYPE_NET_BPS_RX:
metricTags = map[string]string{
"net.rx": "",
}
case cloudprovider.EIP_METRIC_TYPE_NET_BPS_TX:
metricTags = map[string]string{
"net.tx": "",
}
case cloudprovider.EIP_METRIC_TYPE_NET_DROP_SPEED_TX:
metricTags = map[string]string{
"out_ratelimit_drop_speed": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_vpc_eip", metric, opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
}
tags := map[string]string{}
if len(tag) > 0 && len(tagKey) > 0 {
tags[tagKey] = tag
}
for i := range result {
ret = append(ret, cloudprovider.MetricValues{
Id: result[i].InstanceId,
MetricType: opts.MetricType,
Values: []cloudprovider.MetricValue{
{
Timestamp: time.UnixMilli(result[i].Timestamp),
Value: result[i].GetValue(),
Tags: tags,
},
},
})
}
}
return ret, nil
}
+152 -22
View File
@@ -51,6 +51,9 @@ type MetricData struct {
Maximum float64
Diskname string
Device string
State string
ProcessName string
}
func (d MetricData) GetValue() float64 {
@@ -74,6 +77,12 @@ func (d MetricData) GetTags() map[string]string {
if len(d.Device) > 0 {
ret[cloudprovider.METRIC_TAG_DEVICE] = fmt.Sprintf("%s(%s)", d.Device, d.Diskname)
}
if len(d.ProcessName) > 0 {
ret[cloudprovider.METRIC_TAG_PROCESS_NAME] = d.ProcessName
}
if len(d.State) > 0 {
ret[cloudprovider.METRIC_TAG_STATE] = d.State
}
return ret
}
@@ -98,11 +107,11 @@ func (self *SApsaraClient) tryGetDepartments() []string {
return self.departments
}
func (self *SApsaraClient) ListMetrics(ns, metricName string, start, end time.Time) ([]MetricData, error) {
func (self *SApsaraClient) ListMetrics(ns, metricName, dimensions string, start, end time.Time) ([]MetricData, error) {
ret := []MetricData{}
departments := self.tryGetDepartments()
for i := range departments {
part, err := self.listMetrics(departments[i], ns, metricName, start, end)
part, err := self.listMetrics(departments[i], ns, metricName, dimensions, start, end)
if err != nil {
if strings.Contains(err.Error(), "NoPermission") {
continue
@@ -114,11 +123,11 @@ func (self *SApsaraClient) ListMetrics(ns, metricName string, start, end time.Ti
return ret, nil
}
func (self *SApsaraClient) listMetrics(departmentId, ns, metricName string, start, end time.Time) ([]MetricData, error) {
func (self *SApsaraClient) listMetrics(departmentId, ns, metricName, dimensions string, start, end time.Time) ([]MetricData, error) {
result := []MetricData{}
nextToken := ""
for {
part, next, err := self._listMetrics(departmentId, ns, metricName, nextToken, start, end)
part, next, err := self._listMetrics(departmentId, ns, metricName, dimensions, nextToken, start, end)
if err != nil {
return nil, errors.Wrap(err, "listMetrics")
}
@@ -131,7 +140,7 @@ func (self *SApsaraClient) listMetrics(departmentId, ns, metricName string, star
return result, nil
}
func (self *SApsaraClient) _listMetrics(departmentId, ns, metricName, nextToken string, start, end time.Time) ([]MetricData, string, error) {
func (self *SApsaraClient) _listMetrics(departmentId, ns, metricName, dimensions string, nextToken string, start, end time.Time) ([]MetricData, string, error) {
params := make(map[string]string)
params["MetricName"] = metricName
params["Namespace"] = ns
@@ -139,6 +148,9 @@ func (self *SApsaraClient) _listMetrics(departmentId, ns, metricName, nextToken
if len(nextToken) > 0 {
params["NextToken"] = nextToken
}
if len(dimensions) > 0 {
params["Dimensions"] = dimensions
}
params["Department"] = departmentId
params["StartTime"] = fmt.Sprintf("%d", start.UnixMilli())
params["EndTime"] = fmt.Sprintf("%d", end.UnixMilli())
@@ -178,6 +190,8 @@ func (self *SApsaraClient) GetMetrics(opts *cloudprovider.MetricListOptions) ([]
return self.GetRdsMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_LB:
return self.GetElbMetrics(opts)
case cloudprovider.METRIC_RESOURCE_TYPE_EIP:
return self.GetEipMetrics(opts)
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.ResourceType)
}
@@ -226,12 +240,26 @@ func (self *SApsaraClient) GetEcsMetrics(opts *cloudprovider.MetricListOptions)
metricTags = map[string]string{
"diskusage_utilization": "",
}
case cloudprovider.VM_METRIC_TYPE_PROCESS_NUMBER:
metricTags = map[string]string{
"process.number": "",
}
tagKey = cloudprovider.METRIC_TAG_PROCESS_NAME
case cloudprovider.VM_METRIC_TYPE_NET_TCP_CONNECTION:
metricTags = map[string]string{
"net_tcpconnection": "",
}
tagKey = cloudprovider.METRIC_TAG_STATE
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_ecs_dashboard", metric, opts.StartTime, opts.EndTime)
dimensions := ""
if opts.MetricType == cloudprovider.VM_METRIC_TYPE_PROCESS_NUMBER {
dimensions = jsonutils.Marshal(map[string]string{"instanceId": opts.ResourceId}).String()
}
result, err := self.ListMetrics("acs_ecs_dashboard", metric, dimensions, opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
@@ -275,7 +303,7 @@ func (self *SApsaraClient) GetOssMetrics(opts *cloudprovider.MetricListOptions)
"InternetSend": cloudprovider.METRIC_TAG_NET_TYPE_INTERNET,
"IntranetSend": cloudprovider.METRIC_TAG_NET_TYPE_INTRANET,
}
tagKey = cloudprovider.METRIC_TAG_REQUST
tagKey = cloudprovider.METRIC_TAG_NET_TYPE
case cloudprovider.BUCKET_METRIC_TYPE_NET_BPS_RX:
metricTags = map[string]string{
"InternetRecv": cloudprovider.METRIC_TAG_NET_TYPE_INTERNET,
@@ -284,17 +312,41 @@ func (self *SApsaraClient) GetOssMetrics(opts *cloudprovider.MetricListOptions)
tagKey = cloudprovider.METRIC_TAG_NET_TYPE
case cloudprovider.BUCKET_METRYC_TYPE_REQ_COUNT:
metricTags = map[string]string{
"GetObjectCount": cloudprovider.METRIC_TAG_REQUST_GET,
"PostObjectCount": cloudprovider.METRIC_TAG_REQUST_POST,
"ServerErrorCount": cloudprovider.METRIC_TAG_REQUST_5XX,
"TotalRequestCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_5XX_COUNT:
metricTags = map[string]string{
"ServerErrorCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_4XX_COUNT:
metricTags = map[string]string{
"AuthorizationErrorCount": "authorization",
"ClientOtherErrorCount": "other",
"ClientTimeoutErrorCount": "timeout",
}
tagKey = cloudprovider.METRIC_TAG_REQUST
case cloudprovider.BUCKET_METRIC_TYPE_REQ_3XX_COUNT:
metricTags = map[string]string{
"RedirectCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_REQ_2XX_COUNT:
metricTags = map[string]string{
"SuccessCount": "",
}
case cloudprovider.BUCKET_METRIC_TYPE_STORAGE_SIZE:
metricTags = map[string]string{
"MeteringStorageUtilization": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_oss_dashboard", metric, opts.StartTime, opts.EndTime)
dimensions := ""
if opts.MetricType == cloudprovider.BUCKET_METRIC_TYPE_STORAGE_SIZE {
dimensions = jsonutils.Marshal(map[string]string{"BucketName": opts.ResourceId}).String()
}
result, err := self.ListMetrics("acs_oss_dashboard", metric, dimensions, opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
@@ -362,7 +414,7 @@ func (self *SApsaraClient) GetRedisMetrics(opts *cloudprovider.MetricListOptions
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_kvstore", metric, opts.StartTime, opts.EndTime)
result, err := self.ListMetrics("acs_kvstore", metric, "", opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
@@ -417,12 +469,16 @@ func (self *SApsaraClient) GetRdsMetrics(opts *cloudprovider.MetricListOptions)
metrics = map[string]string{
"ConnectionUsage": "",
}
case cloudprovider.RDS_METRIC_TYPE_QPS:
metrics = map[string]string{
"MySQL_QPS": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric := range metrics {
result, err := self.ListMetrics("acs_rds_dashboard", metric, opts.StartTime, opts.EndTime)
result, err := self.ListMetrics("acs_rds_dashboard", metric, "", opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
@@ -448,26 +504,100 @@ func (self *SApsaraClient) GetElbMetrics(opts *cloudprovider.MetricListOptions)
switch opts.MetricType {
case cloudprovider.LB_METRIC_TYPE_NET_BPS_RX:
metricTags = map[string]string{
"InstanceTrafficRX": "",
"TrafficRXNew": "",
}
case cloudprovider.LB_METRIC_TYPE_NET_BPS_TX:
metricTags = map[string]string{
"InstanceTrafficTX": "",
"TrafficTXNew": "",
}
case cloudprovider.LB_METRIC_TYPE_HRSP_COUNT:
case cloudprovider.LB_METRIC_TYPE_DROP_PACKET_RX:
metricTags = map[string]string{
"InstanceStatusCode2xx": cloudprovider.METRIC_TAG_REQUST_2XX,
"InstanceStatusCode3xx": cloudprovider.METRIC_TAG_REQUST_3XX,
"InstanceStatusCode4xx": cloudprovider.METRIC_TAG_REQUST_4XX,
"InstanceStatusCode5xx": cloudprovider.METRIC_TAG_REQUST_5XX,
"InstanceDropPacketRX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_PACKET_TX:
metricTags = map[string]string{
"InstanceDropPacketTX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_TRAFFIC_RX:
metricTags = map[string]string{
"InstanceDropTrafficRX": "",
}
case cloudprovider.LB_METRIC_TYPE_DROP_TRAFFIC_TX:
metricTags = map[string]string{
"InstanceDropTrafficTX": "",
}
case cloudprovider.LB_METRIC_TYPE_UNHEALTHY_SERVER_COUNT:
metricTags = map[string]string{
"UnhealthyServerCount": "",
}
case cloudprovider.LB_METRIC_TYPE_NET_INACTIVE_CONNECTION:
metricTags = map[string]string{
"InactiveConnection": "",
}
case cloudprovider.LB_METRIC_TYPE_MAX_CONNECTION:
metricTags = map[string]string{
"MaxConnection": "",
}
case cloudprovider.LB_METRIC_TYPE_NET_PACKET_RX:
metricTags = map[string]string{
"PacketRX": "",
}
case cloudprovider.LB_METRIC_TYPE_NET_PACKET_TX:
metricTags = map[string]string{
"PacketTX": "",
}
tagKey = cloudprovider.METRIC_TAG_REQUST
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_slb_dashboard", metric, opts.StartTime, opts.EndTime)
result, err := self.ListMetrics("acs_slb_dashboard", metric, "", opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue
}
tags := map[string]string{}
if len(tag) > 0 && len(tagKey) > 0 {
tags[tagKey] = tag
}
for i := range result {
ret = append(ret, cloudprovider.MetricValues{
Id: result[i].InstanceId,
MetricType: opts.MetricType,
Values: []cloudprovider.MetricValue{
{
Timestamp: time.UnixMilli(result[i].Timestamp),
Value: result[i].GetValue(),
Tags: tags,
},
},
})
}
}
return ret, nil
}
func (self *SApsaraClient) GetEipMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
metricTags, tagKey := map[string]string{}, ""
switch opts.MetricType {
case cloudprovider.EIP_METRIC_TYPE_NET_BPS_RX:
metricTags = map[string]string{
"net_rx.rate": "",
}
case cloudprovider.EIP_METRIC_TYPE_NET_BPS_TX:
metricTags = map[string]string{
"net_tx.rate": "",
}
case cloudprovider.EIP_METRIC_TYPE_NET_DROP_SPEED_TX:
metricTags = map[string]string{
"out_ratelimit_drop_speed": "",
}
default:
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "%s", opts.MetricType)
}
ret := []cloudprovider.MetricValues{}
for metric, tag := range metricTags {
result, err := self.ListMetrics("acs_vpc_eip", metric, "", opts.StartTime, opts.EndTime)
if err != nil {
log.Errorf("ListMetric(%s) error: %v", metric, err)
continue