mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 02:37:24 +08:00
fix(region,host,scheduler): isolated device balance with to reserve cpus (#21614)
This commit is contained in:
@@ -552,7 +552,8 @@ type SHostPingInput struct {
|
||||
type HostReserveCpusInput struct {
|
||||
Cpus string
|
||||
Mems string
|
||||
DisableSchedLoadBalance *bool `json:"disable_sched_load_balance"`
|
||||
DisableSchedLoadBalance *bool `json:"disable_sched_load_balance"`
|
||||
ProcessesPrefix []string `json:"processes_prefix"`
|
||||
}
|
||||
|
||||
type HostAutoMigrateInput struct {
|
||||
|
||||
@@ -185,4 +185,5 @@ const (
|
||||
|
||||
const (
|
||||
HOSTMETA_RESERVED_CPUS_INFO = "reserved_cpus_info"
|
||||
HOSTMETA_RESERVED_CPUS_RATE = "reserved_cpus_rate"
|
||||
)
|
||||
|
||||
@@ -4750,6 +4750,50 @@ func (hh *SHost) PerformPing(ctx context.Context, userCred mcclient.TokenCredent
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (host *SHost) getHostNodeReservePercent(reservedCpusStr string) (map[string]float32, error) {
|
||||
reservedCpuset, err := cpuset.Parse(reservedCpusStr)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "cpuset parse reserved cpus")
|
||||
}
|
||||
|
||||
topoObj, err := host.SysInfo.Get("topology")
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "get topology from host sys_info")
|
||||
}
|
||||
info := new(hostapi.HostTopology)
|
||||
if err := topoObj.Unmarshal(info); err != nil {
|
||||
return nil, errors.Wrap(err, "Unmarshal host topology struct")
|
||||
}
|
||||
nodecpus := map[int]int{}
|
||||
nodeReservedCpus := map[int]int{}
|
||||
for i := range info.Nodes {
|
||||
cSet := cpuset.NewBuilder()
|
||||
for j := 0; j < len(info.Nodes[i].Cores); j++ {
|
||||
for k := 0; k < len(info.Nodes[i].Cores[j].LogicalProcessors); k++ {
|
||||
if reservedCpuset.Contains(info.Nodes[i].Cores[j].LogicalProcessors[k]) {
|
||||
if cnt, ok := nodeReservedCpus[info.Nodes[i].ID]; !ok {
|
||||
nodeReservedCpus[info.Nodes[i].ID] = 1
|
||||
} else {
|
||||
nodeReservedCpus[info.Nodes[i].ID] = 1 + cnt
|
||||
}
|
||||
}
|
||||
|
||||
cSet.Add(info.Nodes[i].Cores[j].LogicalProcessors[k])
|
||||
}
|
||||
}
|
||||
nodecpus[info.Nodes[i].ID] = cSet.Result().Size()
|
||||
}
|
||||
reserveRate := map[string]float32{}
|
||||
for nodeId, cnt := range nodecpus {
|
||||
reserveCnt, ok := nodeReservedCpus[nodeId]
|
||||
if !ok {
|
||||
reserveCnt = 0
|
||||
}
|
||||
reserveRate[strconv.Itoa(nodeId)] = float32(reserveCnt) / float32(cnt)
|
||||
}
|
||||
return reserveRate, nil
|
||||
}
|
||||
|
||||
func (host *SHost) getHostLogicalCores() ([]int, error) {
|
||||
cpuObj, err := host.SysInfo.Get("cpu_info")
|
||||
if err != nil {
|
||||
@@ -4805,14 +4849,6 @@ func (hh *SHost) PerformReserveCpus(
|
||||
return nil, httperrors.NewNotSupportedError("host type %s not support reserve cpus", hh.HostType)
|
||||
}
|
||||
|
||||
cnt, err := hh.GetRunningGuestCount()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if cnt > 0 {
|
||||
return nil, httperrors.NewBadRequestError("host %s has %d guests, can't update reserve cpus", hh.Id, cnt)
|
||||
}
|
||||
|
||||
if input.Cpus == "" {
|
||||
return nil, httperrors.NewInputParameterError("missing cpus")
|
||||
}
|
||||
@@ -4847,11 +4883,27 @@ func (hh *SHost) PerformReserveCpus(
|
||||
}
|
||||
}
|
||||
|
||||
if len(input.Cpus) > 0 {
|
||||
reservePercent, err := hh.getHostNodeReservePercent(input.Cpus)
|
||||
if err != nil {
|
||||
return nil, errors.Errorf("failed getHostNodeReservePercent: %s", err)
|
||||
}
|
||||
err = hh.SetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, reservePercent, userCred)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
} else {
|
||||
err = hh.RemoveMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, userCred)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
err = hh.SetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_INFO, input, userCred)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if hh.CpuReserved < cs.Size() {
|
||||
if hh.CpuReserved != cs.Size() {
|
||||
_, err = db.Update(hh, func() error {
|
||||
hh.CpuReserved = cs.Size()
|
||||
return nil
|
||||
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
"math"
|
||||
"reflect"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -652,20 +653,18 @@ func (pq *SorttedGroupDevs) Pop() interface{} {
|
||||
return item
|
||||
}
|
||||
|
||||
func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel(
|
||||
ctx context.Context, guest *SGuest, host *SHost, devConfig *api.IsolatedDeviceConfig,
|
||||
userCred mcclient.TokenCredential, usedDevMap map[string]*SIsolatedDevice, preferNumaNodes []int,
|
||||
) error {
|
||||
if len(devConfig.Model) == 0 {
|
||||
return fmt.Errorf("Not found model from info: %#v", devConfig)
|
||||
}
|
||||
// if dev type is not nic, wire is empty string
|
||||
devs, err := manager.findHostUnusedByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId)
|
||||
type SNodeIsolateDevicesInfo struct {
|
||||
TotalDevCount int
|
||||
ReservedRate float32
|
||||
}
|
||||
|
||||
func (manager *SIsolatedDeviceManager) getDevNodesUsedRate(
|
||||
ctx context.Context, host *SHost, devConfig *api.IsolatedDeviceConfig, topo *hostapi.HostTopology,
|
||||
) (map[string]SNodeIsolateDevicesInfo, error) {
|
||||
devs, err := manager.findHostDevsByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId)
|
||||
if err != nil || len(devs) == 0 {
|
||||
return fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id)
|
||||
return nil, fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id)
|
||||
}
|
||||
// 1. group devices by device_path
|
||||
groupDevs := make(SorttedGroupDevs, 0)
|
||||
mapDevs := map[string][]SIsolatedDevice{}
|
||||
for i := range devs {
|
||||
dev := devs[i]
|
||||
@@ -680,11 +679,206 @@ func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel(
|
||||
}
|
||||
mapDevs[devPath] = gdevs
|
||||
}
|
||||
nodesGroupDevs := map[string]SorttedGroupDevs{}
|
||||
for devPath, mappedDevs := range mapDevs {
|
||||
groupDevs = append(groupDevs, &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
numaNode := strconv.Itoa(int(mappedDevs[0].NumaNode))
|
||||
if _, ok := nodesGroupDevs[numaNode]; ok {
|
||||
nodesGroupDevs[numaNode] = append(nodesGroupDevs[numaNode], &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
} else {
|
||||
groupDevs := make(SorttedGroupDevs, 0)
|
||||
nodesGroupDevs[numaNode] = append(groupDevs, &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
reserveRate := map[string]float32{}
|
||||
reserveRateStr := host.GetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, nil)
|
||||
reserveRateJ, err := jsonutils.ParseString(reserveRateStr)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "parse reserveRateStr")
|
||||
}
|
||||
err = reserveRateJ.Unmarshal(&reserveRate)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "unmarshal reserveRateStr")
|
||||
}
|
||||
|
||||
nodeNoDevIds := map[int]int{}
|
||||
for i := range topo.Nodes {
|
||||
nodeId := strconv.Itoa(topo.Nodes[i].ID)
|
||||
if _, ok := nodesGroupDevs[nodeId]; !ok {
|
||||
nodeInt, _ := strconv.Atoi(nodeId)
|
||||
nodeNoDevIds[nodeInt] = -1
|
||||
}
|
||||
}
|
||||
//
|
||||
//for nodeId, _ := range reserveRate {
|
||||
// if _, ok := nodesGroupDevs[nodeId]; !ok {
|
||||
// nodeInt, _ := strconv.Atoi(nodeId)
|
||||
// nodeNoDevIds[nodeInt] = -1
|
||||
// }
|
||||
//}
|
||||
|
||||
reserveNodes := map[string][]string{}
|
||||
for i := range topo.Nodes {
|
||||
if _, ok := nodeNoDevIds[topo.Nodes[i].ID]; ok {
|
||||
minDistance := int(math.MaxInt16)
|
||||
selectNodeId := ""
|
||||
for nodeId, _ := range nodesGroupDevs {
|
||||
nodeInt, _ := strconv.Atoi(nodeId)
|
||||
if topo.Nodes[i].Distances[nodeInt] < minDistance {
|
||||
selectNodeId = strconv.Itoa(nodeInt)
|
||||
minDistance = topo.Nodes[i].Distances[nodeInt]
|
||||
}
|
||||
}
|
||||
noDevNodeId := strconv.Itoa(topo.Nodes[i].ID)
|
||||
log.Debugf("node %s select node %s", noDevNodeId, selectNodeId)
|
||||
if nodes, ok := reserveNodes[selectNodeId]; ok {
|
||||
reserveNodes[selectNodeId] = append(nodes, noDevNodeId)
|
||||
} else {
|
||||
reserveNodes[selectNodeId] = []string{noDevNodeId}
|
||||
}
|
||||
}
|
||||
}
|
||||
reserveRates := map[string]SNodeIsolateDevicesInfo{}
|
||||
for nodeId, devGroups := range nodesGroupDevs {
|
||||
nodeCnt := 1
|
||||
nodeReserveRate := reserveRate[nodeId]
|
||||
if nodes, ok := reserveNodes[nodeId]; ok {
|
||||
for i := range nodes {
|
||||
nodeReserveRate += reserveRate[nodes[i]]
|
||||
nodeCnt += 1
|
||||
}
|
||||
}
|
||||
nodeReserveRate = nodeReserveRate / float32(nodeCnt)
|
||||
devCnt := 0
|
||||
for i := range devGroups {
|
||||
devCnt += len(devGroups[i].Devs)
|
||||
}
|
||||
reserveRates[nodeId] = SNodeIsolateDevicesInfo{
|
||||
TotalDevCount: devCnt,
|
||||
ReservedRate: nodeReserveRate,
|
||||
}
|
||||
log.Debugf("node %v nodeCnt %v nodeReserveRate %v", nodeId, nodeCnt, nodeReserveRate)
|
||||
}
|
||||
return reserveRates, nil
|
||||
}
|
||||
|
||||
func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel(
|
||||
ctx context.Context, guest *SGuest, host *SHost, devConfig *api.IsolatedDeviceConfig,
|
||||
userCred mcclient.TokenCredential, usedDevMap map[string]*SIsolatedDevice, preferNumaNodes []int,
|
||||
) error {
|
||||
if len(devConfig.Model) == 0 {
|
||||
return fmt.Errorf("Not found model from info: %#v", devConfig)
|
||||
}
|
||||
// if dev type is not nic, wire is empty string
|
||||
devs, err := manager.findHostUnusedByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId)
|
||||
if err != nil || len(devs) == 0 {
|
||||
return fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id)
|
||||
}
|
||||
// 1. group devices by device_path and numa nodes
|
||||
//groupDevs := make(SorttedGroupDevs, 0)
|
||||
mapDevs := map[string][]SIsolatedDevice{}
|
||||
for i := range devs {
|
||||
dev := devs[i]
|
||||
devPath := dev.DevicePath
|
||||
var gdevs []SIsolatedDevice
|
||||
|
||||
gdevs, ok := mapDevs[devPath]
|
||||
if !ok {
|
||||
gdevs = []SIsolatedDevice{dev}
|
||||
} else {
|
||||
gdevs = append(gdevs, dev)
|
||||
}
|
||||
mapDevs[devPath] = gdevs
|
||||
}
|
||||
|
||||
var groupDevs SorttedGroupDevs
|
||||
if len(preferNumaNodes) > 0 {
|
||||
groupDevs = make(SorttedGroupDevs, 0)
|
||||
for devPath, mappedDevs := range mapDevs {
|
||||
groupDevs = append(groupDevs, &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
}
|
||||
} else {
|
||||
nodesGroupDevs := map[int8]SorttedGroupDevs{}
|
||||
for devPath, mappedDevs := range mapDevs {
|
||||
numaNode := mappedDevs[0].NumaNode
|
||||
if _, ok := nodesGroupDevs[numaNode]; ok {
|
||||
nodesGroupDevs[numaNode] = append(nodesGroupDevs[numaNode], &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
} else {
|
||||
groupDevs := make(SorttedGroupDevs, 0)
|
||||
nodesGroupDevs[numaNode] = append(groupDevs, &GroupDevs{
|
||||
DevPath: devPath,
|
||||
Devs: mappedDevs,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
var selectedNode int8 = -1
|
||||
if len(nodesGroupDevs) == 1 {
|
||||
for nodeId := range nodesGroupDevs {
|
||||
selectedNode = nodeId
|
||||
}
|
||||
} else {
|
||||
reservedCpusStr := host.GetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_INFO, nil)
|
||||
if len(reservedCpusStr) > 0 {
|
||||
topoObj, err := host.SysInfo.Get("topology")
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "get topology from host sys_info")
|
||||
}
|
||||
topo := new(hostapi.HostTopology)
|
||||
if err := topoObj.Unmarshal(topo); err != nil {
|
||||
return errors.Wrap(err, "Unmarshal host topology struct")
|
||||
}
|
||||
nodesReserveRate, err := manager.getDevNodesUsedRate(ctx, host, devConfig, topo)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var selectedNodeUtil float32 = 1.0
|
||||
for nodeId, gds := range nodesGroupDevs {
|
||||
freeDevCnt := 0
|
||||
for i := range gds {
|
||||
freeDevCnt += len(gds[i].Devs)
|
||||
}
|
||||
|
||||
nodeTotalCnt := nodesReserveRate[strconv.Itoa(int(nodeId))].TotalDevCount
|
||||
usedDevCnt := nodeTotalCnt - freeDevCnt
|
||||
|
||||
nodeReserveRate := nodesReserveRate[strconv.Itoa(int(nodeId))].ReservedRate
|
||||
nodeCnt := (1 - nodeReserveRate) * float32(nodeTotalCnt)
|
||||
nodeutil := float32(usedDevCnt) / nodeCnt
|
||||
log.Debugf("selectedNodeUtil node %v util %v usedDevCnt %v totalDevCnt %v", nodeId, nodeutil, usedDevCnt, nodeCnt)
|
||||
if nodeutil < selectedNodeUtil {
|
||||
selectedNodeUtil = nodeutil
|
||||
selectedNode = nodeId
|
||||
}
|
||||
}
|
||||
} else {
|
||||
var selectedNodeDevCnt = 0
|
||||
for nodeId, gds := range nodesGroupDevs {
|
||||
devCnt := 0
|
||||
for i := range gds {
|
||||
devCnt += len(gds[i].Devs)
|
||||
}
|
||||
if devCnt > selectedNodeDevCnt {
|
||||
selectedNodeDevCnt = devCnt
|
||||
selectedNode = nodeId
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
log.Debugf("selectedNodeUtil node %v", selectedNode)
|
||||
groupDevs = nodesGroupDevs[selectedNode]
|
||||
}
|
||||
sort.Sort(groupDevs)
|
||||
|
||||
@@ -845,6 +1039,29 @@ func (manager *SIsolatedDeviceManager) findHostUnusedByDevConfig(model, devType,
|
||||
return manager.findHostUnusedByDevAttr(model, "dev_type", devType, hostId, wireId)
|
||||
}
|
||||
|
||||
func (manager *SIsolatedDeviceManager) findHostDevsByDevConfig(model, devType, hostId, wireId string) ([]SIsolatedDevice, error) {
|
||||
return manager.findHostDevsByDevAttr(model, "dev_type", devType, hostId, wireId)
|
||||
}
|
||||
func (manager *SIsolatedDeviceManager) findHostDevsByDevAttr(model, attrKey, attrVal, hostId, wireId string) ([]SIsolatedDevice, error) {
|
||||
devs := make([]SIsolatedDevice, 0)
|
||||
q := manager.Query()
|
||||
q = q.Equals("model", model).Equals("host_id", hostId)
|
||||
if attrVal != "" {
|
||||
q.Equals(attrKey, attrVal)
|
||||
}
|
||||
if wireId != "" {
|
||||
wire := WireManager.FetchWireById(wireId)
|
||||
if wire.VpcId == api.DEFAULT_VPC_ID {
|
||||
q = q.Equals("wire_id", wireId)
|
||||
}
|
||||
}
|
||||
err := db.FetchModelObjects(manager, q, &devs)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return devs, nil
|
||||
}
|
||||
|
||||
func (manager *SIsolatedDeviceManager) findHostUnusedByDevAttr(model, attrKey, attrVal, hostId, wireId string) ([]SIsolatedDevice, error) {
|
||||
devs := make([]SIsolatedDevice, 0)
|
||||
q := manager.findUnusedQuery()
|
||||
|
||||
@@ -219,7 +219,10 @@ type CpuSetCounter struct {
|
||||
Nodes []*NumaNode
|
||||
NumaEnabled bool
|
||||
CPUCmtbound float32
|
||||
Lock sync.Mutex
|
||||
MEMCmtbound float32
|
||||
|
||||
GuestIds map[string]struct{}
|
||||
Lock sync.Mutex
|
||||
}
|
||||
|
||||
func NewGuestCpuSetCounter(
|
||||
@@ -230,6 +233,8 @@ func NewGuestCpuSetCounter(
|
||||
cpuSetCounter.Nodes = make([]*NumaNode, len(info.Nodes))
|
||||
cpuSetCounter.NumaEnabled = numaAllocate
|
||||
cpuSetCounter.CPUCmtbound = cpuCmtbound
|
||||
cpuSetCounter.MEMCmtbound = memCmtBound
|
||||
cpuSetCounter.GuestIds = map[string]struct{}{}
|
||||
hasL3Cache := false
|
||||
nodeReserveMem := reservedMemMb / len(info.Nodes) * 1024
|
||||
for i := 0; i < len(info.Nodes); i++ {
|
||||
@@ -244,6 +249,7 @@ func NewGuestCpuSetCounter(
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
reservedCpuCnt := 0
|
||||
cpuDies := make([]*CPUDie, 0)
|
||||
for j := 0; j < len(info.Nodes[i].Caches); j++ {
|
||||
if info.Nodes[i].Caches[j].Level != 3 {
|
||||
@@ -254,6 +260,7 @@ func NewGuestCpuSetCounter(
|
||||
dieBuilder := cpuset.NewBuilder()
|
||||
for k := 0; k < len(info.Nodes[i].Caches[j].LogicalProcessors); k++ {
|
||||
if reservedCpus != nil && reservedCpus.Contains(int(info.Nodes[i].Caches[j].LogicalProcessors[k])) {
|
||||
reservedCpuCnt += 1
|
||||
continue
|
||||
}
|
||||
dieBuilder.Add(int(info.Nodes[i].Caches[j].LogicalProcessors[k]))
|
||||
@@ -271,6 +278,7 @@ func NewGuestCpuSetCounter(
|
||||
for j := 0; j < len(info.Nodes[i].Cores); j++ {
|
||||
for k := 0; k < len(info.Nodes[i].Cores[j].LogicalProcessors); k++ {
|
||||
if reservedCpus != nil && reservedCpus.Contains(info.Nodes[i].Cores[j].LogicalProcessors[k]) {
|
||||
reservedCpuCnt += 1
|
||||
continue
|
||||
}
|
||||
dieBuilder.Add(info.Nodes[i].Cores[j].LogicalProcessors[k])
|
||||
@@ -286,6 +294,7 @@ func NewGuestCpuSetCounter(
|
||||
|
||||
hasL3Cache = false
|
||||
node.CpuDies = cpuDies
|
||||
node.ReserveCpuCount = reservedCpuCnt
|
||||
sort.Sort(node.CpuDies)
|
||||
cpuSetCounter.Nodes[i] = node
|
||||
}
|
||||
@@ -294,13 +303,14 @@ func NewGuestCpuSetCounter(
|
||||
return cpuSetCounter, nil
|
||||
}
|
||||
|
||||
func (pq *CpuSetCounter) AllocCpusetWithNodeCount(vcpuCount int, memSizeKB int64, nodeCount int) (map[int]SAllocNumaCpus, error) {
|
||||
func (pq *CpuSetCounter) AllocCpusetWithNodeCount(vcpuCount int, memSizeKB int64, nodeCount int, guestId string) (map[int]SAllocNumaCpus, error) {
|
||||
if !pq.NumaEnabled {
|
||||
return pq.AllocCpuset(vcpuCount, memSizeKB, nil)
|
||||
return pq.AllocCpuset(vcpuCount, memSizeKB, nil, guestId)
|
||||
}
|
||||
if len(pq.Nodes) < nodeCount {
|
||||
return nil, nil
|
||||
}
|
||||
pq.GuestIds[guestId] = struct{}{}
|
||||
|
||||
pq.Lock.Lock()
|
||||
defer pq.Lock.Unlock()
|
||||
@@ -339,13 +349,14 @@ func (pq *CpuSetCounter) IsNumaEnabled() bool {
|
||||
return pq.NumaEnabled
|
||||
}
|
||||
|
||||
func (pq *CpuSetCounter) AllocCpuset(vcpuCount int, memSizeKB int64, preferNumaNodes []int8) (map[int]SAllocNumaCpus, error) {
|
||||
func (pq *CpuSetCounter) AllocCpuset(vcpuCount int, memSizeKB int64, preferNumaNodes []int8, guestId string) (map[int]SAllocNumaCpus, error) {
|
||||
pq.Lock.Lock()
|
||||
defer pq.Lock.Unlock()
|
||||
|
||||
if len(pq.Nodes) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
pq.GuestIds[guestId] = struct{}{}
|
||||
|
||||
if pq.NumaEnabled && len(preferNumaNodes) > 0 {
|
||||
sortedNumaDistance := pq.getDistancesSeqByPreferNodes(preferNumaNodes)
|
||||
@@ -416,7 +427,7 @@ func (pq *CpuSetCounter) allocCpuNumaNodesByPreferNodes(
|
||||
} else {
|
||||
log.Infof("node %v not enough", pq.Nodes[i])
|
||||
}
|
||||
log.Infof("node %d, free mems %d", pq.Nodes[nodeIdx].NodeId, pq.Nodes[nodeIdx].NumaNodeFreeMemSizeKB)
|
||||
log.Infof("node %d, free mems %d, vcpuCount %d, GuestCounts %v", pq.Nodes[nodeIdx].NodeId, pq.Nodes[nodeIdx].NumaNodeFreeMemSizeKB, pq.Nodes[nodeIdx].VcpuCount, len(pq.GuestIds))
|
||||
}
|
||||
|
||||
if allocatedNode < nodeCount {
|
||||
@@ -429,6 +440,8 @@ type SSortedNumaDistance struct {
|
||||
NodeIndex int
|
||||
Distance int
|
||||
FreeMemSize int
|
||||
UsedRate float32
|
||||
CpuReserved bool
|
||||
}
|
||||
|
||||
func (pq *CpuSetCounter) getDistancesSeqByPreferNodes(preferNumaNodes []int8) []SSortedNumaDistance {
|
||||
@@ -438,19 +451,37 @@ func (pq *CpuSetCounter) getDistancesSeqByPreferNodes(preferNumaNodes []int8) []
|
||||
for j := range preferNumaNodes {
|
||||
distance += pq.Nodes[i].Distances[preferNumaNodes[j]]
|
||||
}
|
||||
|
||||
var useableCpuRate float32 = 1.0
|
||||
if pq.Nodes[i].ReserveCpuCount > 0 {
|
||||
useableCpuRate = float32(pq.Nodes[i].CpuCount) / float32(pq.Nodes[i].CpuCount+pq.Nodes[i].ReserveCpuCount)
|
||||
}
|
||||
|
||||
usedMems := float32(pq.Nodes[i].NumaNodeMemSizeKB - pq.Nodes[i].NumaNodeFreeMemSizeKB)
|
||||
usedRate := usedMems / (float32(pq.Nodes[i].MemTotalSizeKB) * pq.MEMCmtbound * useableCpuRate)
|
||||
|
||||
//memCmt := float32(usedMems / pq.Nodes[i].NumaNodeMemSizeKB)
|
||||
//cpuPro := float32(pq.Nodes[i].CpuCount) * pq.CPUCmtbound / (float32(pq.Nodes[i].CpuCount)*pq.CPUCmtbound - float32(pq.Nodes[i].VcpuCount))
|
||||
sortedNumaDistance[i] = SSortedNumaDistance{
|
||||
NodeIndex: i,
|
||||
Distance: distance,
|
||||
FreeMemSize: int(pq.Nodes[i].NumaNodeFreeMemSizeKB),
|
||||
UsedRate: usedRate,
|
||||
CpuReserved: pq.Nodes[i].ReserveCpuCount > 0,
|
||||
}
|
||||
}
|
||||
sort.Slice(sortedNumaDistance, func(i, j int) bool {
|
||||
// 7 is tolerant max distances
|
||||
if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance {
|
||||
if sortedNumaDistance[i].Distance > (7 + sortedNumaDistance[j].Distance) {
|
||||
return false
|
||||
} else if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance {
|
||||
return true
|
||||
} else {
|
||||
return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize
|
||||
}
|
||||
|
||||
if sortedNumaDistance[i].CpuReserved {
|
||||
return sortedNumaDistance[i].UsedRate < sortedNumaDistance[j].UsedRate
|
||||
}
|
||||
return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize
|
||||
})
|
||||
return sortedNumaDistance
|
||||
}
|
||||
@@ -670,10 +701,12 @@ type NumaNode struct {
|
||||
LogicalProcessors cpuset.CPUSet
|
||||
VcpuCount int
|
||||
CpuCount int
|
||||
ReserveCpuCount int
|
||||
|
||||
NodeId int
|
||||
Distances []int
|
||||
NumaNodeMemSizeKB int64
|
||||
MemTotalSizeKB int64
|
||||
NumaNodeFreeMemSizeKB int64
|
||||
}
|
||||
|
||||
@@ -696,6 +729,7 @@ func NewNumaNode(
|
||||
return nil, errors.Errorf("node %d no memory info: %#v", nodeInfo.ID, nodeInfo)
|
||||
}
|
||||
n.NumaNodeMemSizeKB = int64(float32(nodeInfo.Memory.TotalUsableBytes/1024-int64(reservedMemSizeKB)) * memCmtBound)
|
||||
n.MemTotalSizeKB = nodeInfo.Memory.TotalUsableBytes / 1024
|
||||
} else {
|
||||
nodeHugepagePath := fmt.Sprintf("/sys/devices/system/node/node%d/hugepages/hugepages-%dkB", n.NodeId, hugepageSizeKB)
|
||||
if !fileutils2.Exists(nodeHugepagePath) {
|
||||
|
||||
@@ -2488,7 +2488,7 @@ func (task *SGuestHotplugCpuMemTask) startAddCpusWithFreeVcpuSet(vcpuSet []int)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, nil)
|
||||
cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, nil, task.GetId())
|
||||
for _, cpus := range cpus {
|
||||
//pcpus := cpuset.NewCPUSet(cpus.Cpuset...).String()
|
||||
//vcpus := fmt.Sprintf("%d-%d", vcpuId, vcpuId)
|
||||
|
||||
@@ -1117,7 +1117,7 @@ func (s *sPodGuestInstance) allocateCpuNumaPin() error {
|
||||
}
|
||||
}
|
||||
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes)
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes, s.GetId())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -1140,6 +1140,12 @@ func (s *sPodGuestInstance) allocateCpuNumaPin() error {
|
||||
vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset))
|
||||
for i := range numaCpus.Cpuset {
|
||||
vcpuPin[i].Pcpu = numaCpus.Cpuset[i]
|
||||
if i < int(s.Desc.Cpu) {
|
||||
vcpuPin[i].Vcpu = i
|
||||
} else {
|
||||
vcpuPin[i].Vcpu = -1
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
memPin := &desc.SCpuNumaPin{
|
||||
|
||||
@@ -177,7 +177,7 @@ func (s *SKVMGuestInstance) reallocateNumaNodes(isMigrate bool) error {
|
||||
}
|
||||
|
||||
func (s *SKVMGuestInstance) reallocateMigrateNumaNodes() error {
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(s.Desc.Cpu), s.Desc.Mem*1024, len(s.Desc.MemDesc.Mem.Mems)+1)
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(s.Desc.Cpu), s.Desc.Mem*1024, len(s.Desc.MemDesc.Mem.Mems)+1, s.GetId())
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "AllocCpusetWithNodeCount")
|
||||
}
|
||||
@@ -371,7 +371,7 @@ func (s *SKVMGuestInstance) initLiveDescFromSourceGuest(srcDesc *desc.SGuestDesc
|
||||
cpuNumaPin = s.Desc.CpuNumaPin
|
||||
} else {
|
||||
// allocate cpu numa pin local
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1)
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1, s.GetId())
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "AllocCpusetWithNodeCount")
|
||||
}
|
||||
@@ -2700,7 +2700,7 @@ func (s *SKVMGuestInstance) allocGuestNumaCpuset() error {
|
||||
}
|
||||
}
|
||||
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes)
|
||||
nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes, s.GetId())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -173,6 +173,8 @@ func LoadGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) error {
|
||||
if s.IsRunning() {
|
||||
m.cpuSet.Lock.Lock()
|
||||
defer m.cpuSet.Lock.Unlock()
|
||||
m.cpuSet.GuestIds[s.GetId()] = struct{}{}
|
||||
|
||||
for _, vcpuPin := range guestDesc.VcpuPin {
|
||||
pcpuSet, err := cpuset.Parse(vcpuPin.Pcpus)
|
||||
if err != nil {
|
||||
@@ -186,13 +188,18 @@ func LoadGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) error {
|
||||
}
|
||||
m.cpuSet.LoadCpus(pcpuSet.ToSlice(), vcpuSet.Size())
|
||||
}
|
||||
|
||||
for _, numaCpuset := range guestDesc.CpuNumaPin {
|
||||
pcpus := make([]int, 0)
|
||||
vcpus := make([]int, 0)
|
||||
for i := range numaCpuset.VcpuPin {
|
||||
pcpus = append(pcpus, numaCpuset.VcpuPin[i].Pcpu)
|
||||
if numaCpuset.VcpuPin[i].Vcpu >= 0 {
|
||||
vcpus = append(vcpus, numaCpuset.VcpuPin[i].Vcpu)
|
||||
}
|
||||
}
|
||||
|
||||
m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, int(*numaCpuset.NodeId), pcpus, len(numaCpuset.VcpuPin))
|
||||
m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, int(*numaCpuset.NodeId), pcpus, len(vcpus))
|
||||
}
|
||||
}
|
||||
return nil
|
||||
@@ -205,16 +212,22 @@ func ReleaseCpuNumaPin(m *SGuestManager, cpuNumaPin []*desc.SCpuNumaPin) {
|
||||
|
||||
for _, numaCpus := range cpuNumaPin {
|
||||
pcpus := make([]int, 0)
|
||||
vcpus := make([]int, 0)
|
||||
for i := range numaCpus.VcpuPin {
|
||||
pcpus = append(pcpus, numaCpus.VcpuPin[i].Pcpu)
|
||||
if numaCpus.VcpuPin[i].Vcpu >= 0 {
|
||||
vcpus = append(vcpus, numaCpus.VcpuPin[i].Vcpu)
|
||||
}
|
||||
}
|
||||
m.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, int(*numaCpus.NodeId), pcpus, len(numaCpus.VcpuPin))
|
||||
m.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, int(*numaCpus.NodeId), pcpus, len(vcpus))
|
||||
}
|
||||
}
|
||||
|
||||
func ReleaseGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) {
|
||||
m.cpuSet.Lock.Lock()
|
||||
defer m.cpuSet.Lock.Unlock()
|
||||
delete(m.cpuSet.GuestIds, s.GetId())
|
||||
|
||||
guestDesc := s.GetDesc()
|
||||
ReleaseCpuNumaPin(m, guestDesc.CpuNumaPin)
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"math"
|
||||
"net"
|
||||
"os"
|
||||
@@ -833,6 +834,9 @@ func (h *SHostInfo) initCgroup() error {
|
||||
!reservedCpusTask.CustomConfig(cgrouputils.CPUSET_SCHED_LOAD_BALANCE, "0") {
|
||||
return fmt.Errorf("failed init host reserved cpuset sched load balance")
|
||||
}
|
||||
if len(h.reservedCpusInfo.ProcessesPrefix) > 0 {
|
||||
go h.startBindReservedCpus(h.reservedCpusInfo.ProcessesPrefix)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -2670,6 +2674,58 @@ func (h *SHostInfo) MemCmtBound() float32 {
|
||||
return h.memCmtBound
|
||||
}
|
||||
|
||||
func (h *SHostInfo) getProcessesPids(processesPrefix []string) (map[string]string, error) {
|
||||
files, err := ioutil.ReadDir("/proc")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
res := map[string]string{}
|
||||
re := regexp.MustCompile(`^\d+$`)
|
||||
for _, f := range files {
|
||||
if re.MatchString(f.Name()) {
|
||||
cmdline, err := fileutils2.FileGetContents(path.Join("/proc", f.Name(), "cmdline"))
|
||||
if err != nil {
|
||||
log.Errorf("failed read proc %s cmdline: %s", f.Name(), err)
|
||||
continue
|
||||
}
|
||||
segs := strings.Split(cmdline, "\x00")
|
||||
if utils.IsInStringArray(segs[0], processesPrefix) {
|
||||
res[segs[0]] = f.Name()
|
||||
log.Infof("getProcessesPids append %s %s", segs[0], f.Name())
|
||||
}
|
||||
}
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
func (h *SHostInfo) startBindReservedCpus(processesPrefix []string) {
|
||||
for {
|
||||
processPids, err := h.getProcessesPids(processesPrefix)
|
||||
if err != nil {
|
||||
log.Errorf("getProcessesPids %s", err)
|
||||
} else {
|
||||
for process, pid := range processPids {
|
||||
cgroupName := path.Join(hostconsts.HOST_RESERVED_CPUSET, strings.ReplaceAll(process, "/", "_"))
|
||||
task := cgrouputils.NewCGroupCPUSetTask(pid, cgroupName, 0, "")
|
||||
if !task.Configure() {
|
||||
log.Errorf("process failed init reserved cpuset %s %s", process, pid)
|
||||
continue
|
||||
}
|
||||
|
||||
if !task.CustomConfig(cgrouputils.CPUSET_CLONE_CHILDREN, "1") {
|
||||
log.Errorf("process failed set host reserved cpuset clone children %s %s", process, pid)
|
||||
continue
|
||||
}
|
||||
if !task.SetTask() {
|
||||
log.Errorf("process %s %s failed set cgroup cpuset", process, pid)
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Second * 100)
|
||||
}
|
||||
}
|
||||
|
||||
func NewHostInfo() (*SHostInfo, error) {
|
||||
var res = new(SHostInfo)
|
||||
res.sysinfo = &SSysInfo{}
|
||||
|
||||
@@ -92,10 +92,11 @@ type HostReserveCpusOptions struct {
|
||||
Cpus string
|
||||
Mems string
|
||||
DisableSchedLoadBalance bool
|
||||
ProcessesPrefix []string `help:"Processes prefix bind reserved cpus"`
|
||||
}
|
||||
|
||||
func (o *HostReserveCpusOptions) Params() (jsonutils.JSONObject, error) {
|
||||
return options.StructToParams(o)
|
||||
return jsonutils.Marshal(o), nil
|
||||
}
|
||||
|
||||
type HostAutoMigrateOnHostDownOptions struct {
|
||||
|
||||
+5
-3
@@ -612,11 +612,13 @@ func (h *SHostTopo) getDistancesSeqByPreferNodes(preferNumaNodes []int) []SSorte
|
||||
}
|
||||
sort.Slice(sortedNumaDistance, func(i, j int) bool {
|
||||
// 7 is tolerant max distances
|
||||
if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance {
|
||||
if sortedNumaDistance[i].Distance > (7 + sortedNumaDistance[j].Distance) {
|
||||
return false
|
||||
} else if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance {
|
||||
return true
|
||||
} else {
|
||||
return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize
|
||||
}
|
||||
|
||||
return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize
|
||||
})
|
||||
return sortedNumaDistance
|
||||
}
|
||||
|
||||
@@ -625,6 +625,7 @@ const (
|
||||
CPUSET_CPUS = "cpuset.cpus"
|
||||
CPUSET_MEMS = "cpuset.mems"
|
||||
CPUSET_SCHED_LOAD_BALANCE = "cpuset.sched_load_balance"
|
||||
CPUSET_CLONE_CHILDREN = "cgroup.clone_children"
|
||||
)
|
||||
|
||||
func (c *CGroupCPUSetTask) Module() string {
|
||||
|
||||
Reference in New Issue
Block a user