mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #1295 in YUNIONIO/onecloud from ~LIZEXI/onecloud:feature/lzx-sku-schedule to release/2.8.0
* commit '38da40461d4a1bce81e5113baafc8951207e9f76': scheduler: server sku predicate
This commit is contained in:
@@ -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"`
|
||||
|
||||
@@ -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:"<SERVER_CPU_COUNT>"`
|
||||
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:"<SERVER_CPU_COUNT>"`
|
||||
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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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{}),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
+1
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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 (
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user