From d25ed285048a7f5138205703366c0ff56a4e49df Mon Sep 17 00:00:00 2001 From: mhf Date: Tue, 25 Oct 2022 10:31:42 +0800 Subject: [PATCH] controlMetricsGoroutineCount --- pkg/cloudmon/options/options.go | 3 +++ pkg/cloudmon/providerdriver/base.go | 24 ++++++++++++++++++++++++ pkg/cloudmon/resources/resources.go | 4 ++++ 3 files changed, 31 insertions(+) diff --git a/pkg/cloudmon/options/options.go b/pkg/cloudmon/options/options.go index 38c0dc80e3..90cd48a80e 100644 --- a/pkg/cloudmon/options/options.go +++ b/pkg/cloudmon/options/options.go @@ -35,6 +35,9 @@ type CloudMonOptions struct { SkipMetricPullProviders string `help:"Skip indicate provider metric pull" default:""` InfluxDatabase string `help:"influxdb database name, default telegraf" default:"telegraf"` + + CloudAccountCollectMetricsBatchCount int `help:"Cloud Account Collect Metrics Batch Count" default:"10"` + CloudResourceCollectMetricsBatchCount int `help:"Cloud Resource Collect Metrics BatchC ount" default:"40"` } type PingProbeOptions struct { diff --git a/pkg/cloudmon/providerdriver/base.go b/pkg/cloudmon/providerdriver/base.go index c5e41e9eab..b9dcc3c6fd 100644 --- a/pkg/cloudmon/providerdriver/base.go +++ b/pkg/cloudmon/providerdriver/base.go @@ -137,14 +137,18 @@ type SCollectByResourceIdDriver struct { } func (self *SCollectByResourceIdDriver) CollectDBInstanceMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.DBInstanceDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(rds api.DBInstanceDetails) { defer func() { wg.Done() + <-ch }() opts := &cloudprovider.MetricListOptions{ ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_RDS, @@ -207,14 +211,18 @@ func (self *SCollectByResourceIdDriver) CollectDBInstanceMetrics(ctx context.Con } func (self *SCollectByResourceIdDriver) CollectServerMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ServerDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(vm api.ServerDetails) { defer func() { wg.Done() + <-ch }() opts := &cloudprovider.MetricListOptions{ ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_SERVER, @@ -288,14 +296,18 @@ func (self *SCollectByResourceIdDriver) CollectServerMetrics(ctx context.Context } func (self *SCollectByResourceIdDriver) CollectHostMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.HostDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(vm api.HostDetails) { defer func() { wg.Done() + <-ch }() opts := &cloudprovider.MetricListOptions{ ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_HOST, @@ -367,14 +379,18 @@ func (self *SCollectByResourceIdDriver) CollectHostMetrics(ctx context.Context, } func (self *SCollectByResourceIdDriver) CollectRedisMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.ElasticcacheDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(vm api.ElasticcacheDetails) { defer func() { wg.Done() + <-ch }() opts := &cloudprovider.MetricListOptions{ ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_REDIS, @@ -446,14 +462,18 @@ func (self *SCollectByResourceIdDriver) CollectRedisMetrics(ctx context.Context, } func (self *SCollectByResourceIdDriver) CollectBucketMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.BucketDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(vm api.BucketDetails) { defer func() { wg.Done() + <-ch }() opts := &cloudprovider.MetricListOptions{ ResourceType: cloudprovider.METRIC_RESOURCE_TYPE_BUCKET, @@ -524,15 +544,19 @@ func (self *SCollectByResourceIdDriver) CollectBucketMetrics(ctx context.Context } func (self *SCollectByResourceIdDriver) CollectK8sMetrics(ctx context.Context, manager api.CloudproviderDetails, provider cloudprovider.ICloudProvider, res map[string]api.KubeClusterDetails, start, end time.Time) error { + ch := make(chan struct{}, options.Options.CloudResourceCollectMetricsBatchCount) + defer close(ch) metrics := []influxdb.SMetricData{} s := auth.GetAdminSession(ctx, options.Options.Region) var wg sync.WaitGroup var mu sync.Mutex for i := range res { + ch <- struct{}{} wg.Add(1) go func(vm api.KubeClusterDetails) { defer func() { wg.Done() + <-ch }() // 未同步到本地k8s集群 if len(vm.ExternalClusterId) == 0 { diff --git a/pkg/cloudmon/resources/resources.go b/pkg/cloudmon/resources/resources.go index da688371af..e49515f8e5 100644 --- a/pkg/cloudmon/resources/resources.go +++ b/pkg/cloudmon/resources/resources.go @@ -582,6 +582,8 @@ func (self *SResources) CollectMetrics(ctx context.Context, userCred mcclient.To if isStart { return } + ch := make(chan struct{}, options.Options.CloudAccountCollectMetricsBatchCount) + defer close(ch) s := auth.GetAdminSession(context.Background(), options.Options.Region) resources := self.Cloudproviders.getResources("") cloudproviders := map[string]api.CloudproviderDetails{} @@ -591,10 +593,12 @@ func (self *SResources) CollectMetrics(ctx context.Context, userCred mcclient.To _startTime := _endTime.Add(-1 * time.Minute * time.Duration(options.Options.CollectMetricInterval)) var wg sync.WaitGroup for i := range cloudproviders { + ch <- struct{}{} wg.Add(1) go func(manager api.CloudproviderDetails) { defer func() { wg.Done() + <-ch }() if strings.Contains(strings.ToLower(options.Options.SkipMetricPullProviders), strings.ToLower(manager.Provider)) {