mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix(scheduler): consider storage backend and mediumType size (#15470)
This commit is contained in:
@@ -19,10 +19,14 @@ import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/tristate"
|
||||
"yunion.io/x/pkg/util/sets"
|
||||
"yunion.io/x/pkg/utils"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates"
|
||||
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
|
||||
"yunion.io/x/onecloud/pkg/scheduler/core"
|
||||
)
|
||||
|
||||
@@ -48,6 +52,104 @@ func (p *StoragePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []co
|
||||
return true, nil
|
||||
}
|
||||
|
||||
type diskSizeRequest struct {
|
||||
backend string
|
||||
max int64
|
||||
total int64
|
||||
disks []*compute.DiskConfig
|
||||
}
|
||||
|
||||
func newDiskSizeRequest(backend string) *diskSizeRequest {
|
||||
return &diskSizeRequest{
|
||||
backend: backend,
|
||||
max: -1,
|
||||
total: 0,
|
||||
disks: make([]*compute.DiskConfig, 0),
|
||||
}
|
||||
}
|
||||
|
||||
func (req *diskSizeRequest) GetMax() int64 {
|
||||
return req.max
|
||||
}
|
||||
|
||||
func (req *diskSizeRequest) GetTotal() int64 {
|
||||
return req.total
|
||||
}
|
||||
|
||||
func (req *diskSizeRequest) Add(disk *compute.DiskConfig) *diskSizeRequest {
|
||||
req.total += int64(disk.SizeMb)
|
||||
if req.max < int64(disk.SizeMb) {
|
||||
req.max = int64(disk.SizeMb)
|
||||
}
|
||||
|
||||
req.disks = append(req.disks, disk)
|
||||
return req
|
||||
}
|
||||
|
||||
func (req *diskSizeRequest) NewByMediumType(mt string) *diskSizeRequest {
|
||||
newReq := newDiskSizeRequest(req.backend)
|
||||
for _, d := range req.disks {
|
||||
if d.Medium == mt {
|
||||
newReq.Add(d)
|
||||
}
|
||||
}
|
||||
return newReq
|
||||
}
|
||||
|
||||
type diskBackendSizeRequest struct {
|
||||
reqs map[string]*diskSizeRequest
|
||||
// beMdm is a map that use backend as key and medium types as value
|
||||
beMdm map[string]sets.String
|
||||
}
|
||||
|
||||
func newDiskBackendSizeRequest() *diskBackendSizeRequest {
|
||||
return &diskBackendSizeRequest{
|
||||
reqs: make(map[string]*diskSizeRequest),
|
||||
beMdm: make(map[string]sets.String),
|
||||
}
|
||||
}
|
||||
|
||||
func (ds *diskBackendSizeRequest) get(backend string) (*diskSizeRequest, bool) {
|
||||
req, ok := ds.reqs[backend]
|
||||
return req, ok
|
||||
}
|
||||
|
||||
func (ds *diskBackendSizeRequest) set(backend string, req *diskSizeRequest) *diskBackendSizeRequest {
|
||||
ds.reqs[backend] = req
|
||||
return ds
|
||||
}
|
||||
|
||||
func (ds *diskBackendSizeRequest) Add(disk *compute.DiskConfig) *diskBackendSizeRequest {
|
||||
backend := disk.Backend
|
||||
req, ok := ds.get(backend)
|
||||
if !ok {
|
||||
req = newDiskSizeRequest(backend)
|
||||
}
|
||||
req.Add(disk)
|
||||
ds.set(backend, req)
|
||||
|
||||
mds, ok := ds.beMdm[backend]
|
||||
if !ok {
|
||||
mds = sets.NewString()
|
||||
}
|
||||
mds.Insert(disk.Medium)
|
||||
ds.beMdm[backend] = mds
|
||||
|
||||
return ds
|
||||
}
|
||||
|
||||
func (ds *diskBackendSizeRequest) GetBackendMediumMap() map[string]sets.String {
|
||||
return ds.beMdm
|
||||
}
|
||||
|
||||
func (ds *diskBackendSizeRequest) Get(backend string, mediumType string) (*diskSizeRequest, error) {
|
||||
req, ok := ds.get(backend)
|
||||
if !ok {
|
||||
return nil, errors.Errorf("Not found diskBackendSizeRequest by backend %q", backend)
|
||||
}
|
||||
return req.NewByMediumType(mediumType), nil
|
||||
}
|
||||
|
||||
func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
|
||||
h := predicates.NewPredicateHelper(p, u, c)
|
||||
|
||||
@@ -87,12 +189,15 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can
|
||||
}
|
||||
}
|
||||
|
||||
getStorageCapacity := func(backend string, reqMaxSize int64, reqTotalSize int64, useRsvd bool) (*storageCapacity, *storageCapacity) {
|
||||
totalFree, actualFree := getter.GetFreeStorageSizeOfType(backend, useRsvd)
|
||||
getStorageCapacity := func(backend string, mediumType string, reqMaxSize int64, reqTotalSize int64, useRsvd bool) (*storageCapacity, *storageCapacity, error) {
|
||||
totalFree, actualFree, err := getter.GetFreeStorageSizeOfType(backend, mediumType, useRsvd, reqMaxSize)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
reqTotalSize = utils.Max(reqTotalSize, 1)
|
||||
capacity := totalFree / reqTotalSize
|
||||
actualCapacity := actualFree / reqTotalSize
|
||||
return newStorageCapacity(capacity, totalFree, false), newStorageCapacity(actualCapacity, actualFree, true)
|
||||
return newStorageCapacity(capacity, totalFree, false), newStorageCapacity(actualCapacity, actualFree, true), nil
|
||||
}
|
||||
|
||||
getReqSizeStr := func(backend string) string {
|
||||
@@ -106,27 +211,27 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can
|
||||
return strings.Join(ss, "+")
|
||||
}
|
||||
|
||||
getStorageFreeStr := func(backend string, useRsvd bool, isActual bool) string {
|
||||
getStorageFreeStr := func(backend string, mediumType string, useRsvd bool, isActual bool) string {
|
||||
ss := []string{}
|
||||
for _, s := range getter.Storages() {
|
||||
if s.StorageType == backend {
|
||||
if candidate.IsStorageBackendMediumMatch(s, backend, mediumType) {
|
||||
if isActual {
|
||||
total := s.Capacity
|
||||
free := total - s.ActualCapacityUsed
|
||||
ss = append(ss, fmt.Sprintf("actual_total:%d - actual_used:%d = free:%d", total, s.ActualCapacityUsed, free))
|
||||
ss = append(ss, fmt.Sprintf("storage %q, actual_total:%d - actual_used:%d = free:%d", s.GetName(), total, s.ActualCapacityUsed, free))
|
||||
} else {
|
||||
total := int64(float32(s.Capacity) * s.Cmtbound)
|
||||
used := s.GetUsedCapacity(tristate.True)
|
||||
waste := s.GetUsedCapacity(tristate.False)
|
||||
free := total - int64(used) - int64(waste)
|
||||
ss = append(ss, fmt.Sprintf("total:%d - used:%d - waste:%d = free:%d", total, used, waste, free))
|
||||
ss = append(ss, fmt.Sprintf("storage %q, total:%d - used:%d - waste:%d = free:%d", s.GetName(), total, used, waste, free))
|
||||
}
|
||||
}
|
||||
}
|
||||
return strings.Join(ss, " + ")
|
||||
}
|
||||
|
||||
sizeRequest := make(map[string]map[string]int64, 0)
|
||||
sizeRequest := newDiskBackendSizeRequest()
|
||||
storeRequest := make(map[string]int64, 0)
|
||||
for _, disk := range d.Disks {
|
||||
if isMigrate() && !isLocalhostBackend(disk.Backend) {
|
||||
@@ -136,14 +241,7 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can
|
||||
storeRequest[disk.Storage] = 1
|
||||
} else if !isMigrate() || (isMigrate() && isLocalhostBackend(disk.Backend)) {
|
||||
// if migrate, only local storage need check capacity constraint
|
||||
if _, ok := sizeRequest[disk.Backend]; !ok {
|
||||
sizeRequest[disk.Backend] = map[string]int64{"max": -1, "total": 0}
|
||||
}
|
||||
max := sizeRequest[disk.Backend]["max"]
|
||||
if max < int64(disk.SizeMb) {
|
||||
sizeRequest[disk.Backend]["max"] = int64(disk.SizeMb)
|
||||
}
|
||||
sizeRequest[disk.Backend]["total"] += int64(disk.SizeMb)
|
||||
sizeRequest.Add(disk)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -157,27 +255,38 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can
|
||||
useRsvd := h.UseReserved()
|
||||
minCapacity := int64(0xFFFFFFFF)
|
||||
|
||||
appendFailMsg := func(backend string, req map[string]int64, useRsvd bool, capacity *storageCapacity) {
|
||||
reqStr := fmt.Sprintf("no enough %q storage, req=%v(%v)", backend, req["total"], getReqSizeStr(backend))
|
||||
appendFailMsg := func(backend string, mediumType string, req *diskSizeRequest, useRsvd bool, capacity *storageCapacity) {
|
||||
reqStr := fmt.Sprintf("no enough backend %q, mediumType %q storage, req=%v(%v)", backend, mediumType, req.GetTotal(), getReqSizeStr(backend))
|
||||
freePrex := "free"
|
||||
isActual := capacity.isActual
|
||||
if isActual {
|
||||
freePrex = "actual_free"
|
||||
}
|
||||
freeStr := fmt.Sprintf("%s=%v(%v)", freePrex, capacity.free, getStorageFreeStr(backend, useRsvd, isActual))
|
||||
freeStr := fmt.Sprintf("%s=%v(%v)", freePrex, capacity.free, getStorageFreeStr(backend, mediumType, useRsvd, isActual))
|
||||
msg := reqStr + ", " + freeStr
|
||||
h.AppendPredicateFailMsg(msg)
|
||||
}
|
||||
|
||||
for be, req := range sizeRequest {
|
||||
capacity, actualCapacity := getStorageCapacity(be, req["max"], req["total"], useRsvd)
|
||||
tmpCap := utils.Min(capacity.capacity, actualCapacity.capacity)
|
||||
if capacity.capacity <= 0 {
|
||||
appendFailMsg(be, req, useRsvd, capacity)
|
||||
} else if actualCapacity.capacity <= 0 {
|
||||
appendFailMsg(be, req, useRsvd, actualCapacity)
|
||||
for be, mds := range sizeRequest.GetBackendMediumMap() {
|
||||
for _, medium := range mds.List() {
|
||||
req, err := sizeRequest.Get(be, medium)
|
||||
if err != nil {
|
||||
h.Exclude(fmt.Sprintf("get request size by backend %q, medium %q: %v", be, medium, err))
|
||||
break
|
||||
}
|
||||
capacity, actualCapacity, err := getStorageCapacity(be, medium, req.GetMax(), req.GetTotal(), useRsvd)
|
||||
if err != nil {
|
||||
h.Exclude(err.Error())
|
||||
continue
|
||||
}
|
||||
tmpCap := utils.Min(capacity.capacity, actualCapacity.capacity)
|
||||
if capacity.capacity <= 0 {
|
||||
appendFailMsg(be, medium, req, useRsvd, capacity)
|
||||
} else if actualCapacity.capacity <= 0 {
|
||||
appendFailMsg(be, medium, req, useRsvd, actualCapacity)
|
||||
}
|
||||
minCapacity = utils.Min(minCapacity, tmpCap)
|
||||
}
|
||||
minCapacity = utils.Min(minCapacity, tmpCap)
|
||||
}
|
||||
|
||||
h.SetCapacity(minCapacity)
|
||||
|
||||
+37
-5
@@ -223,16 +223,48 @@ func (b baseHostGetter) TotalMemorySize(_ bool) int64 {
|
||||
return int64(b.h.MemSize)
|
||||
}
|
||||
|
||||
func (b baseHostGetter) GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) {
|
||||
func checkStorageSize(s *api.CandidateStorage, reqMaxSize int64, useRsvd bool) error {
|
||||
storageSize := s.FreeCapacity
|
||||
if useRsvd {
|
||||
storageSize += s.GetReserved()
|
||||
}
|
||||
minSize := utils.Min(storageSize, s.ActualFreeCapacity)
|
||||
if minSize >= reqMaxSize {
|
||||
return nil
|
||||
}
|
||||
return errors.Errorf("storage %q free size %d less than max request %d, use reserverd %v, free_capacity(%d), actual_free_capacity(%d)", s.GetName(), minSize, reqMaxSize, useRsvd, storageSize, s.ActualFreeCapacity)
|
||||
}
|
||||
|
||||
func IsStorageBackendMediumMatch(s *api.CandidateStorage, backend string, mediumType string) bool {
|
||||
if s.StorageType != backend {
|
||||
return false
|
||||
}
|
||||
if mediumType == "" {
|
||||
return true
|
||||
}
|
||||
return s.MediumType == mediumType
|
||||
}
|
||||
|
||||
func (b baseHostGetter) GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) {
|
||||
var size int64
|
||||
var actualSize int64
|
||||
foundLEReqStore := false
|
||||
errs := make([]error, 0)
|
||||
for _, s := range b.Storages() {
|
||||
if s.StorageType == storageType {
|
||||
size += int64(float32(s.Capacity) * s.Cmtbound)
|
||||
actualSize += s.Capacity - s.ActualCapacityUsed
|
||||
if IsStorageBackendMediumMatch(s, storageType, mediumType) {
|
||||
size += s.FreeCapacity
|
||||
actualSize += s.ActualFreeCapacity
|
||||
if err := checkStorageSize(s, reqMaxSize, false); err != nil {
|
||||
errs = append(errs, err)
|
||||
} else {
|
||||
foundLEReqStore = true
|
||||
}
|
||||
}
|
||||
}
|
||||
return size, actualSize
|
||||
if foundLEReqStore {
|
||||
return size, actualSize, nil
|
||||
}
|
||||
return size, actualSize, errors.NewAggregate(errs)
|
||||
}
|
||||
|
||||
func (b baseHostGetter) GetFreePort(netId string) int {
|
||||
|
||||
+21
-10
@@ -23,7 +23,7 @@ import (
|
||||
"time"
|
||||
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/util/errors"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/util/sets"
|
||||
"yunion.io/x/pkg/util/workqueue"
|
||||
"yunion.io/x/pkg/utils"
|
||||
@@ -85,8 +85,8 @@ func (h *hostGetter) StorageInfo() []*baremetal.BaremetalStorage {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *hostGetter) GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) {
|
||||
return h.h.GetFreeStorageSizeOfType(storageType, useRsvd)
|
||||
func (h *hostGetter) GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) {
|
||||
return h.h.GetFreeStorageSizeOfType(storageType, mediumType, useRsvd, reqMaxSize)
|
||||
}
|
||||
|
||||
func (h *hostGetter) GetFreePort(netId string) int {
|
||||
@@ -317,18 +317,25 @@ func (h *HostDesc) freeStorageSize(onlyLocal, useRsvd bool) int64 {
|
||||
return total
|
||||
}
|
||||
|
||||
func (h *HostDesc) GetFreeStorageSizeOfType(sType string, useRsvd bool) (int64, int64) {
|
||||
return h.freeStorageSizeOfType(sType, useRsvd)
|
||||
func (h *HostDesc) GetFreeStorageSizeOfType(sType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) {
|
||||
return h.freeStorageSizeOfType(sType, mediumType, useRsvd, reqMaxSize)
|
||||
}
|
||||
|
||||
func (h *HostDesc) freeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) {
|
||||
func (h *HostDesc) freeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) {
|
||||
var total int64
|
||||
var actualTotal int64
|
||||
foundLEReqStore := false
|
||||
errs := make([]error, 0)
|
||||
|
||||
for _, storage := range h.Storages {
|
||||
if storage.StorageType == storageType {
|
||||
if IsStorageBackendMediumMatch(storage, storageType, mediumType) {
|
||||
total += int64(storage.FreeCapacity)
|
||||
actualTotal += int64(storage.ActualFreeCapacity)
|
||||
if err := checkStorageSize(storage, reqMaxSize, useRsvd); err != nil {
|
||||
errs = append(errs, err)
|
||||
} else {
|
||||
foundLEReqStore = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if utils.IsLocalStorage(storageType) {
|
||||
@@ -338,11 +345,15 @@ func (h *HostDesc) freeStorageSizeOfType(storageType string, useRsvd bool) (int6
|
||||
total += sizeSub
|
||||
}
|
||||
}
|
||||
if useRsvd {
|
||||
return reservedResourceAddCal(total, h.GuestReservedStorageSizeFree(), useRsvd), actualTotal
|
||||
if !foundLEReqStore {
|
||||
return 0, 0, errors.NewAggregate(errs)
|
||||
}
|
||||
|
||||
return total - int64(h.GetPendingUsage().DiskUsage.Get(storageType)), actualTotal
|
||||
if useRsvd {
|
||||
return reservedResourceAddCal(total, h.GuestReservedStorageSizeFree(), useRsvd), actualTotal, nil
|
||||
}
|
||||
|
||||
return total - int64(h.GetPendingUsage().DiskUsage.Get(storageType)), actualTotal, nil
|
||||
}
|
||||
|
||||
func (h *HostDesc) GetFreePort(netId string) int {
|
||||
|
||||
@@ -100,7 +100,7 @@ type CandidatePropertyGetter interface {
|
||||
FreeMemorySize(useRsvd bool) int64
|
||||
|
||||
StorageInfo() []*baremetal.BaremetalStorage
|
||||
GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64)
|
||||
GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error)
|
||||
|
||||
GetFreePort(netId string) int
|
||||
|
||||
|
||||
@@ -180,17 +180,17 @@ func (mr *MockCandidatePropertyGetterMockRecorder) GetFreePort(arg0 interface{})
|
||||
}
|
||||
|
||||
// GetFreeStorageSizeOfType mocks base method
|
||||
func (m *MockCandidatePropertyGetter) GetFreeStorageSizeOfType(arg0 string, arg1 bool) (int64, int64) {
|
||||
func (m *MockCandidatePropertyGetter) GetFreeStorageSizeOfType(arg0 string, arg1 string, arg2 bool, arg3 int64) (int64, int64, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "GetFreeStorageSizeOfType", arg0, arg1)
|
||||
ret := m.ctrl.Call(m, "GetFreeStorageSizeOfType", arg0, arg1, arg2, arg3)
|
||||
ret0, _ := ret[0].(int64)
|
||||
return ret0, math.MaxInt64
|
||||
return ret0, math.MaxInt64, nil
|
||||
}
|
||||
|
||||
// GetFreeStorageSizeOfType indicates an expected call of GetFreeStorageSizeOfType
|
||||
func (mr *MockCandidatePropertyGetterMockRecorder) GetFreeStorageSizeOfType(arg0, arg1 interface{}) *gomock.Call {
|
||||
func (mr *MockCandidatePropertyGetterMockRecorder) GetFreeStorageSizeOfType(arg0, arg1, arg2, arg3 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetFreeStorageSizeOfType", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).GetFreeStorageSizeOfType), arg0, arg1)
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetFreeStorageSizeOfType", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).GetFreeStorageSizeOfType), arg0, arg1, arg2, arg3)
|
||||
}
|
||||
|
||||
// GetIpmiInfo mocks base method
|
||||
|
||||
@@ -171,7 +171,7 @@ func buildGetter(ctrl *gomock.Controller, param sGetterParams) *mock.MockCandida
|
||||
cg.EXPECT().FreeCPUCount(gomock.Any()).AnyTimes().Return(param.FreeCPUCount)
|
||||
cg.EXPECT().TotalMemorySize(gomock.Any()).AnyTimes().Return(param.TotalMemorySize)
|
||||
cg.EXPECT().FreeMemorySize(gomock.Any()).AnyTimes().Return(param.FreeMemorySize)
|
||||
cg.EXPECT().GetFreeStorageSizeOfType(gomock.Any(), gomock.Any()).AnyTimes().Return(param.FreeStorageSizeAnyType, int64(0))
|
||||
cg.EXPECT().GetFreeStorageSizeOfType(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes().Return(param.FreeStorageSizeAnyType, int64(0), nil)
|
||||
if param.QuotaKeys != nil {
|
||||
cg.EXPECT().GetQuotaKeys(gomock.Any()).AnyTimes().Return(param.QuotaKeys)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user