Merge pull request #3235 from Mjoycarry/feature/unforced_dispersion

Feature: Unforced Instance Group
This commit is contained in:
Zexi Li
2019-10-17 13:34:28 +08:00
committed by GitHub
11 changed files with 263 additions and 42 deletions
+33 -4
View File
@@ -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
})
}
+10 -8
View File
@@ -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) {
+7 -6
View File
@@ -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)
+5 -9
View File
@@ -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 {
+1 -1
View File
@@ -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)
@@ -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()
}
@@ -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)
+2 -1
View File
@@ -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{}),
)
}
+1 -3
View File
@@ -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 {
+101 -4
View File
@@ -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
}
+10 -5
View File
@@ -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