vcsdplugin
# 1.plugin
# 1.1.注册
volcano启动会触发plugin调用RegisterPluginBuilder注册,相关的plugin会基于builder初始化及注册到session plugins。func init() { // Plugins for Jobs framework.RegisterPluginBuilder(drf.PluginName, drf.New) framework.RegisterPluginBuilder(gang.PluginName, gang.New) framework.RegisterPluginBuilder(deviceshare.PluginName, deviceshare.New) framework.RegisterPluginBuilder(predicates.PluginName, predicates.New) framework.RegisterPluginBuilder(priority.PluginName, priority.New) framework.RegisterPluginBuilder(nodeorder.PluginName, nodeorder.New) framework.RegisterPluginBuilder(conformance.PluginName, conformance.New) framework.RegisterPluginBuilder(binpack.PluginName, binpack.New) framework.RegisterPluginBuilder(tdm.PluginName, tdm.New) framework.RegisterPluginBuilder(overcommit.PluginName, overcommit.New) framework.RegisterPluginBuilder(sla.PluginName, sla.New) framework.RegisterPluginBuilder(tasktopology.PluginName, tasktopology.New) framework.RegisterPluginBuilder(numaaware.PluginName, numaaware.New) framework.RegisterPluginBuilder(cdp.PluginName, cdp.New) framework.RegisterPluginBuilder(rescheduling.PluginName, rescheduling.New) framework.RegisterPluginBuilder(usage.PluginName, usage.New) framework.RegisterPluginBuilder(pdb.PluginName, pdb.New) framework.RegisterPluginBuilder(nodegroup.PluginName, nodegroup.New) // Plugins for Queues framework.RegisterPluginBuilder(proportion.PluginName, proportion.New) framework.RegisterPluginBuilder(capacity.PluginName, capacity.New) // Plugins for Extender framework.RegisterPluginBuilder(extender.PluginName, extender.New) // Plugins for ResourceQuota framework.RegisterPluginBuilder(resourcequota.PluginName, resourcequota.New) } // RegisterPluginBuilder register the plugin func RegisterPluginBuilder(name string, pc PluginBuilder) { pluginMutex.Lock() defer pluginMutex.Unlock() pluginBuilders[name] = pc } // GetPluginBuilder get the pluginbuilder by name func GetPluginBuilder(name string) (PluginBuilder, bool) { pluginMutex.RLock() defer pluginMutex.RUnlock() pb, found := pluginBuilders[name] return pb, found }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48注意
这里注册的都是内部插件,
scheduler启动会基于.so动态链接加载外部实现插件注册到pluginBuilders
# 1.2.加载
volcano scheduler每轮调度循环都会创建framework session会话,内部会基于pluginBuilders实例化一遍所有的plugin。// OpenSession start the session func OpenSession(cache cache.Cache, tiers []conf.Tier, configurations []conf.Configuration) *Session { ssn := openSession(cache) ssn.Tiers = tiers ssn.Configurations = configurations ssn.NodeMap = GenerateNodeMapAndSlice(ssn.Nodes) ssn.PodLister = NewPodLister(ssn) for _, tier := range tiers { for _, plugin := range tier.Plugins { if pb, found := GetPluginBuilder(plugin.Name); found { plugin := pb(plugin.Arguments) ssn.plugins[plugin.Name()] = plugin ... // plugin向session注册回调 plugin.OnSessionOpen(ssn) ... } } } return ssn }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23注意
pluginBuilder就是内部plugin和加载的.so外部plugin对应的初始化回调,下面会介绍一些常用插件
# 2.gang
# 2.1.open
session初始化会依次执行plugin.OnSessionOpen,内部会向session注册各种调用链,scheduler action会调用这些回调进行相关检查。func (gp *gangPlugin) OnSessionOpen(ssn *framework.Session) { // job有效性检查 validJobFn := func(obj interface{}) *api.ValidateResult { job, ok := obj.(*api.JobInfo) if !ok { return err } // (bound+binding+running+allocated+succeeded+pipelined+pending)[taskRole]≥task minAvailable if valid := job.CheckTaskValid(); !valid { return err } // bound+binding+running+allocated+succeeded+pipelined+pending vtn := job.ValidTaskNum() if vtn < job.MinAvailable { return err } return nil } // 向session注册validJobFn ssn.AddJobValidFn(gp.Name(), validJobFn) preemptableFn := func(preemptor *api.TaskInfo, preemptees []*api.TaskInfo) ([]*api.TaskInfo, int) { ... // 遍历待抢占task for _, preemptee := range preemptees { job := ssn.Jobs[preemptee.Job] // job实际运行的task数量 if _, found := jobOccupiedMap[job.UID]; !found { // bound+binding+running+allocated+succeeded jobOccupiedMap[job.UID] = job.ReadyTaskNum() } // 运行的task数量满足最小可用,多出的task可抢占 if jobOccupiedMap[job.UID] > job.MinAvailable { jobOccupiedMap[job.UID]-- victims = append(victims, preemptee) } } return victims, util.Permit } // 向session注册preempt/reclaim回调 ssn.AddReclaimableFn(gp.Name(), preemptableFn) ssn.AddPreemptableFn(gp.Name(), preemptableFn) ... // 向session注册JobOrderFn,未满足最小可用任务的优先 ssn.AddJobOrderFn(gp.Name(), jobOrderFn) // 注册JobReadyFn ssn.AddJobReadyFn(gp.Name(), func(obj interface{}) bool { ji := obj.(*api.JobInfo) // (bound+binding+running+allocated+succeeded+bestEffort pending)[taskRole]≥task minAvailable // (bound+binding+running+allocated+succeeded+bestEffort pending)≥job minAvailable if ji.CheckTaskReady() && ji.IsReady() { return true } return false }) // 注册pipelinedFn pipelinedFn := func(obj interface{}) int { ji := obj.(*api.JobInfo) // (bound+binding+running+allocated+succeeded+pipelined+bestEffort pending)[taskRole]≥task minAvailable // (bound+binding+running+allocated+succeeded+pipelined+bestEffort pending)≥job minAvailable if ji.CheckTaskPipelined() && ji.IsPipelined() { return util.Permit } return util.Reject } ssn.AddJobPipelinedFn(gp.Name(), pipelinedFn) // 注册饥饿检查Fn jobStarvingFn := func(obj interface{}) bool { ji := obj.(*api.JobInfo) // (bound+binding+running+allocated+succeeded+pipelined)<job minAvailable return ji.IsStarving() } ssn.AddJobStarvingFns(gp.Name(), jobStarvingFn) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86注意
gang注册的callback基本上都是根据task状态检查是否满足最小副本任务要求
# 2.2.close
gang.OnSessionClose()用于会话结束对job进行最终状态评估,基于job minAvailable更新job podGroup condition状态。func (gp *gangPlugin) OnSessionClose(ssn *framework.Session) { ... // 遍历session job for _, job := range ssn.Jobs { // (bound+binding+running+allocated+succeeded+bestEffort pending)<job minAvailable if !job.IsReady() { schedulableTaskNum := func() (num int32) { for _, task := range job.TaskStatusIndex[api.Pending] { ctx := task.GetTransactionContext() // task参与过事务 if task.LastTransaction != nil { ctx = *task.LastTransaction } // pending task分配资源还未运行,视为准ready if api.AllocatedStatus(ctx.Status) { num++ } } // bound+binding+running+allocated+succeeded+已分配资源的pending task return num + job.ReadyTaskNum() } // 未就绪的task数量 unreadyTaskCount = job.MinAvailable - schedulableTaskNum() ... // 未调度的job unScheduleJobCount++ ... // 更新job podGroup condition jc := &scheduling.PodGroupCondition{ Type: scheduling.PodGroupUnschedulableType, Status: v1.ConditionTrue, LastTransitionTime: metav1.Now(), TransitionID: string(ssn.UID), Reason: v1beta1.NotEnoughResourcesReason, Message: msg, } ssn.UpdatePodGroupCondition(job, jc) ... // job就绪 } else { // 更新job podGroup condition jc := &scheduling.PodGroupCondition{ Type: scheduling.PodGroupScheduled, Status: v1.ConditionTrue, LastTransitionTime: metav1.Now(), TransitionID: string(ssn.UID), Reason: "tasks in gang are ready to be scheduled", Message: "", } ssn.UpdatePodGroupCondition(job, jc) ... } ... unreadyTaskCount = 0 } ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61注意
job状态评估还是基于minAvailable可用检查,基于评估结果更新job podGroup condition
# 3.drf
# 3.1.open
drf基于资源公平调度算法和分层公平调度算法调整集群业务的吞吐量,尽量规避胖业务饿死众多小业务场景,确保多种资源共存时满足分配的公平原则。func (drf *drfPlugin) OnSessionOpen(ssn *framework.Session) { // 记录资源总量 drf.totalResource.Add(ssn.TotalResource) ... // 遍历session job for _, job := range ssn.Jobs { ... for status, tasks := range job.TaskStatusIndex { // bound/binding/running/allocated视为资源有效分配 if api.AllocatedStatus(status) { // 记录job分配的资源 for _, t := range tasks { attr.allocated.Add(t.Resreq) } } } // 计算drf share max(r_alloc/r_total) drf.updateJobShare(job.Namespace, job.Name, attr) // 记录attr drf.jobAttrs[job.UID] = attr // 层级公平 if hierarchyEnabled { queue := ssn.Queues[job.Queue] // 已分配总资源 drf.totalAllocated.Add(attr.allocated) // 沿Tree更新层级树的node share drf.UpdateHierarchicalShare(root, drf.totalAllocated, job, attr, queue.Hierarchy, queue.Weights) } } // 注册抢占回调 preemptableFn := func(preemptor *api.TaskInfo, preemptees []*api.TaskInfo) ([]*api.TaskInfo, int) { ... // 计算抢占成功后的share latt := drf.jobAttrs[preemptor.Job] lalloc := latt.allocated.Clone().Add(preemptor.Resreq) _, ls := drf.calculateShare(lalloc, drf.totalResource) ... // 遍历被抢占的task for _, preemptee := range preemptees { // 未处理过,记录job attr if _, found := allocations[preemptee.Job]; !found { ratt := drf.jobAttrs[preemptee.Job] allocations[preemptee.Job] = ratt.allocated.Clone() } // 计算受害者释放资源的share ralloc := allocations[preemptee.Job].Sub(preemptee.Resreq) _, rs := drf.calculateShare(ralloc, drf.totalResource) // 抢占share更高的或相近的任务 if ls < rs || math.Abs(ls-rs) <= 0.000001 { addVictim(preemptee) } } return victims, util.Permit } ssn.AddPreemptableFn(drf.Name(), preemptableFn) // 层级公平 if hierarchyEnabled { // 注册queueOrderFn回调 queueOrderFn := func(l interface{}, r interface{}) int { lv := l.(*api.QueueInfo) rv := r.(*api.QueueInfo) // 由root到queue,避免资源分配倾斜 // 1.未饱和的node queue优先 // 2.share/wright小的优先 return drf.compareQueues(drf.hierarchicalRoot, lv, rv) } ssn.AddQueueOrderFn(drf.Name(), queueOrderFn) // 注册reclaimFn reclaimFn := func(reclaimer *api.TaskInfo, reclaimees []*api.TaskInfo) ([]*api.TaskInfo, int) { ... // update reclaim drf attr := drf.jobAttrs[ljob.UID] lattr := &drfAttr{ allocated: attr.allocated.Clone() } lattr.allocated.Add(reclaimer.Resreq) // 假设抢占完成,reclaimer资源计入totalAllocated totalAllocated.Add(reclaimer.Resreq) // 模拟计算reclaimer获取资源的share drf.updateShare(lattr) // 沿Tree更新层级树的node share drf.UpdateHierarchicalShare(root, totalAllocated, ljob, lattr, lqueue.Hierarchy, lqueue.Weights) // 寻找受害者 for _, preemptee := range reclaimees { ... // 由totalAllocated清理被抢占的资源 totalAllocated.Sub(preemptee.Resreq) ... // 模拟计算被抢占后受害者share rattr := &drfAttr{ allocated: attr.allocated.Clone() } rattr.allocated.Sub(preemptee.Resreq) drf.updateShare(rattr) // 沿Tree更新层级树的node share drf.UpdateHierarchicalShare(root, totalAllocated, rjob, rattr, rqueue.Hierarchy, rqueue.Weights) // 由root到queue对比,避免资源分配倾斜 // 1.未饱和的node queue优先 // 2.share/wright小的优先 ret := drf.compareQueues(root, lqueue, rqueue) // 恢复share现场(刚才是模拟,还未真正抢占) totalAllocated.Add(preemptee.Resreq) rattr.allocated.Add(preemptee.Resreq) drf.updateShare(rattr) drf.UpdateHierarchicalShare(root, totalAllocated, rjob, rattr, rqueue.Hierarchy, rqueue.Weights) // 资源由更富的队列流向reclaimer if ret < 0 { victims = append(victims, preemptee) } ... } return victims, util.Permit } ssn.AddReclaimableFn(drf.Name(), reclaimFn) } // 注册jobOrderFn jobOrderFn := func(l interface{}, r interface{}) int { ... // share相同 if drf.jobAttrs[lv.UID].share == drf.jobAttrs[rv.UID].share { return 0 } // share小的job优先 if drf.jobAttrs[lv.UID].share < drf.jobAttrs[rv.UID].share { return -1 } return 1 } ssn.AddJobOrderFn(drf.Name(), jobOrderFn) // 注册event handler. ssn.AddEventHandler(&framework.EventHandler{ // 申请资源 AllocateFunc: func(event *framework.Event) { // 更新job attr资源占用 attr := drf.jobAttrs[event.Task.Job] attr.allocated.Add(event.Task.Resreq) // 更新jobshare job := ssn.Jobs[event.Task.Job] drf.updateJobShare(job.Namespace, job.Name, attr) // 层级公平 if hierarchyEnabled { // 沿curNode向上更新node share queue := ssn.Queues[job.Queue] drf.totalAllocated.Add(event.Task.Resreq) drf.UpdateHierarchicalShare(drf.hierarchicalRoot, drf.totalAllocated, job, attr, ...) } }, // 释放资源 DeallocateFunc: func(event *framework.Event) { // 更新job attr资源占用 attr := drf.jobAttrs[event.Task.Job] attr.allocated.Sub(event.Task.Resreq) // 更新job share job := ssn.Jobs[event.Task.Job] drf.updateJobShare(job.Namespace, job.Name, attr) // 层级公平 if hierarchyEnabled { // 沿curNode向上更新node share queue := ssn.Queues[job.Queue] drf.totalAllocated.Sub(event.Task.Resreq) drf.UpdateHierarchicalShare(drf.hierarchicalRoot, drf.totalAllocated, job, attr, queue.Hierarchy, queue.Weights) } }, }) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190注意
share=max(rtotalralloc)
# 3.2.compare
drf。compareQueues()会根据前面的node share检查job/queue优先级,share/weight更小的job/queue优先调度,避免资源倾斜。func (drf *drfPlugin) compareQueues(root *hierarchicalNode, lqueue *api.QueueInfo, rqueue *api.QueueInfo) ... { // 拆分path node lnode := root lpaths := strings.Split(lqueue.Hierarchy, "/") rnode := root rpaths := strings.Split(rqueue.Hierarchy, "/") // 由root->queue层级检查 for i, depth := 0, min(len(lpaths), len(rpaths)); i < depth; i++ { // rnode用满资源份额 if !lnode.saturated && rnode.saturated { return -1 } // lnode用满资源份额 if lnode.saturated && !rnode.saturated { return 1 } // share/weight小的优先,相同转到下一级节点 if lnode.attr.share/lnode.weight == rnode.attr.share/rnode.weight { if i < depth-1 { lnode = lnode.children[lpaths[i+1]] rnode = rnode.children[rpaths[i+1]] } } else { return lnode.attr.share/lnode.weight - rnode.attr.share/rnode.weight } } return 0 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28注意
node资源用满会标记saturated,优先调度的总是饥饿的或share/weight小的,避免小任务长期阻塞
# 3.3.hierarchy
hierarchy会沿curNode向上检索层级树的parent,相关路径的node依次更新share,作为后续compareQueues优先级检查基准。func (drf *drfPlugin) UpdateHierarchicalShare(root *hierarchicalNode, totalAllocated *api.Resource, ...) { // 过滤活跃的资源(未用满) for _, rn := range drf.totalResource.ResourceNames() { if totalAllocated.Get(rn) < drf.totalResource.Get(rn) { demandingResources[rn] = true } } // 构建层级结构 drf.buildHierarchy(root, job, attr, hierarchy, hierarchicalWeights) // 更新层级share drf.updateHierarchicalShare(root, demandingResources) } // build hierarchy if the node does not exist func (drf *drfPlugin) buildHierarchy(root *hierarchicalNode, job *api.JobInfo, ...) { inode := root paths := strings.Split(hierarchy, "/") weights := strings.Split(hierarchicalWeights, "/") // 遍历path node for i := 1; i < len(paths); i++ { // curNode已经挂到树 if child, ok := inode.children[paths[i]]; ok { // 向下走一层 inode = child } else { // share weight fweight, _ := strconv.ParseFloat(weights[i], 64) if fweight < 1 { fweight = 1 } // 初始化child node child = &hierarchicalNode{ ... } // 注册child inode.children[paths[i]] = child child.parent = inode inode = child } } // job挂到树的叶子节点 child := &hierarchicalNode{ ... } inode.children[string(job.UID)] = child } // updateHierarchicalShare updates the node attribute recursively func (drf *drfPlugin) updateHierarchicalShare(node *hierarchicalNode, demanding map[v1.ResourceName]bool) { // 叶子节点 if node.children == nil { // 检查job node资源用满没有 node.saturated = resourceSaturated(node.attr.allocated, node.request, demanding) // 内部节点 } else { ... // 遍历child node for _, child := range node.children { // 递归更新share drf.updateHierarchicalShare(child, demanding) // skip empty child and saturated child if child.attr.share != 0 && !child.saturated { // 计算更新的share _, resShare := drf.calculateShare(child.attr.allocated, drf.totalResource) // 记录最小的share if resShare < mdr { mdr = resShare } } } ... // 遍历child node for _, child := range node.children { // 资源未用满 if !child.saturated { saturated = false } // only consider non-empty children if child.attr.share != 0 { // child用满资源 if child.saturated { // 直接更新node allocated t := child.attr.allocated node.attr.allocated.Add(t) // child未用满资源 } else { // child_allocated = allocated*(mdr/share),累计到node // 这样算,越富足的node计算的share越小 t := child.attr.allocated.Clone().Multi(mdr / child.attr.share) node.attr.allocated.Add(t) } } } // 更新node share node.attr.dominantResource, node.attr.share = drf.calculateShare(node.attr.allocated, drf.totalResource) // 更新node资源是否用满 node.saturated = saturated } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
注意
hierarchy更新会累计child allocated资源,若child node资源用满会直接累加,否则基于minShare/child.share作为影响因子调整
# 4.predicate
# 4.1.open
ssn.AddEventHandler()用于节点事件同步和资源管理,ssn.AddPrePredicateFn()会进行前置计算和preFilter注册,以提高调度效率。func (pp *predicatesPlugin) OnSessionOpen(ssn *framework.Session) { ... // 注册事件处理器 ssn.AddEventHandler(&framework.EventHandler{ // pod分配到node AllocateFunc: func(event *framework.Event) { // 更新task及构造pod pod := pl.UpdateTask(event.Task, event.Task.NodeName) nodeName := event.Task.NodeName // 获取sched node node, found := nodeMap[nodeName] if !found { return } // 获取nodeInfo nodeInfo, ok := ssn.Nodes[nodeName] if !ok { return } // GPU设备预选 for _, val := range api.RegisteredDevices { // node注册过该设备 if devices, ok := nodeInfo.Others[val].(api.Devices); ok { // pod未请求该设备 if !devices.HasDeviceRequest(pod) { continue } // 尝试分配GPU及更新Pod注解 devices.Allocate(ssn.KubeClient(), pod) ... } } // Pod注册到node(资源占用) node.AddPod(pod) }, // pod由node释放 DeallocateFunc: func(event *framework.Event) { // 更新task及构造Pod pod := pl.UpdateTask(event.Task, "") nodeName := event.Task.NodeName // pod关联node存在 node, found := nodeMap[nodeName] if !found { return } nodeInfo, ok := ssn.Nodes[nodeName] if !ok { return } // GPU设备释放 for _, val := range api.RegisteredDevices { // node注册过该设备 if devices, ok := nodeInfo.Others[val].(api.Devices); ok { // pod未请求该设备 if !devices.HasDeviceRequest(pod) { continue } // 释放GPU设备,清理Pod注解及deviceMap devices.Release(ssn.KubeClient(), pod) ... } } // Pod由node释放(资源释放) node.RemovePod(klog.FromContext(context.TODO()), pod) ... }, }) ... // 注册前置检查处理器 ssn.AddPrePredicateFn(pp.Name(), func(task *api.TaskInfo) error { // nodePort if predicate.nodePortEnable { // nodePort预计算&preFilter注册 _, status := nodePortFilter.PreFilter(context.TODO(), state, task.Pod) // 记录plugin是否跳过 handleSkipPrePredicatePlugin(status, task, skipPlugins, nodeports.Name) ... } // sidecar容器未支持&task使用sidecar(会引起资源计算模型混乱) if !features.EnableSidecarContainers && task.HasRestartableInitContainer { return fmt.Errorf("pod has a restartable init container and SidecarContainers feature is disabled") } // podAffinity if predicate.podAffinityEnable { // pod亲和预计算&preFilter注册 _, status := podAffinityFilter.PreFilter(context.TODO(), state, task.Pod) // 记录plugin是否跳过 handleSkipPrePredicatePlugin(status, task, skipPlugins, interpodaffinity.Name) ... } // podTopologySpread if predicate.podTopologySpreadEnable { // pod拓扑预计算&preFilter注册 _, status := podTopologySpreadFilter.PreFilter(context.TODO(), state, task.Pod) // 记录plugin是否跳过 handleSkipPrePredicatePlugin(status, task, skipPlugins, podTopologySpreadFilter.Name()) ... } return nil }) ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
注意
eventHandler会作为callbackFn最后调用进行资源更新或回收,prePredicateFn本质注册的都是kubernetes内置的plugin
# 4.2.predicate
predicate用于Pod调度,选中节点前会进行一系列的调度可行性检查,只有全部检查通过的节点才会进入后续的优选阶段,未完成则判定为不可调度。func (pp *predicatesPlugin) OnSessionOpen(ssn *framework.Session) { ... // 注册预选处理器 ssn.AddPredicateFn(pp.Name(), func(task *api.TaskInfo, node *api.NodeInfo) error { predicateStatus := make([]*api.Status, 0) // 关联node必须存在 nodeInfo, found := nodeMap[node.Name] if !found { return api.NewFitErrWithStatus(task, node, predicateStatus...) } // 调度Pod超出最大允许数量 if node.Allocatable.MaxTaskNum <= len(nodeInfo.Pods) { podsNumStatus := &api.Status{ Code: api.Unschedulable, ... Plugin: pp.Name() } predicateStatus = append(predicateStatus, podsNumStatus) } // 调度约束检查 predicateByStablefilter := func(nodeInfo *k8sframework.NodeInfo) ([]*api.Status, bool, error) { ... // 检查node不可调度标记&pod容忍不可调度污点 status := nodeUnscheduleFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, nodeUnscheduleStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(nodeUnscheduleStatus) { return predicateStatus, false, fmt.Errorf("plugin predicates failed") } } // nodeAffinity检查 if predicate.nodeAffinityEnable { // plugin设置的节点亲和检查/pod节点亲和性检查 status := nodeAffinityFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, nodeAffinityStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(nodeAffinityStatus) { return predicateStatus, false, fmt.Errorf("plugin predicates failed") } } } // taintToleration检查 if predicate.taintTolerationEnable { // pod对于节点污点容忍检查 status := tolerationFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, tolerationStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(tolerationStatus) { return predicateStatus, false, fmt.Errorf("plugin predicates failed") } } } return predicateStatus, true, nil } ... // predicate cache if predicate.cacheEnable { // 读取缓存结果 fit, err = pCache.PredicateWithCache(node.Name, task.Pod) // 未命中 if err != nil { // 重新检查调度约束 predicateCacheStatus, fit, _ = predicateByStablefilter(nodeInfo) // 更新缓存 pCache.UpdateCache(node.Name, task.Pod, fit) } ... // 未启用cache } else { // 检查调度约束 predicateCacheStatus, fit, _ = predicateByStablefilter(nodeInfo) } predicateStatus = append(predicateStatus, predicateCacheStatus...) // 节点不可调度 if !fit { return api.NewFitErrWithStatus(task, node, predicateStatus...) } // 端口冲突检查 if predicate.nodePortEnable { // 预检查plugin是否跳过 isSkip := handleSkipPredicatePlugin(task, skipPlugins, nodePortFilter.Name(), node) // 未跳过 if !isSkip { // pod端口和节点已占用端口冲突检测 status := nodePortFilter.Filter(context.TODO(), state, nil, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, nodePortStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(nodePortStatus) { return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } } // podAffinity检查 if predicate.podAffinityEnable { // 预检查plugin是否跳过 isSkip := handleSkipPredicatePlugin(task, skipPlugins, podAffinityFilter.Name(), node) if !isSkip { // pod亲和性检查/反亲和性检查/待调度node已有反亲和性检查 status := podAffinityFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { // 标记不可调度 podAffinityStatus.Code = api.UnschedulableAndUnresolvable predicateStatus = append(predicateStatus, podAffinityStatus) return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } // nodeVolumeLimit检查 if predicate.nodeVolumeLimitsEnable { // 检查CSI卷数量上限/Pod申请卷资源上限 status := nodeVolumeLimitsCSIFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, nodeVolumeStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(nodeVolumeStatus) { return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } // volumeZone检查 if predicate.volumeZoneEnable { // node设置分区标签&podPV所属分区与节点匹配 status := volumeZoneFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, volumeZoneStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(volumeZoneStatus) { return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } // podTopology检查 if predicate.podTopologySpreadEnable { // 预检查plugin是否跳过 isSkip := handleSkipPredicatePlugin(task, skipPlugins, podTopologySpreadFilter.Name(), node) if !isSkip { // 根据preFilter阶段计算的pod拓扑状态进行匹配: // 1.节点必须包含拓扑约束指定的 tpKey 标签 // 2.pod匹配拓扑约束的标签选择器才会计入分布数量 // 3.某个拓扑单元的分布数量偏差不允许超出最大限制 status := podTopologySpreadFilter.Filter(context.TODO(), state, task.Pod, nodeInfo) if api.ConvertPredicateStatus(status).Code != api.Success { predicateStatus = append(predicateStatus, podTopologyStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(podTopologyStatus) { return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } } // 资源比例检查 if predicate.proportionalEnable { // 资源比例限制,避免GPU节点被CPU/Mem应用吃光: // 1.Pod请求特殊资源,直接放行 // 2.基于特殊资源计算CPU/Mem预期空闲资源,Pod分配完成必须预留足额的空闲资源 proportionalStatus, _ := checkNodeResourceIsProportional(task, node, predicate.proportional) if proportionalStatus.Code != api.Success { predicateStatus = append(predicateStatus, proportionalStatus) // 终止检查(skip/wait/error/unschedulableAndUnresolvable) if ShouldAbort(proportionalStatus) { return api.NewFitErrWithStatus(task, node, predicateStatus...) } } } if len(predicateStatus) > 0 { return api.NewFitErrWithStatus(task, node, predicateStatus...) } return nil }) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187注意
predicate基本执行的还是kubernetes plugin进行节点筛选,仅实现了涉及的GPU检查
# 4.remaining
# 4.1.binpack
binpack会优先填满资源利用率高的节点,避免节点资源碎片化,资源利用率越高节点分数给的越高,打分会基于配置文件设置的资源权重进行。func (bp *binpackPlugin) OnSessionOpen(ssn *framework.Session) { ... // 节点打分 nodeOrderFn := func(task *api.TaskInfo, node *api.NodeInfo) (float64, error) { binPackingScore := BinPackingScore(task, node, bp.weight) return binPackingScore, nil } // 设置binpack资源权重,注册nodeOrderFn回调 if bp.weight.BinPackingWeight != 0 { ssn.AddNodeOrderFn(bp.Name(), nodeOrderFn) } } // BinPackingScore use the best fit polices during scheduling. // Goals: // - Schedule Jobs using BestFit Policy using Resource Bin Packing Priority Function // - Reduce Fragmentation of scarce resources on the Cluster func BinPackingScore(task *api.TaskInfo, node *api.NodeInfo, weight priorityWeight) float64 { ... for _, resource := range requested.ResourceNames() { // 获取某个资源的请求量 request := requested.Get(resource) if request == 0 { continue } // 获取某个资源的node可用量 allocate := allocatable.Get(resource) // 获取某个资源的节点已使用量 nodeUsed := used.Get(resource) // 获取某个资源的权重 resourceWeight, found := weight.BinPackingResources[resource] if !found { continue } // 计算资源分数 resourceScore, err := ResourceBinPackingScore(request, allocate, nodeUsed, resourceWeight) if err != nil { return 0 } // 更新总分和权重值 score += resourceScore weightSum += resourceWeight } // 归一化 if weightSum > 0 { score /= float64(weightSum) } // score= score *= float64(k8sFramework.MaxNodeScore * int64(weight.BinPackingWeight)) return score } // ResourceBinPackingScore calculate the binpack score for resource with provided info func ResourceBinPackingScore(requested, capacity, used float64, weight int) (float64, error) { if capacity == 0 || weight == 0 { return 0, nil } usedFinally := requested + used ... score := usedFinally * float64(weight) / capacity return score, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75注意
score=∑weight∑(capacity(request+used)⋅weight)×100×binPackingWeight
# 4.2.priority
priority会向session注册task/job优先级检查及抢占回调,preemptableFn根据task优先级决定被抢占资源的受害者列表。func (pp *priorityPlugin) OnSessionOpen(ssn *framework.Session) { // 注册taskOrderFn,优先级高的先调度 taskOrderFn := func(l interface{}, r interface{}) int { lv := l.(*api.TaskInfo) rv := r.(*api.TaskInfo) if lv.Priority == rv.Priority { return 0 } if lv.Priority > rv.Priority { return -1 } return 1 } ssn.AddTaskOrderFn(pp.Name(), taskOrderFn) // 注册jobOrderFn,优先级高的先调度 jobOrderFn := func(l, r interface{}) int { lv := l.(*api.JobInfo) rv := r.(*api.JobInfo) if lv.Priority > rv.Priority { return -1 } if lv.Priority < rv.Priority { return 1 } return 0 } ssn.AddJobOrderFn(pp.Name(), jobOrderFn) // 注册preemptableFn,过滤可被抢占的受害者task preemptableFn := func(preemptor *api.TaskInfo, preemptees []*api.TaskInfo) ([]*api.TaskInfo, int) { preemptorJob := ssn.Jobs[preemptor.Job] ... // 遍历被抢占的task for _, preemptee := range preemptees { // task所属的job preempteeJob := ssn.Jobs[preemptee.Job] // 受害者和抢占者job不同 if preempteeJob.UID != preemptorJob.UID { // 只抢占优先级更低的job task if preempteeJob.Priority < preemptorJob.Priority { victims = append(victims, preemptee) } // 受害者和抢占者job相同 } else { // 只抢占优先级更低的task if preemptee.Priority < preemptor.Priority { victims = append(victims, preemptee) } } } return victims, util.Permit } ssn.AddPreemptableFn(pp.Name(), preemptableFn) // 注册饥饿检查回调 jobStarvingFn := func(obj interface{}) bool { ji := obj.(*api.JobInfo) // (bound+binding+running+allocated+succeeded+pipelined)<job minAvailable return ji.ReadyTaskNum()+ji.WaitingTaskNum() < int32(len(ji.Tasks)) } ssn.AddJobStarvingFns(pp.Name(), jobStarvingFn) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70注意
priority基于优先级决定task调度顺序及被抢占的task列表
# 4.3.overcmt
overcommit基于overCommitFactor参数设置集群资源的容忍度,实际将集群的资源量扩大overCommitFactor倍,类似于资源超卖效果。/* actions: "enqueue, allocate, backfill" tiers: - plugins: - name: overcommit arguments: overcommit-factor: 1.0 */ func (op *overcommitPlugin) OnSessionOpen(ssn *framework.Session) { // 加载overCommitFactor因子 op.pluginArguments.GetFloat64(&op.overCommitFactor, overCommitFactor) if op.overCommitFactor < 1.0 { op.overCommitFactor = 1.2 } // 计算空闲资源 op.totalResource.Add(ssn.TotalResource) used := api.EmptyResource() for _, node := range ssn.Nodes { used.Add(node.Used) } // idle=total*factor-used op.idleResource = op.totalResource.Clone().Multi(op.overCommitFactor).SubWithoutAssert(used) // 遍历session job for _, job := range ssn.Jobs { // inqueue job resources(排除门控限制的task) if job.PodGroup.Status.Phase == scheduling.PodGroupInqueue && job.PodGroup.Spec.MinResources != nil { op.inqueueResource.Add(job.DeductSchGatedResources(job.GetMinResources())) continue } // running job resources if job.PodGroup.Status.Phase == PodGroupRunning && job.PodGroup.Spec.MinResources != nil && // 实际运行成员满足PodGroup最小要求 int32(util.CalculateAllocatedTaskNum(job)) >= job.PodGroup.Spec.MinMember { inqueued := util.GetInqueueResource(job, job.Allocated) op.inqueueResource.Add(job.DeductSchGatedResources(inqueued)) } } ssn.AddJobEnqueueableFn(op.Name(), func(obj interface{}) int { job := obj.(*api.JobInfo) idle := op.idleResource inqueue := api.EmptyResource() // 已分配资源 inqueue.Add(op.inqueueResource) // job未声明资源需求,放行入队 if job.PodGroup.Spec.MinResources == nil { return util.Permit } // job申请资源+已分配资源未超出空闲资源 jobMinReq := job.GetMinResources() if inqueue.Add(jobMinReq).LessEqualWithDimension(idle, jobMinReq) { // only compare requested resource return util.Permit } return util.Reject }) ssn.AddJobEnqueuedFn(op.Name(), func(obj interface{}) { job := obj.(*api.JobInfo) // job未声明资源需求 if job.PodGroup.Spec.MinResources == nil { return } // 记录到已分配资源 jobMinReq := job.GetMinResources() op.inqueueResource.Add(job.DeductSchGatedResources(jobMinReq)) }) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73注意
overcommit主要注册enqueueableFn和enqueueFn,检查job申请资源合法性,未溢出才允许入队