From 38da40461d4a1bce81e5113baafc8951207e9f76 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Thu, 21 Mar 2019 16:36:11 +0800 Subject: [PATCH] scheduler: server sku predicate --- pkg/apis/compute/api.go | 2 +- pkg/mcclient/options/schedulers.go | 10 +- .../predicates/disk_schedtag_predicate.go | 59 +++++---- .../algorithm/predicates/sku_predicate.go | 44 +++++++ pkg/scheduler/algorithmprovider/defaults.go | 1 + pkg/scheduler/cache/candidate/hosts.go | 1 + pkg/scheduler/data_manager/sku/sku.go | 121 ++++++++++++++++++ pkg/scheduler/options/options.go | 2 + pkg/scheduler/service/service.go | 4 + 9 files changed, 212 insertions(+), 32 deletions(-) create mode 100644 pkg/scheduler/algorithm/predicates/sku_predicate.go create mode 100644 pkg/scheduler/data_manager/sku/sku.go diff --git a/pkg/apis/compute/api.go b/pkg/apis/compute/api.go index a753ce43a5..c034a98412 100644 --- a/pkg/apis/compute/api.go +++ b/pkg/apis/compute/api.go @@ -101,6 +101,7 @@ type ServerConfigs struct { Hypervisor string `json:"hypervisor"` // ResourceType "shared|prepaid|dedicated"` ResourceType string `json:"resource_type"` + InstanceType string `json:"instance_type"` Project string `json:"project"` Backup bool `json:"backup"` Count int `json:"count"` @@ -130,7 +131,6 @@ type ServerCreateInput struct { GenerateName string `json:"generate_name"` VmemSize int `json:"vmem_size"` VcpuCount int `json:"vcpu_count"` - InstanceType string `json:"instance_type"` UserData string `json:"user_data"` Keypair string `json:"keypair"` diff --git a/pkg/mcclient/options/schedulers.go b/pkg/mcclient/options/schedulers.go index d4f3785b44..bf369d6a18 100644 --- a/pkg/mcclient/options/schedulers.go +++ b/pkg/mcclient/options/schedulers.go @@ -8,9 +8,10 @@ import ( type SchedulerTestBaseOptions struct { ServerConfigs - Mem int `help:"Memory size (MB), default 512" metavar:"MEMORY" default:"512"` - Ncpu int `help:"#CPU cores of VM server, default 1" default:"1" metavar:""` - Log bool `help:"Record to schedule history"` + Mem int `help:"Memory size (MB), default 512" metavar:"MEMORY" default:"512"` + Ncpu int `help:"#CPU cores of VM server, default 1" default:"1" metavar:""` + Sku string `help:"Server SKU instance type"` + Log bool `help:"Record to schedule history"` } func (o SchedulerTestBaseOptions) data(s *mcclient.ClientSession) (*scheduler.ServerConfig, error) { @@ -28,6 +29,9 @@ func (o SchedulerTestBaseOptions) data(s *mcclient.ClientSession) (*scheduler.Se if o.Ncpu > 0 { data.Ncpu = o.Ncpu } + if o.Sku != "" { + data.InstanceType = o.Sku + } return data, nil } diff --git a/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go b/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go index ae8c26f893..0011c84d79 100644 --- a/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go +++ b/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go @@ -6,7 +6,6 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/util/errors" - "yunion.io/x/pkg/utils" computeapi "yunion.io/x/onecloud/pkg/apis/compute" schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" @@ -189,10 +188,22 @@ func (p *DiskSchedtagPredicate) Execute(u *core.Unit, c core.Candidater) (bool, h := NewPredicateHelper(p, u, c) storages := c.Getter().Storages() + ds := p.GetDiskStoragesMap(c.IndexKey()) disks := u.SchedData().Disks for idx, d := range disks { - matchedStorages, err := p.checkStorages(d, storages) + fitStorages := make([]*api.CandidateStorage, 0) + for _, s := range storages { + if p.isStorageFitDisk(s, d) { + fitStorages = append(fitStorages, s) + } + } + if len(fitStorages) == 0 { + h.Exclude(fmt.Sprintf("Not found available storages for disk backend %q", d.Backend)) + break + } + + matchedStorages, err := p.checkStorages(d, fitStorages) if err != nil { h.Exclude(err.Error()) } @@ -245,15 +256,13 @@ func (p *DiskSchedtagPredicate) selectStorage(d *computeapi.DiskConfig, storages noTagStorages := []*api.CandidateStorage{} avoidStorages := []*api.CandidateStorage{} for _, storage := range storages { - if p.isStorageFitDisk(storage, d) { - candi := storage.CandidateStorage - if storage.isNoTag() { - noTagStorages = append(noTagStorages, candi) - } else if storage.hasPreferTags() { - preferStorages = append(preferStorages, candi) - } else if storage.hasAvoidTags() { - avoidStorages = append(avoidStorages, candi) - } + candi := storage.CandidateStorage + if storage.isNoTag() { + noTagStorages = append(noTagStorages, candi) + } else if storage.hasPreferTags() { + preferStorages = append(preferStorages, candi) + } else if storage.hasAvoidTags() { + avoidStorages = append(avoidStorages, candi) } } sortStorages := []*api.CandidateStorage{} @@ -264,28 +273,22 @@ func (p *DiskSchedtagPredicate) selectStorage(d *computeapi.DiskConfig, storages } func (p *DiskSchedtagPredicate) GetLeastUsedStorage(storages []*api.CandidateStorage, backend string) *api.CandidateStorage { - var backends []string - if backend == computeapi.STORAGE_LOCAL { - backends = []string{computeapi.STORAGE_NAS, computeapi.STORAGE_LOCAL} - } else if len(backend) > 0 { - backends = []string{backend} - } else { - backends = []string{} - } - return p.getLeastUsedStorage(storages, backends) + //var backends []string + //if backend == computeapi.STORAGE_LOCAL { + //backends = []string{computeapi.STORAGE_NAS, computeapi.STORAGE_LOCAL} + //} else if len(backend) > 0 { + //backends = []string{backend} + //} else { + //backends = []string{} + //} + return p.getLeastUsedStorage(storages) } -func (p *DiskSchedtagPredicate) getLeastUsedStorage(storages []*api.CandidateStorage, backends []string) *api.CandidateStorage { +func (p *DiskSchedtagPredicate) getLeastUsedStorage(storages []*api.CandidateStorage) *api.CandidateStorage { var best *api.CandidateStorage var bestCap int for i := 0; i < len(storages); i++ { s := storages[i] - if len(backends) > 0 { - in, _ := utils.InStringArray(s.StorageType, backends) - if !in { - continue - } - } capa := s.GetFreeCapacity() if best == nil || bestCap < capa { bestCap = capa @@ -299,7 +302,7 @@ func (p *DiskSchedtagPredicate) GetHypervisorDriver() models.IGuestDriver { return models.GetDriver(p.Hypervisor) } -func (p *DiskSchedtagPredicate) isStorageFitDisk(storage *PredicatedStorage, d *computeapi.DiskConfig) bool { +func (p *DiskSchedtagPredicate) isStorageFitDisk(storage *api.CandidateStorage, d *computeapi.DiskConfig) bool { if d.Storage != "" { if storage.Id == d.Storage || storage.Name == d.Storage { return true diff --git a/pkg/scheduler/algorithm/predicates/sku_predicate.go b/pkg/scheduler/algorithm/predicates/sku_predicate.go new file mode 100644 index 0000000000..5e6273fb1f --- /dev/null +++ b/pkg/scheduler/algorithm/predicates/sku_predicate.go @@ -0,0 +1,44 @@ +package predicates + +import ( + "fmt" + + "yunion.io/x/onecloud/pkg/scheduler/core" + skuman "yunion.io/x/onecloud/pkg/scheduler/data_manager/sku" +) + +type InstanceTypePredicate struct { + BasePredicate +} + +func (p *InstanceTypePredicate) Name() string { + return "instance_type" +} + +func (p *InstanceTypePredicate) Clone() core.FitPredicate { + return &InstanceTypePredicate{} +} + +func (p *InstanceTypePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) { + if u.SchedData().InstanceType == "" || !u.IsPublicCloudProvider() { + return false, nil + } + return true, nil +} + +func (p *InstanceTypePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { + h := NewPredicateHelper(p, u, c) + + d := u.SchedData() + + zoneId := c.Getter().Zone().Id + zoneName := c.Getter().Zone().Name + instanceType := d.InstanceType + + sku := skuman.GetByZone(instanceType, zoneId) + if sku == nil { + h.Exclude(fmt.Sprintf("Not found server sku %s at zone %s", instanceType, zoneName)) + } + + return h.GetResult() +} diff --git a/pkg/scheduler/algorithmprovider/defaults.go b/pkg/scheduler/algorithmprovider/defaults.go index 7a1b8ec849..e14acd00e1 100644 --- a/pkg/scheduler/algorithmprovider/defaults.go +++ b/pkg/scheduler/algorithmprovider/defaults.go @@ -28,6 +28,7 @@ func defaultPredicates() sets.String { factory.RegisterFitPredicate("k-GuestIsolatedDeviceFilter", &predicateguest.IsolatedDevicePredicate{}), factory.RegisterFitPredicate("l-GuestResourceTypeFilter", &predicates.ResourceTypePredicate{}), factory.RegisterFitPredicate("m-GuestDiskschedtagFilter", &predicates.DiskSchedtagPredicate{}), + factory.RegisterFitPredicate("n-ServerSkuFilter", &predicates.InstanceTypePredicate{}), ) } diff --git a/pkg/scheduler/cache/candidate/hosts.go b/pkg/scheduler/cache/candidate/hosts.go index a77fe42764..b11cc73692 100644 --- a/pkg/scheduler/cache/candidate/hosts.go +++ b/pkg/scheduler/cache/candidate/hosts.go @@ -213,6 +213,7 @@ type HostBuilder struct { cpuIOLoads map[string]map[string]float64 schedtags []computemodels.SSchedtag + zoneSkus map[string][]computemodels.SServerSku } func (h *HostDesc) String() string { diff --git a/pkg/scheduler/data_manager/sku/sku.go b/pkg/scheduler/data_manager/sku/sku.go new file mode 100644 index 0000000000..bea51fe809 --- /dev/null +++ b/pkg/scheduler/data_manager/sku/sku.go @@ -0,0 +1,121 @@ +package sku + +import ( + "fmt" + "sync" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/util/wait" + "yunion.io/x/sqlchemy" + + "yunion.io/x/onecloud/pkg/compute/models" +) + +var ( + skuManager *SSkuManager +) + +func Start(refreshInterval time.Duration) { + skuManager = &SSkuManager{ + skuMap: newSkuMap(), + refreshInterval: refreshInterval, + } + skuManager.sync() +} + +func GetByZone(instanceType, zoneId string) *models.SServerSku { + return skuManager.GetByZone(instanceType, zoneId) +} + +type skuMap struct { + *sync.Map +} + +type skuList []*models.SServerSku + +func (l skuList) Has(newSku *models.SServerSku) (int, bool) { + for i, oldSku := range l { + if oldSku.Id == newSku.Id { + return i, true + } + } + return -1, false +} + +func (l skuList) DebugString() string { + return fmt.Sprintf("%s", jsonutils.Marshal(l).String()) +} + +func (l skuList) GetByZone(zoneId string) *models.SServerSku { + for _, s := range l { + if s.ZoneId == zoneId { + return s + } + } + return nil +} + +func newSkuMap() *skuMap { + return &skuMap{ + Map: new(sync.Map), + } +} + +func (cache *skuMap) Get(instanceType string) skuList { + value, ok := cache.Load(instanceType) + if ok { + return value.(skuList) + } + return nil +} + +func (cache *skuMap) Add(instanceType string, sku *models.SServerSku) { + skus := cache.Get(instanceType) + if skus == nil { + skus = make([]*models.SServerSku, 0) + } + skus = append(skus, sku) + cache.Store(instanceType, skus) +} + +type SSkuManager struct { + // skus cache all server skus in database, key is InstanceType, value is []models.SServerSku + skuMap *skuMap + refreshInterval time.Duration +} + +func (m *SSkuManager) syncOnce() { + log.Infof("SkuManager start sync") + startTime := time.Now() + + skus := make([]models.SServerSku, 0) + q := models.ServerSkuManager.Query() + q = q.Filter( + sqlchemy.OR( + sqlchemy.Equals(q.Field("prepaid_status"), models.SkuStatusAvailable), + sqlchemy.Equals(q.Field("postpaid_status"), models.SkuStatusAvailable))) + if err := q.All(&skus); err != nil { + log.Errorf("SkuManager query all available skus error: %v", err) + return + } + m.skuMap = newSkuMap() + for _, sku := range skus { + tmp := sku + m.skuMap.Add(sku.Name, &tmp) + } + log.Infof("SkuManager end sync, consume %s", time.Since(startTime)) +} + +func (m *SSkuManager) sync() { + wait.Forever(m.syncOnce, m.refreshInterval) +} + +func (m *SSkuManager) GetByZone(instanceType, zoneId string) *models.SServerSku { + l := m.skuMap.Get(instanceType) + if l == nil { + return nil + } + return l.GetByZone(zoneId) +} diff --git a/pkg/scheduler/options/options.go b/pkg/scheduler/options/options.go index ef57948781..ddb0d9caf5 100644 --- a/pkg/scheduler/options/options.go +++ b/pkg/scheduler/options/options.go @@ -80,6 +80,8 @@ type SchedOptions struct { WireDBCacheTTL string `help:"Wire database cache TTL" default:"0s"` WireDBCachePeriod string `help:"Wire database cache period" default:"5m"` + + SkuRefreshInterval string `help:"Server SKU refresh interval" default:"12h"` } var ( diff --git a/pkg/scheduler/service/service.go b/pkg/scheduler/service/service.go index 75d06c2bc2..563142ecbb 100644 --- a/pkg/scheduler/service/service.go +++ b/pkg/scheduler/service/service.go @@ -16,12 +16,15 @@ import ( app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" "yunion.io/x/onecloud/pkg/cloudcommon/db" computemodels "yunion.io/x/onecloud/pkg/compute/models" + skuman "yunion.io/x/onecloud/pkg/scheduler/data_manager/sku" "yunion.io/x/onecloud/pkg/scheduler/db/models" schedhandler "yunion.io/x/onecloud/pkg/scheduler/handler" schedman "yunion.io/x/onecloud/pkg/scheduler/manager" o "yunion.io/x/onecloud/pkg/scheduler/options" "yunion.io/x/onecloud/pkg/util/gin/middleware" + _ "yunion.io/x/onecloud/pkg/compute/guestdrivers" + _ "yunion.io/x/onecloud/pkg/compute/hostdrivers" _ "yunion.io/x/onecloud/pkg/scheduler/algorithmprovider" ) @@ -43,6 +46,7 @@ func StartService() error { } stopEverything := make(chan struct{}) + go skuman.Start(utils.ToDuration(opts.SkuRefreshInterval)) schedman.InitAndStart(stopEverything) }