Skip to content
Closed
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 11 additions & 1 deletion pkg/device/nvidia/device.go
Original file line number Diff line number Diff line change
Expand Up @@ -942,6 +942,16 @@ func cordonedDevices(nodeInfo *device.NodeInfo) map[string]struct{} {
}

func (nv *NvidiaGPUDevices) Fit(devices []*device.DeviceUsage, request device.ContainerDeviceRequest, pod *corev1.Pod, nodeInfo *device.NodeInfo, allocated *device.PodDevices) (bool, map[string]device.ContainerDevices, string) {
return nv.fit(devices, request, pod, nodeInfo, allocated, allocated)
}

// FitWithQuota keeps placement constraints separate from the complete Pod
// allocation history used to calculate effective init and app container usage.
func (nv *NvidiaGPUDevices) FitWithQuota(devices []*device.DeviceUsage, request device.ContainerDeviceRequest, pod *corev1.Pod, nodeInfo *device.NodeInfo, allocated, quotaAllocated *device.PodDevices) (bool, map[string]device.ContainerDevices, string) {
return nv.fit(devices, request, pod, nodeInfo, allocated, quotaAllocated)
}

func (nv *NvidiaGPUDevices) fit(devices []*device.DeviceUsage, request device.ContainerDeviceRequest, pod *corev1.Pod, nodeInfo *device.NodeInfo, allocated, quotaAllocated *device.PodDevices) (bool, map[string]device.ContainerDevices, string) {
k := request
originReq := k.Nums
prevnuma := -1
Expand Down Expand Up @@ -1026,7 +1036,7 @@ func (nv *NvidiaGPUDevices) Fit(devices []*device.DeviceUsage, request device.Co
}
usedmem, usedcores = profile.MemoryMB, profile.Core
}
if !fitQuota(pod, tmpDevs, allocated, pod.Namespace, dev.ID, int64(usedmem), int64(usedcores)) {
if !fitQuota(pod, tmpDevs, quotaAllocated, pod.Namespace, dev.ID, int64(usedmem), int64(usedcores)) {
reason[common.ResourceQuotaNotFit]++
klog.V(3).InfoS(common.ResourceQuotaNotFit, "pod", pod.Name, "memreq", memreq, "coresreq", k.Coresreq)
continue
Expand Down
63 changes: 53 additions & 10 deletions pkg/scheduler/score.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,15 @@ func nodeDeviceBaseTypes(list policy.DeviceUsageList) map[string]struct{} {
return types
}

type quotaAwareDevice interface {
FitWithQuota(devices []*device.DeviceUsage, request device.ContainerDeviceRequest, pod *corev1.Pod, nodeInfo *device.NodeInfo, allocated, quotaAllocated *device.PodDevices) (bool, map[string]device.ContainerDevices, string)
}

func fitInDevices(node *NodeUsage, requests device.ContainerDeviceRequests, pod *corev1.Pod, nodeInfo *device.NodeInfo, devinput *device.PodDevices, weights util.DeviceScoringWeights) (bool, string) {
return fitInDevicesWithQuota(node, requests, pod, nodeInfo, devinput, devinput, weights)
}

func fitInDevicesWithQuota(node *NodeUsage, requests device.ContainerDeviceRequests, pod *corev1.Pod, nodeInfo *device.NodeInfo, devinput, quotaInput *device.PodDevices, weights util.DeviceScoringWeights) (bool, string) {
// Compute scores for all devices based on the request.
for index := range node.Devices.DeviceLists {
node.Devices.DeviceLists[index].ComputeScore(requests, weights)
Expand All @@ -110,7 +118,14 @@ func fitInDevices(node *NodeUsage, requests device.ContainerDeviceRequests, pod
return false, common.GenReason(map[string]int{common.NodeInsufficientDevice: len(typeDevices)}, int(k.Nums))
}

fit, tmpDevs, reason := devPlugin.Fit(typeDevices, k, pod, nodeInfo, devinput)
var fit bool
var tmpDevs map[string]device.ContainerDevices
var reason string
if quotaAware, ok := devPlugin.(quotaAwareDevice); ok {
fit, tmpDevs, reason = quotaAware.FitWithQuota(typeDevices, k, pod, nodeInfo, devinput, quotaInput)
} else {
fit, tmpDevs, reason = devPlugin.Fit(typeDevices, k, pod, nodeInfo, devinput)
}
if !fit {
return false, reason
}
Expand Down Expand Up @@ -264,6 +279,32 @@ func sidecarInitIndexes(task *corev1.Pod) map[int]struct{} {
return idx
}

func activeInitAllocations(initAllocs device.PodDevices, sidecarIdx map[int]struct{}) device.PodDevices {
active := make(device.PodDevices, len(initAllocs))
for devType, rows := range initAllocs {
active[devType] = make(device.PodSingleDevice, len(rows))
for idx := range rows {
if _, isSidecar := sidecarIdx[idx]; isSidecar {
active[devType][idx] = rows[idx].DeepCopy()
}
}
}
return active
}

func completeAllocationHistory(initAllocs, activeAllocs device.PodDevices, numInitContainers int) device.PodDevices {
complete := initAllocs.DeepCopy()
for devType, rows := range activeAllocs {
if _, ok := complete[devType]; !ok {
complete[devType] = make(device.PodSingleDevice, numInitContainers)
}
if len(rows) > numInitContainers {
complete[devType] = append(complete[devType], rows[numInitContainers:].DeepCopy()...)
}
}
return complete
}

func allocateInitContainers(appNodeCopy *NodeUsage, nodeID string, resourceReqs device.PodDeviceRequests, task *corev1.Pod, nodeInfo *device.NodeInfo, allocTypes map[string]struct{}, sidecarIdx map[int]struct{}, numInitContainers int, peakUsage map[string]peakUsageSnapshot, weights util.DeviceScoringWeights) (device.PodDevices, bool, string) {
initAllocs := make(device.PodDevices)

Expand Down Expand Up @@ -304,8 +345,8 @@ func allocateInitContainers(appNodeCopy *NodeUsage, nodeID string, resourceReqs
return initAllocs, true, ""
}

func allocateAppContainers(score *policy.NodeScore, appNodeCopy *NodeUsage, resourceReqs device.PodDeviceRequests, task *corev1.Pod, nodeInfo *device.NodeInfo, allocTypes map[string]struct{}, numInitContainers int, nodeID string, weights util.DeviceScoringWeights) (string, bool) {
appIndex := 0
func allocateAppContainers(score *policy.NodeScore, appNodeCopy *NodeUsage, resourceReqs device.PodDeviceRequests, task *corev1.Pod, nodeInfo *device.NodeInfo, initAllocs device.PodDevices, allocTypes map[string]struct{}, numInitContainers int, nodeID string, weights util.DeviceScoringWeights) (string, bool) {
appIndex := numInitContainers
for ctrid, n := range resourceReqs {
if ctrid < numInitContainers {
continue
Expand All @@ -321,7 +362,11 @@ func allocateAppContainers(score *policy.NodeScore, appNodeCopy *NodeUsage, reso
appIndex++
continue
}
fit, reason := fitInDevices(appNodeCopy, n, task, nodeInfo, &score.Devices, weights)
quotaAllocs := score.Devices
if numInitContainers > 0 {
quotaAllocs = completeAllocationHistory(initAllocs, score.Devices, numInitContainers)
}
fit, reason := fitInDevicesWithQuota(appNodeCopy, n, task, nodeInfo, &score.Devices, &quotaAllocs, weights)
if !fit {
klog.V(4).InfoS(common.NodeUnfitPod, "pod", klog.KObj(task), "node", nodeID, "reason", reason)
return reason, false
Expand Down Expand Up @@ -379,16 +424,14 @@ func (s *Scheduler) scoreNode(nodeID string, node *NodeUsage, resourceReqs devic
return nodeScoreResult{reason: reason}
}
initAllocs = allocs
score.Devices = activeInitAllocations(initAllocs, sidecarIdx)
}

if reason, fit := allocateAppContainers(&score, appNodeCopy, resourceReqs, task, nodeInfo, allocTypes, numInitContainers, nodeID, weights); !fit {
if reason, fit := allocateAppContainers(&score, appNodeCopy, resourceReqs, task, nodeInfo, initAllocs, allocTypes, numInitContainers, nodeID, weights); !fit {
return nodeScoreResult{reason: reason}
}

if numInitContainers > 0 && initAllocs != nil {
for devType, initConList := range initAllocs {
score.Devices[devType] = append(initConList, score.Devices[devType]...)
}
if numInitContainers > 0 {
score.Devices = completeAllocationHistory(initAllocs, score.Devices, numInitContainers)
}

applyPeakUsage(node, appNodeCopy, peakUsage)
Expand Down
108 changes: 108 additions & 0 deletions pkg/scheduler/score_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import (
"github.com/Project-HAMi/HAMi/pkg/device/hygon"
"github.com/Project-HAMi/HAMi/pkg/device/kunlun"
"github.com/Project-HAMi/HAMi/pkg/device/metax"
"github.com/Project-HAMi/HAMi/pkg/device/mthreads"
"github.com/Project-HAMi/HAMi/pkg/device/nvidia"
"github.com/Project-HAMi/HAMi/pkg/scheduler/config"
"github.com/Project-HAMi/HAMi/pkg/scheduler/policy"
Expand All @@ -58,6 +59,113 @@ func TestMain(m *testing.M) {
m.Run()
}

func TestScoreNodeQuotaPreservesInitContainerPositions(t *testing.T) {
const namespace = "score-init-position-quota"
quota := &corev1.ResourceQuota{
ObjectMeta: metav1.ObjectMeta{Name: "gpu-memory", Namespace: namespace},
Spec: corev1.ResourceQuotaSpec{Hard: corev1.ResourceList{
"limits.hami.io/gpumem": resource.MustParse("24000"),
}},
}
manager := device.NewQuotaManager()
manager.AddQuota(quota)
t.Cleanup(func() { manager.DelQuota(quota) })

gpuContainer := func(name string, percentage int64, sidecar bool) corev1.Container {
container := corev1.Container{
Name: name,
Resources: corev1.ResourceRequirements{Limits: corev1.ResourceList{
"hami.io/gpu": resource.MustParse("1"),
"hami.io/gpumem-percentage": *resource.NewQuantity(percentage, resource.DecimalSI),
}},
}
if sidecar {
always := corev1.ContainerRestartPolicyAlways
container.RestartPolicy = &always
}
return container
}

tests := []struct {
name string
inits []corev1.Container
apps []corev1.Container
fits bool
usedMB int32
}{
{name: "no init over quota", apps: []corev1.Container{gpuContainer("app-a", 40, false), gpuContainer("app-b", 40, false)}},
{name: "non-GPU init does not hide an app", inits: []corev1.Container{{Name: "setup"}}, apps: []corev1.Container{gpuContainer("app-a", 40, false), gpuContainer("app-b", 40, false)}},
{name: "GPU init does not hide an app", inits: []corev1.Container{gpuContainer("init", 30, false)}, apps: []corev1.Container{gpuContainer("app-a", 40, false), gpuContainer("app-b", 40, false)}},
{name: "native sidecar and app are concurrent", inits: []corev1.Container{gpuContainer("sidecar", 40, true)}, apps: []corev1.Container{gpuContainer("app", 40, false)}},
{name: "sequential GPU init and app reuse", inits: []corev1.Container{gpuContainer("init", 50, false)}, apps: []corev1.Container{gpuContainer("app", 40, false)}, fits: true, usedMB: 20000},
{name: "non-GPU init leaves valid apps unchanged", inits: []corev1.Container{{Name: "setup"}}, apps: []corev1.Container{gpuContainer("app-a", 20, false), gpuContainer("app-b", 20, false)}, fits: true, usedMB: 16000},
}

for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
pod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{Name: "quota-test", Namespace: namespace},
Spec: corev1.PodSpec{InitContainers: tc.inits, Containers: tc.apps},
}
node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "gpu-node"}}
usage := &NodeUsage{
Node: node,
NodeInfo: &device.NodeInfo{ID: node.Name, Node: node},
Devices: policy.DeviceUsageList{DeviceLists: []*policy.DeviceListsScore{{
Device: &device.DeviceUsage{ID: "gpu-0", Type: nvidia.NvidiaGPUDevice, Health: true, Count: 10, Totalmem: 40000, Totalcore: 100},
}}},
}
nodes := map[string]*NodeUsage{node.Name: usage}
failedNodes := map[string]string{}
got, err := (&Scheduler{}).calcScoreWithOptions(&nodes, device.Resourcereqs(pod), pod, failedNodes, false, false)
assert.NilError(t, err)
if !tc.fits {
assert.Equal(t, len(got.NodeList), 0)
assert.Assert(t, strings.Contains(failedNodes[node.Name], common.ResourceQuotaNotFit), "failure reason: %q", failedNodes[node.Name])
assert.Equal(t, usage.Devices.DeviceLists[0].Device.Usedmem, int32(0))
return
}
assert.Equal(t, len(got.NodeList), 1)
assert.Equal(t, len(got.NodeList[0].Devices[nvidia.NvidiaGPUDevice]), len(tc.inits)+len(tc.apps))
assert.Equal(t, usage.Devices.DeviceLists[0].Device.Usedmem, tc.usedMB)
})
}
}

func TestScoreNodeRegularMthreadsInitDoesNotConstrainAppPlacement(t *testing.T) {
pod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{Name: "mthreads-init", Namespace: "default"},
Spec: corev1.PodSpec{
InitContainers: []corev1.Container{{Name: "init"}},
Containers: []corev1.Container{{Name: "app"}},
},
}
requests := device.PodDeviceRequests{
{mthreads.MthreadsGPUDevice: {Nums: 1, Type: mthreads.MthreadsGPUDevice, Memreq: 8000, Coresreq: 10}},
{mthreads.MthreadsGPUDevice: {Nums: 1, Type: mthreads.MthreadsGPUDevice, Memreq: 2000, Coresreq: 50}},
}
node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "mthreads-node"}}
usage := &NodeUsage{
Node: node,
NodeInfo: &device.NodeInfo{ID: node.Name, Node: node},
Devices: policy.DeviceUsageList{DeviceLists: []*policy.DeviceListsScore{
{Device: &device.DeviceUsage{ID: "gpu-a", Type: mthreads.MthreadsGPUDevice, Health: true, Count: 10, Totalmem: 4000, Totalcore: 100}},
{Device: &device.DeviceUsage{ID: "gpu-b", Type: mthreads.MthreadsGPUDevice, Health: true, Count: 10, Totalmem: 10000, Totalcore: 20}},
}},
}
nodes := map[string]*NodeUsage{node.Name: usage}
failedNodes := map[string]string{}

got, err := (&Scheduler{}).calcScoreWithOptions(&nodes, requests, pod, failedNodes, false, false)
assert.NilError(t, err)
assert.Equal(t, len(failedNodes), 0)
assert.Equal(t, len(got.NodeList), 1)
allocations := got.NodeList[0].Devices[mthreads.MthreadsGPUDevice]
assert.Equal(t, len(allocations), 2)
assert.Equal(t, allocations[0][0].UUID, "gpu-b")
assert.Equal(t, allocations[1][0].UUID, "gpu-a")
}

// test case matrix
/**
| node num | per node device | pod use device | device having use | score |
Expand Down
Loading