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