mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #15233 from gouqi11/controlMetricsGoroutineCount
controlMetricsGoroutineCount
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)) {
|
||||
|
||||
Reference in New Issue
Block a user