From 0c93b0c7d464651d73226fa1e9b5a65be2e9af08 Mon Sep 17 00:00:00 2001 From: Rain Date: Fri, 11 Oct 2019 20:27:39 +0800 Subject: [PATCH] Feature: Unforced Group which prioritize scheduling guest to host with enough resources of InstanceGroup and then to host even if this hosts has no enough resources. Add groups desc in SGuest.ToSchedDesc. The backup guest also occupies InstanceGroup resources now. --- cmd/climc/shell/instance_group.go | 37 +++++- pkg/compute/models/groups.go | 18 +-- pkg/compute/models/guests.go | 13 ++- pkg/compute/models/hosts.go | 14 +-- pkg/mcclient/modules/mod_instance_group.go | 2 +- .../predicates/instance_group_predicate.go | 86 +++++++++++++- .../algorithm/predicates/predicates.go | 8 ++ pkg/scheduler/algorithmprovider/defaults.go | 3 +- pkg/scheduler/api/sched.go | 4 +- pkg/scheduler/core/context.go | 105 +++++++++++++++++- pkg/scheduler/core/generic_scheduler.go | 15 ++- 11 files changed, 263 insertions(+), 42 deletions(-) diff --git a/cmd/climc/shell/instance_group.go b/cmd/climc/shell/instance_group.go index 2f0b0dbad4..a4ae8096b6 100644 --- a/cmd/climc/shell/instance_group.go +++ b/cmd/climc/shell/instance_group.go @@ -15,6 +15,8 @@ package shell import ( + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modules" "yunion.io/x/onecloud/pkg/mcclient/options" @@ -62,10 +64,11 @@ func init() { NAME string `help:"name of instance group"` ZONEID string `help:"zone id" json:"zone_id"` - ServiceType string `help:"service type"` - ParentId string `help:"parent id"` - SchedStrategy string `help:"scheduler strategy"` - Granularity string `help:"the upper limit number of guests with this group in a host"` + ServiceType string `help:"service type"` + ParentId string `help:"parent id"` + SchedStrategy string `help:"scheduler strategy"` + Granularity string `help:"the upper limit number of guests with this group in a host"` + ForceDispersion bool `help:"force to make guest dispersion"` } R(&InstanceGroupCreateOptions{}, "instance-group-create", "Create a instance group", @@ -94,4 +97,30 @@ func init() { }, ) + type InstanceGroupUpdateOptions struct { + ID string `help:"ID or Name of servers to update" json:"-"` + Name string `help:"New name to change"` + Granularity string `help:"the upper limit number of guests with this group in a host"` + ForceDispersion string `help:"force to make guest dispersion" choices:"yes|no" json:"-"` + } + + R(&InstanceGroupUpdateOptions{}, "instance-group-update", "update a instance group", + func(s *mcclient.ClientSession, args *InstanceGroupUpdateOptions) error { + params, err := options.StructToParams(args) + if err != nil { + return err + } + if args.ForceDispersion == "yes" { + params.Set("force_dispersion", jsonutils.JSONTrue) + } else { + params.Set("force_dispersion", jsonutils.JSONFalse) + } + ret, err := modules.InstanceGroup.Update(s, args.ID, params) + if err != nil { + return err + } + printObject(ret) + return nil + }) + } diff --git a/pkg/compute/models/groups.go b/pkg/compute/models/groups.go index c932398e5b..90bb33aca0 100644 --- a/pkg/compute/models/groups.go +++ b/pkg/compute/models/groups.go @@ -14,7 +14,11 @@ package models -import "yunion.io/x/onecloud/pkg/cloudcommon/db" +import ( + "yunion.io/x/pkg/tristate" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" +) const ( REDIS_TYPE = "REDIS" @@ -44,16 +48,14 @@ func init() { type SGroup struct { db.SVirtualResourceBase - ServiceType string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` // Column(VARCHAR(36, charset='ascii'), nullable=True) - - ParentId string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` // Column(VARCHAR(36, charset='ascii'), nullable=True) - - ZoneId string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"required"` // Column(VARCHAR(36, charset='ascii'), nullable=True) - + ServiceType string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` // Column(VARCHAR(36, charset='ascii'), nullable=True) + ParentId string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` // Column(VARCHAR(36, charset='ascii'), nullable=True) + ZoneId string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"required"` // Column(VARCHAR(36, charset='ascii'), nullable=True) SchedStrategy string `width:"16" charset:"ascii" nullable:"true" default:"" list:"user" update:"user" create:"optional"` // Column(VARCHAR(16, charset='ascii'), nullable=True, default='') // the upper limit number of guests with this group in a host - Granularity int `nullable:"false" list:"user" get:"user" create:"optional" default:"1"` + Granularity int `nullable:"false" list:"user" get:"user" create:"optional" update:"user" default:"1"` + ForceDispersion tristate.TriState `list:"user" get:"user" create:"optional" update:"user" default:"true"` } func (group *SGroup) GetNetworks() ([]SGroupnetwork, error) { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index b09d31aede..5190ca41a7 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -4378,7 +4378,7 @@ func (self *SGuest) ToSchedDesc() *schedapi.ScheduleInput { ServerConfigs: new(api.ServerConfigs), } desc.Id = self.Id - //self.FillGroupSchedDesc(desc) + self.FillGroupSchedDesc(config.ServerConfigs) self.FillDiskSchedDesc(config.ServerConfigs) self.FillNetSchedDesc(config.ServerConfigs) if len(self.HostId) > 0 && regutils.MatchUUID(self.HostId) { @@ -4396,18 +4396,19 @@ func (self *SGuest) ToSchedDesc() *schedapi.ScheduleInput { return desc } -/*func (self *SGuest) FillGroupSchedDesc(desc *schedapi.ServerConfig) { +func (self *SGuest) FillGroupSchedDesc(desc *api.ServerConfigs) { groups := make([]SGroupguest, 0) err := GroupguestManager.Query().Equals("guest_id", self.Id).All(&groups) if err != nil { log.Errorln(err) return } - for i := 0; i < len(groups); i++ { - desc.Set(fmt.Sprintf("srvtag.%d", i), - jsonutils.NewString(fmt.Sprintf("%s:%s", groups[i].SrvtagId, groups[i].Tag))) + groupids := make([]string, len(groups)) + for i := range groups { + groupids[i] = groups[i].GroupId } -}*/ + desc.InstanceGroupIds = groupids +} func (self *SGuest) FillDiskSchedDesc(desc *api.ServerConfigs) { guestDisks := make([]SGuestdisk, 0) diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 279d5d4117..80dfcba67f 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -4340,15 +4340,11 @@ func (host *SHost) IsMaintaining() bool { // InstanceGroups returns the group of guest in host and their frequency of occurrence func (host *SHost) InstanceGroups() ([]SGroup, map[string]int, error) { - guests := host.GetGuests() - if len(guests) == 0 { - return []SGroup{}, make(map[string]int), nil - } - guestIds := make([]string, len(guests)) - for i := range guests { - guestIds[i] = guests[i].GetId() - } - q := GroupguestManager.Query().In("guest_id", guestIds) + q := GuestManager.Query("id") + guestQ := q.Filter(sqlchemy.OR(sqlchemy.Equals(q.Field("host_id"), host.Id), + sqlchemy.Equals(q.Field("backup_host_id"), host.Id))).SubQuery() + groupQ := GroupguestManager.Query().SubQuery() + q = groupQ.Query().Join(guestQ, sqlchemy.Equals(guestQ.Field("id"), groupQ.Field("guest_id"))) groupguests := make([]SGroupguest, 0, 1) err := db.FetchModelObjects(GroupguestManager, q, &groupguests) if err != nil { diff --git a/pkg/mcclient/modules/mod_instance_group.go b/pkg/mcclient/modules/mod_instance_group.go index 8b5c816659..bfbef3ffe7 100644 --- a/pkg/mcclient/modules/mod_instance_group.go +++ b/pkg/mcclient/modules/mod_instance_group.go @@ -23,7 +23,7 @@ var ( func init() { InstanceGroup = NewComputeManager("instancegroup", "instancegroups", []string{"ID", "Name", "Service_Type", "Parent_Id", "Zone_Id", "Sched_Strategy", "Domain_Id", "Project_Id", - "Granularity"}, + "Granularity", "Is_Froced_Sep"}, []string{}) registerCompute(&InstanceGroup) diff --git a/pkg/scheduler/algorithm/predicates/instance_group_predicate.go b/pkg/scheduler/algorithm/predicates/instance_group_predicate.go index decdbec59d..474b7bf7ec 100644 --- a/pkg/scheduler/algorithm/predicates/instance_group_predicate.go +++ b/pkg/scheduler/algorithm/predicates/instance_group_predicate.go @@ -42,12 +42,35 @@ func (p *InstanceGroupPredicate) PreExecute(u *core.Unit, cs []core.Candidater) } func (p *InstanceGroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { + return true, nil, nil +} + +type SForcedGroupPredicate struct { + InstanceGroupPredicate +} + +func (p *SForcedGroupPredicate) Name() string { + return "forced_instance_group" +} + +func (p *SForcedGroupPredicate) Clone() core.FitPredicate { + return &SForcedGroupPredicate{} +} + +// SForcedGroupPredicate make sure that there is no more guest with same group whose IsForcedSpe is ture in a host +// for all forced groups in u.SchedData.InstanceGroupIds, so that the capacity is the min value of the FreeGroupCounts +func (p *SForcedGroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { h := NewPredicateHelper(p, u, c) schedDate := u.SchedData() instanceGroups := c.Getter().InstanceGroups() minFree := math.MaxInt16 for _, id := range schedDate.InstanceGroupIds { + detail := schedDate.InstanceGroupsDetail[id] + // SForcedGroupPredicate only deal with group whose ForceDispersion is ture + if detail.ForceDispersion.IsFalse() { + continue + } var free int if _, ok := instanceGroups[id]; ok { free, _ = c.Getter().GetFreeGroupCount(id) @@ -57,7 +80,6 @@ func (p *InstanceGroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, break } } else { - detail := schedDate.InstanceGroupsDetail[id] free = detail.Granularity } if free < minFree { @@ -68,3 +90,65 @@ func (p *InstanceGroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, h.SetCapacity(int64(minFree)) return h.GetResult() } + +type SUnForcedGroupPredicate struct { + InstanceGroupPredicate +} + +func (p *SUnForcedGroupPredicate) Name() string { + return "unforced_instance_group" +} + +func (p *SUnForcedGroupPredicate) Clone() core.FitPredicate { + return &SUnForcedGroupPredicate{} +} + +func (p *SUnForcedGroupPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) { + ret, err := p.InstanceGroupPredicate.PreExecute(u, cs) + if err != nil || !ret { + return ret, err + } + u.RegisterSelectPriorityUpdater(p.Name(), func(u *core.Unit, origin core.SSelectPriorityValue, + hostID string) core.SSelectPriorityValue { + + return origin.SubOne() + }) + return ret, err +} + +// SUnForcedGroupPredicate make sure that the guests are assigned to these hosts who has enough FreeGroupCount of +// unforced groups, so that it will improve the priority of these hosts meet the conditions and the priority should +// be the max value of the FreeGroupCounts +func (p *SUnForcedGroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, + error) { + + h := NewPredicateHelper(p, u, c) + schedDate := u.SchedData() + + instanceGroups := c.Getter().InstanceGroups() + maxPriority := 0 + for _, id := range schedDate.InstanceGroupIds { + detail := schedDate.InstanceGroupsDetail[id] + // SUnForcedGroupPredicate only deal with group whose ForceDispersion is false + if detail.ForceDispersion.IsTrue() { + continue + } + var priority int + if _, ok := instanceGroups[id]; ok { + free, _ := c.Getter().GetFreeGroupCount(id) + if free < 1 { + priority = 0 + } + priority = free + } else { + priority = detail.Granularity + } + if priority > maxPriority { + maxPriority = priority + } + } + + // set priority + h.SetSelectPriority(maxPriority) + return h.GetResult() +} diff --git a/pkg/scheduler/algorithm/predicates/predicates.go b/pkg/scheduler/algorithm/predicates/predicates.go index 6a070e408b..f0b16b9a20 100644 --- a/pkg/scheduler/algorithm/predicates/predicates.go +++ b/pkg/scheduler/algorithm/predicates/predicates.go @@ -165,6 +165,14 @@ func (h *PredicateHelper) SetCapacityCounter(counter core.Counter) { h.Unit.SetCapacity(h.Candidate.IndexKey(), h.predicate.Name(), counter) } +func (h *PredicateHelper) SetSelectPriority(sp int) { + if sp < 0 { + sp = 0 + } + + h.Unit.SetSelectPriorityWithLock(h.Candidate.IndexKey(), h.predicate.Name(), core.SSelectPriorityValue(sp)) +} + func (h *PredicateHelper) Exclude(reason string) { h.SetCapacity(0) h.AppendPredicateFailMsg(reason) diff --git a/pkg/scheduler/algorithmprovider/defaults.go b/pkg/scheduler/algorithmprovider/defaults.go index 5536f41db7..64450f1595 100644 --- a/pkg/scheduler/algorithmprovider/defaults.go +++ b/pkg/scheduler/algorithmprovider/defaults.go @@ -45,7 +45,8 @@ func defaultPredicates() sets.String { factory.RegisterFitPredicate("m-GuestDiskschedtagFilter", &predicates.DiskSchedtagPredicate{}), factory.RegisterFitPredicate("n-ServerSkuFilter", &predicates.InstanceTypePredicate{}), factory.RegisterFitPredicate("o-GuestNetschedtagFilter", &predicates.NetworkSchedtagPredicate{}), - factory.RegisterFitPredicate("p-GuestDispersionFilter", &predicates.InstanceGroupPredicate{}), + factory.RegisterFitPredicate("p-GuestForcedDispersionFilter", &predicates.SForcedGroupPredicate{}), + factory.RegisterFitPredicate("p-GuestUnForcedDispersionFilter", &predicates.SUnForcedGroupPredicate{}), ) } diff --git a/pkg/scheduler/api/sched.go b/pkg/scheduler/api/sched.go index 37b33ca96b..e6cc3a77d9 100644 --- a/pkg/scheduler/api/sched.go +++ b/pkg/scheduler/api/sched.go @@ -143,9 +143,7 @@ func (data *SchedInfo) reviseData() { } func (d *SchedInfo) SkipDirtyMarkHost() bool { - skipByHypervisor := d.IsContainer || d.Hypervisor == SchedTypeContainer - skipByBackup := d.Backup - return skipByHypervisor || skipByBackup + return d.IsContainer || d.Hypervisor == SchedTypeContainer } func (d *SchedInfo) GetCandidateHostTypes() []string { diff --git a/pkg/scheduler/core/context.go b/pkg/scheduler/core/context.go index 13f63679bd..383c9a26ea 100644 --- a/pkg/scheduler/core/context.go +++ b/pkg/scheduler/core/context.go @@ -19,12 +19,12 @@ import ( "sort" "strings" "sync" - "yunion.io/x/onecloud/pkg/apis/compute" - "yunion.io/x/onecloud/pkg/compute/models" - "yunion.io/x/pkg/tristate" "yunion.io/x/log" + "yunion.io/x/pkg/tristate" + "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/scheduler/api" "yunion.io/x/onecloud/pkg/scheduler/core/score" ) @@ -35,7 +35,8 @@ const ( ) var ( - EmptyCapacities = make(map[string]Counter) + EmptyCapacities = make(map[string]Counter) + EmptySelectPriorityValue = SSelectPriorityValue(0) ) type SharedResourceManager struct { @@ -402,11 +403,17 @@ type Unit struct { LogManager *SchedLogManager AllocatedResources map[string]*AllocatedResource + + SelectPriorityMap map[string]SSelectPriority + SelectPriorityUpdaterMap map[string]SSelectPriorityUpdater + SelectPriorityLock sync.Mutex } func NewScheduleUnit(info *api.SchedInfo, schedManager interface{}) *Unit { cmap := make(map[string]*Capacity) // candidate_id, Capacity smap := make(map[string]Score) // candidate_id, Score + spmap := make(map[string]SSelectPriority) + spumap := make(map[string]SSelectPriorityUpdater) unit := &Unit{ SchedInfo: info, FailedCandidateMap: make(map[string]*FailedCandidates), @@ -421,6 +428,9 @@ func NewScheduleUnit(info *api.SchedInfo, schedManager interface{}) *Unit { LogManager: NewSchedLogManager(), SchedulerManager: schedManager, AllocatedResources: make(map[string]*AllocatedResource), + + SelectPriorityMap: spmap, + SelectPriorityUpdaterMap: spumap, } return unit } @@ -569,6 +579,52 @@ func (u *Unit) SetCapacity(id string, name string, capacity Counter) error { return nil } +func (u *Unit) GetSelectPriority(id string) SSelectPriorityValue { + if sp, ok := u.SelectPriorityMap[id]; ok { + return sp.Value() + } + return EmptySelectPriorityValue +} + +func (u *Unit) SetSelectPriorityWithLock(id string, name string, spv SSelectPriorityValue) { + u.SelectPriorityLock.Lock() + defer u.SelectPriorityLock.Unlock() + + sp, ok := u.SelectPriorityMap[id] + if !ok { + sp = NewSSelctPriority() + u.SelectPriorityMap[id] = sp + } + + sp[name] = spv +} + +func (u *Unit) UpdateSelectPriority() { + for hostID, sp := range u.SelectPriorityMap { + for name, spv := range sp { + sp[name] = u.SelectPriorityUpdaterMap[name](u, spv, hostID) + } + } +} + +func (u *Unit) GetMaxSelectPriority() (max SSelectPriorityValue) { + max = EmptySelectPriorityValue + for _, sp := range u.SelectPriorityMap { + val := sp.Value() + if max.Less(val) { + max = val + } + } + return +} + +func (u *Unit) RegisterSelectPriorityUpdater(name string, f SSelectPriorityUpdater) { + u.SelectPriorityLock.Lock() + defer u.SelectPriorityLock.Unlock() + + u.SelectPriorityUpdaterMap[name] = f +} + func validateCapacityInput(c Counter) bool { if c != nil && c.GetCount() >= 0 { return true @@ -698,3 +754,44 @@ func (u *Unit) GetAllocatedResource(candidateId string) *AllocatedResource { } return ret } + +type SSelectPriority map[string]SSelectPriorityValue + +func (s SSelectPriority) Value() (val SSelectPriorityValue) { + val = EmptySelectPriorityValue + for _, v := range s { + if v > val { + val = v + } + } + return +} + +func NewSSelctPriority() SSelectPriority { + return make(map[string]SSelectPriorityValue) +} + +// SSelectPriorityUpdater will call to update the specified host after each round of selection +type SSelectPriorityUpdater func(u *Unit, origin SSelectPriorityValue, hostID string) SSelectPriorityValue + +type SSelectPriorityValue int + +func (s SSelectPriorityValue) Less(sp SSelectPriorityValue) bool { + return s < sp +} + +func (s SSelectPriorityValue) Sub(sp SSelectPriorityValue) (ret SSelectPriorityValue) { + ret = s - sp + if ret.Less(EmptySelectPriorityValue) { + ret = EmptySelectPriorityValue + } + return +} + +func (s SSelectPriorityValue) SubOne() SSelectPriorityValue { + return s.Sub(SSelectPriorityValue(1)) +} + +func (s SSelectPriorityValue) IsEmpty() bool { + return s == EmptySelectPriorityValue +} diff --git a/pkg/scheduler/core/generic_scheduler.go b/pkg/scheduler/core/generic_scheduler.go index d7c0c9e7a9..79b48f4845 100644 --- a/pkg/scheduler/core/generic_scheduler.go +++ b/pkg/scheduler/core/generic_scheduler.go @@ -349,12 +349,17 @@ func SelectHosts(unit *Unit, priorityList HostPriorityList) ([]*SelectedCandidat completed: for len(priorityList) > 0 { log.V(10).Debugf("PriorityList: %#v", priorityList) + currentPriority := unit.GetMaxSelectPriority() priorityList0 := HostPriorityList{} for _, it := range priorityList { if count <= 0 { break completed } hostID := it.Host + if !currentPriority.IsEmpty() && unit.GetSelectPriority(hostID).Less(currentPriority) { + priorityList0 = append(priorityList0, it) + continue + } var ( selectedItem *SelectedCandidate ok bool @@ -373,9 +378,12 @@ completed: priorityList0 = append(priorityList0, it) } } + if !currentPriority.IsEmpty() { + unit.UpdateSelectPriority() + } // sort by score priorityList = priorityList0 - sort.Sort(sort.Reverse(priorityList)) + //sort.Sort(sort.Reverse(priorityList)) } for _, sc := range selectedMap { @@ -581,10 +589,7 @@ func PrioritizeCandidates( } wg := sync.WaitGroup{} - results := make([]HostPriorityList, 0, len(priorities)) - for range priorities { - results = append(results, nil) - } + results := make([]HostPriorityList, len(priorities)) newPriorities, err := preExecPriorities(priorities, unit, candidates) if err != nil { return nil, err