vcsdaction
# 1.allocate
# 1.1.execute
allocate是调度的核心实现,基于一系列的预选和优选算法为job task pod选择最合适的节点进行绑定,内部基于plugin的算法实现进行筛选。func (alloc *Action) Execute(ssn *framework.Session) { // 解析allocAction配置(predicateErrorCacheEnable) alloc.parseArguments(ssn) // queues sort queues by QueueOrderFn(plugin.QueueOrderFn-->先创建的优先-->uid小的优先) queues := util.NewPriorityQueue(ssn.QueueOrderFn) // used to find job with the highest priority in given queue(plugin.JobOrderFn-->先创建的优先-->uid小的优先) jobsMap := map[api.QueueID]*util.PriorityQueue{} alloc.session = ssn // 归类Queue和Job alloc.pickUpQueuesAndJobs(queues, jobsMap) alloc.allocateResources(queues, jobsMap) } func (alloc *Action) pickUpQueuesAndJobs(...) { ssn := alloc.session for _, job := range ssn.Jobs { // If not config enqueue action, change Pending pg into Inqueue state to avoid blocking job scheduling. if job.IsPending() { if conf.EnabledActionMap["enqueue"] { continue } else { job.PodGroup.Status.Phase = scheduling.PodGroupInqueue } } // job合法性检查未通过 if vr := ssn.JobValid(job); vr != nil && !vr.Pass { continue } // 关联Queue不存在 if _, found := ssn.Queues[job.Queue]; !found { continue } // 初始化Job堆及关联Queue入堆 if _, found := jobsMap[job.Queue]; !found { jobsMap[job.Queue] = util.NewPriorityQueue(ssn.JobOrderFn) queues.Push(ssn.Queues[job.Queue]) } // Job入堆 jobsMap[job.Queue].Push(job) } } // 1. 选出优先级最高Queue // 2. 选出Queue优先级最高Job // 3. 选出Job优先级最高Task // 4. 用predicateFn过滤不满足的节点 // 5. 选出剩余节点中最佳的节点 func (alloc *Action) allocateResources(...) { ssn := alloc.session pendingTasks := map[api.JobID]*util.PriorityQueue{} allNodes := ssn.NodeList // this action would make the resource usage among namespace balanced. for { if queues.Empty() { break } // 优先级最高的Queue queue := queues.Pop().(*api.QueueInfo) // plugin.Overused if ssn.Overused(queue) { // queue超出资源限制 continue } // queue关联job jobs, found := jobsMap[queue.UID] if !found || jobs.Empty() { continue } // 获取优先级最高job job := jobs.Pop().(*api.JobInfo) if _, found = pendingTasks[job.UID]; !found { // task堆(plugin.TaskOrderFn-->podIndex小的优先-->先创建的优先-->uid小的优先) tasks := util.NewPriorityQueue(ssn.TaskOrderFn) // job关联的pending task for _, task := range job.TaskStatusIndex[api.Pending] { // Skip tasks whose pod are scheduling gated if task.SchGated { continue } // Skip BestEffort task in allocate action. if task.Resreq.IsEmpty() { continue } // 收集pending task tasks.Push(task) } pendingTasks[job.UID] = tasks } tasks := pendingTasks[job.UID] if tasks.Empty() { // put queue back again and try other jobs in this queue queues.Push(queue) continue } // 申请task资源 alloc.allocateResourcesForTasks(tasks, job, jobs, queue, allNodes) // Put back the queue to priority queue after job's resource allocating finished. queues.Push(queue) } }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
注意
pickUpQueuesAndJobs负责收集可调度的queue/job,allocateResources负责基于job pending task预选和优先节点
# 1.2.allocres
alloc.allocateResourcesForTasks()负责基于预选优选算法筛选最合适的task node,基于bestNode提交或回滚资源申请确保Pod运行。// checks whether it can continue on allocating for current job func (ji *JobInfo) NeedContinueAllocating() bool { // all tasks running. if int(ji.MinAvailable) == len(ji.Tasks) { return false } failedRoles := ji.FitFailedRoles() // pendingTask with task role. pending := map[string]int32{} for _, task := range ji.TaskStatusIndex[Pending] { pending[task.TaskRole]++ } // 宽松模式(达到minAvailable就可以) if ji.MinAvailable < ji.TaskMinAvailableTotal { left := int32(0) // 未失败的pendingTask数量 for role, cnt := range pending { if _, ok := failedRoles[role]; !ok { left += cnt } } // ready+pending ≥ minAvailable return ji.ReadyTaskNum()+left >= ji.MinAvailable } // 严格模式 // 基于role划分的已分配task数量 allocated := ji.getJobAllocatedRoles() for role := range failedRoles { // role的独立min要求 min := ji.TaskMinAvailable[role] if min == 0 { continue } // the left task with same role can not be allocated, allocated number less than minAvailable if allocated[role] < min { return false } } return true } func (alloc *Action) allocateResourcesForTasks(...) { ... // 还有待调度任务 for !tasks.Empty() { // 优先级最高的task task := tasks.Pop().(*api.TaskInfo) // plugin.allocatableFn if !ssn.Allocatable(queue, task) { continue } // the task with its spec has already predicates failed if job.TaskHasFitErrors(task) { continue } // plugin.prePredicateFn前置检查 if err := ssn.PrePredicateFn(task); err != nil { ... // task无法预选 job.NodesFitErrors[task.UID] = fitErrors break } // 节点预选 predicateNodes, fitErrors := ph.PredicateNodes(task, allNodes, alloc.predicate...) // 无可调度节点 if len(predicateNodes) == 0 { // 更新job task err job.NodesFitErrors[task.UID] = fitErrors // job需继续分配 if job.NeedContinueAllocating() { continue } else { break } } ... for _, n := range predicateNodes { // 节点idle资源充足 if task.InitResreq.LessEqual(n.Idle, api.Zero) { idleCandidateNodes = append(idleCandidateNodes, n) // 节点idle+releasing-pipline资源充足 } else if task.InitResreq.LessEqual(n.FutureIdle(), api.Zero) { futureIdleCandidateNodes = append(futureIdleCandidateNodes, n) } } // 候选节点 candidateNodes = append(candidateNodes, idleCandidateNodes) candidateNodes = append(candidateNodes, futureIdleCandidateNodes) ... // 节点优选 for index, nodes := range candidateNodes { switch { case len(nodes) == 0: ... case len(nodes) == 1: // If only one node after predicate, just use it. bestNode = nodes[0] case len(nodes) > 1: // If more than one node after predicate, using "the best" one // 节点打分 nodeScores := util.PrioritizeNodes(task, nodes, ssn.BatchNodeOrderFn, ssn.NodeOrderMapFn, ssn.NodeOrderReduceFn) // plugin.BestNodeFn bestNode = ssn.BestNodeFn(task, nodeScores) if bestNode == nil { // 取score最高的某一个 bestNode = util.SelectBestNode(nodeScores) } } // If a proper node is found in idleCandidateNodes, skip futureIdleCandidateNodes if bestNode != nil { break } } // bestNode idle资源充足 if task.InitResreq.LessEqual(bestNode.Idle, api.Zero) { // 尝试由bestNode分配资源 if err := stmt.Allocate(task, bestNode); err != nil { // 失败回退 stmt.UnAllocate(task) ... } // 补充releasing资源充足 } else { // Allocate releasing resource to the task if any. if task.InitResreq.LessEqual(bestNode.FutureIdle(), api.Zero) { if err := stmt.Pipeline(task, bestNode.Name, false); err != nil { stmt.UnPipeline(task) ... } } } // plugin.JobReadyFn & 仍有task if ssn.JobReady(job) && !tasks.Empty() { // job重入队 jobs.Push(job) break } } // plugin.JobReadyFn if ssn.JobReady(job) { // 提交本次申请 stmt.Commit() } else { // plugin.JobPipelinedFn if !ssn.JobPipelined(job) { // 回滚本次申请 stmt.Discard() } } }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注意
bestNode资源分配存在两种情况,直接基于idle资源申请或idle+relesing-pipline申请,分配期间会调用plugin各类方法
# 1.3.bestnode
ssh.BestNodeFn()和util.SelectBestNode()用于选出最合适task运行的节点,前者基于plugin完成,后者取score最大的节点。// BestNodeFn invoke bestNode function of the plugins func (ssn *Session) BestNodeFn(task *api.TaskInfo, nodeScores map[float64][]*api.NodeInfo) *api.NodeInfo { for _, tier := range ssn.Tiers { for _, plugin := range tier.Plugins { ... // plugin.BestNodeFn pfn, found := ssn.bestNodeFns[plugin.Name] ... if bestNode := pfn(task, nodeScores); bestNode != nil { return bestNode } } } return nil } // returns best node whose score is highest, pick one randomly if there are many nodes with same score. func SelectBestNode(nodeScores map[float64][]*api.NodeInfo) *api.NodeInfo { ... // 直接取score最大的 for score, nodes := range nodeScores { if score > maxScore { maxScore = score bestNodes = nodes } } if len(bestNodes) == 0 { return nil } // 选其中一个节点 return bestNodes[rand.Intn(len(bestNodes))] }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
注意
bestNode匹配还是优先基于plugin实现,无法满足才基于score直接取分数最大的节点组
# 2.preempt
# 2.1.execute
preempt action负责根据优先级规则为同一队列中高优先级任务进行抢占式调度,queue内的不同job和job内的task均可以参与抢占调度。func (pmpt *Action) Execute(ssn *framework.Session) { ... for _, job := range ssn.Jobs { // pending忽略,wait enqueue action标记 if job.IsPending() { continue } // job检查未通过 if vr := ssn.JobValid(job); vr != nil && !vr.Pass { continue } // 记录关联queue信息 if queue, found := ssn.Queues[job.Queue]; !found { continue } else if _, existed := queues[queue.UID]; !existed { queues[queue.UID] = queue } // 检查job是否处于饥饿状态(plugin.jobStarvingFn) if ssn.JobStarving(job) { // 基于Queue分类管理抢占job(plugin.JobOrder-->先创建的优先-->UID小的优先) if _, found := preemptorsMap[job.Queue]; !found { preemptorsMap[job.Queue] = util.NewPriorityQueue(ssn.JobOrderFn) } preemptorsMap[job.Queue].Push(job) // 记录资源不足的job列表 underRequest = append(underRequest, job) // 基于Queue分类管理job的抢占task(plugin.TaskOrder-->index小的优先-->先创建的优先-->UID小的优先) preemptorTasks[job.UID] = util.NewPriorityQueue(ssn.TaskOrderFn) for _, task := range job.TaskStatusIndex[api.Pending] { if task.SchGated { continue } // 记录pending task preemptorTasks[job.UID].Push(task) } } } ph := util.NewPredicateHelper() // Preemption between Jobs within Queue. for _, queue := range queues { for { // queue对应job队列 preemptors := preemptorsMap[queue.UID] // If no preemptors, no preemption. if preemptors == nil || preemptors.Empty() { break } // 优先级最高的job preemptorJob := preemptors.Pop().(*api.JobInfo) stmt := framework.NewStatement(ssn) ... for { // job饥饿状态检查 if !ssn.JobStarving(preemptorJob) { break } // job没有抢占的pending task if preemptorTasks[preemptorJob.UID].Empty() { break } // 优先级最高的pending task preemptor := preemptorTasks[preemptorJob.UID].Pop().(*api.TaskInfo) // 进行task抢占 assigned, err = pmpt.preempt(ssn, stmt, preemptor, func(task *api.TaskInfo) bool { // 只抢占bounding/running任务 if !api.PreemptableStatus(task.Status) { return false } // bestEffort pod不能抢占非bestEffort pod. if preemptor.BestEffort && !task.BestEffort { return false } // task不允许抢占 if !task.Preemptable { return false } job, found := ssn.Jobs[task.Job] ... // 仅允许抢占同Queue的不同Job return job.Queue == preemptorJob.Queue && preemptor.Job != task.Job }, ph) ... } // plugin.JobPipelinedFn(未否决) if ssn.JobPipelined(preemptorJob) { // 提交抢占结果 stmt.Commit() } else { // 回滚抢占结果 stmt.Discard() continue } // 分配到资源 if assigned { // 入队继续抢占 preemptors.Push(preemptorJob) } } // Preemption between Task within Job. for _, job := range underRequest { // preemptor numbers lose when in same job preemptorTasks[job.UID] = util.NewPriorityQueue(ssn.TaskOrderFn) // job pending task for _, task := range job.TaskStatusIndex[api.Pending] { // skip scheduling gated tasks if task.SchGated { continue } preemptorTasks[job.UID].Push(task) } for { if _, found := preemptorTasks[job.UID]; !found { break } if preemptorTasks[job.UID].Empty() { break } // job内优先级最高的task preemptor := preemptorTasks[job.UID].Pop().(*api.TaskInfo) stmt := framework.NewStatement(ssn) assigned, err := pmpt.preempt(ssn, stmt, preemptor, func(task *api.TaskInfo) bool { // 只抢占bounding/running任务 if !api.PreemptableStatus(task.Status) { return false } // bestEffort pod不能抢占非bestEffort pod. if preemptor.BestEffort && !task.BestEffort { return false } // task不允许抢占 if !task.Preemptable { return false } // 仅允许抢占同Job的不同Task return preemptor.Job == task.Job }, ph) ... // 提交抢占结果 stmt.Commit() // 未分配到资源,继续处理下一个job if !assigned { break } } } } // 驱逐被抢占的task victimTasks(ssn) }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
注意
preempt会基于queue/job尝试抢占调度job/task,核心逻辑均为pmpt.preempt
# 2.2.preempt
pmpt.preempt()用于尝试为高优先级任务腾出足够资源,资源不足会抢占低优先级任务实现资源释放,实现高优先级任务抢占调度至目标节点。func (pmpt *Action) preempt(...) (bool, error) { // 抢占策略检查(Never不允许抢占) if err := pmpt.taskEligibleToPreempt(preemptor); err != nil { return false, err } ... // 前置检查——PVC绑定/节点亲和/拓扑限制(plugin.prePredicateFn) ssn.PrePredicateFn(preemptor) ... // plugin.PredicateFn(检测UnschedulableAndUnresolvable和ErrorSkipOrWait) predicateFn := ssn.PredicateForPreemptAction //filter out those nodes that are UnschedulableAndUnresolvable status got in allocate action allNodes := ssn.GetUnschedulableAndUnresolvableNodesForTask(preemptor) // 预选 predicateNodes, _ := predicateHelper.PredicateNodes(preemptor, allNodes, predicateFn, ...) // 优选 nodeScores := util.PrioritizeNodes(preemptor, predicateNodes, ssn.BatchNodeOrderFn, ssn.NodeOrderMapFn, ssn.NodeOrderReduceFn) selectedNodes := util.SortNodes(nodeScores) // job所在queue job, found := ssn.Jobs[preemptor.Job] ... currentQueue := ssn.Queues[job.Queue] // 由最优节点开始遍历 for _, node := range selectedNodes { ... // 检查node可抢占task for _, task := range node.Tasks { if filter == nil { preemptees = append(preemptees, task.Clone()) // 抢占条件检查 } else if filter(task) { preemptees = append(preemptees, task.Clone()) } } // plugin.PreemptableFn(插件决定可抢占的task,所有plugin结果取交集) victims := ssn.Preemptable(preemptor, preemptees) ... // 检查被抢占task释放的资源是否满足需求 if err := util.ValidateVictims(preemptor, node, victims); err != nil { continue } // 排序victims(task优先级低的-->queue优先级低的-->job优先级低的) victimsQueue := ssn.BuildVictimsPriorityQueue(victims, preemptor) // Preempt victims for tasks, pick lowest priority task first. preempted := api.EmptyResource() for !victimsQueue.Empty() { // plugin允许Queue继续向抢占task分配资源&节点剩余资源充足 if ssn.Allocatable(currentQueue, preemptor) && preemptor.InitResreq.LessEqual(node.FutureIdle(), 0){ break } // 获取优先级最低的task preemptee := victimsQueue.Pop().(*api.TaskInfo) // 驱逐task if err := stmt.Evict(preemptee, "preempt"); err != nil { continue } // 累计释放资源 preempted.Add(preemptee.Resreq) } evictionOccurred := false // 发生过驱逐 if !preempted.IsEmpty() { evictionOccurred = true } // plugin允许Queue继续向抢占task分配资源&节点剩余资源充足 if ssn.Allocatable(currentQueue, preemptor) && preemptor.InitResreq.LessEqual(node.FutureIdle(), 0) { // 进行调度资源扣减 if err := stmt.Pipeline(preemptor, node.Name, evictionOccurred); err != nil { // 失败回滚 stmt.UnPipeline(preemptor) ... } // Ignore pipeline error, will be corrected in next scheduling loop. assigned = true break } } return assigned, 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注意
pmpt.preempt本质也会进行预选和优先,抢占期间会考虑节点资源情况、任务优先级,还会基于plugin过滤可抢占任务及检查分配许可
# 3.backfill
# 3.1.execute
backfill是调度流程的回填步骤,处理待调度Pod列表未设置资源申请的任务,调度期间遍历所有预选节点,满足任务调度需求就将Pod调度到节点。func (backfill *Action) Execute(ssn *framework.Session) { ... // 节点承载Task可行性检查Hook predicateFunc := ssn.PredicateForAllocateAction // 可以backfill的pendingTask pendingTasks := backfill.pickUpPendingTasks(ssn) for _, task := range pendingTasks { job := ssn.Jobs[task.Job] ph := util.NewPredicateHelper() ... // 前置检查(plugin.PrePredicateFn) if err := ssn.PrePredicateFn(task); err != nil { ... job.NodesFitErrors[task.UID] = fe break } // 节点预选 predicateNodes, fitErrors := ph.PredicateNodes(task, ssn.NodeList, predicateFunc, ...) if len(predicateNodes) == 0 { job.NodesFitErrors[task.UID] = fitErrors break } node := predicateNodes[0] if len(predicateNodes) > 1 { // 节点打分 nodeScores := util.PrioritizeNodes(task, predicateNodes, ssn.BatchNodeOrderFn, ...) // plugin.BestNodeFn选出第一个插件通过的节点 node = ssn.BestNodeFn(task, nodeScores) if node == nil { // 取分值最大的其中一个节点 node = util.SelectBestNode(nodeScores) } } // 扣减资源占用及分发任务(sc.BindFlowChannel <- taskInfo) if err := ssn.Allocate(task, node); err != nil { fe.SetNodeError(node.Name, err) continue } ... } }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注意
backfill处理的是未声明资源需求的任务,属于调度最后的回填逻辑
# 3.2.picktask
backfill.pickUpPendingTasks()用于遍历session job,基于优先级组织queue/job/task,将pending task基于优先级排序到列表。func (backfill *Action) pickUpPendingTasks(ssn *framework.Session) []*api.TaskInfo { ... // 遍历session job for _, job := range ssn.Jobs { if job.IsPending() { continue } // job条件检查 if vr := ssn.JobValid(job); vr != nil && !vr.Pass { continue } queue, found := ssn.Queues[job.Queue] if !found { continue } // job相关pendingTask for _, task := range job.TaskStatusIndex[api.Pending] { // 仅关心bestEffort类型任务 if !task.BestEffort { continue } // 未准备好调度 if task.SchGated { continue } // task入队 if _, existed := tasks[job.UID]; !existed { tasks[job.UID] = util.NewPriorityQueue(ssn.TaskOrderFn) } tasks[job.UID].Push(task) } // job相关pipelinedTask for _, task := range job.TaskStatusIndex[api.Pipelined] { // // 仅关心bestEffort类型任务 if !task.BestEffort { continue } // 回滚资源预留重新计算(低优任务) stmt := framework.NewStatement(ssn) stmt.UnPipeline(task) ... // task入队 if _, existed := tasks[job.UID]; !existed { tasks[job.UID] = util.NewPriorityQueue(ssn.TaskOrderFn) } tasks[job.UID].Push(task) } if _, existed := tasks[job.UID]; !existed { continue } // job相关queue入队 if _, existed := jobs[queue.UID]; !existed { queues.Push(queue) jobs[job.Queue] = util.NewPriorityQueue(ssn.JobOrderFn) } // job入队 jobs[job.Queue].Push(job) } for !queues.Empty() { // 优先级最高的queue queue, ok := queues.Pop().(*api.QueueInfo) ... // queue相关job处理 for !jobs[queue.UID].Empty() { // 优先级最高的job job, ok := jobs[queue.UID].Pop().(*api.JobInfo) ... // job相关的pendingTask,基于优先级排序 for !tasks[job.UID].Empty() { pendingTasks = append(pendingTasks, tasks[job.UID].Pop().(*api.TaskInfo)) } } } return pendingTasks }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注意
backfill.pickUpPendingTasks()本质上基于优先级梳理queue/job/task,将task根据优先级排序放到列表
# 4.remaining
# 4.1.enqueue
enqueue action是调度的准备阶段,集群资源满足作业调度的最小需求,作业状态会由pending转为enqueue,资源检查由关联的plugin完成。func (enqueue *Action) Execute(ssn *framework.Session) { // queue堆(plugin.QueueOrderFn-->先创建的优先-->uid小的优先) queues := util.NewPriorityQueue(ssn.QueueOrderFn) queueSet := sets.NewString() // job堆(plugin.JobOrderFn-->先创建的优先-->uid小的优先) jobsMap := map[api.QueueID]*util.PriorityQueue{} for _, job := range ssn.Jobs { if job.ScheduleStartTimestamp.IsZero() { ssn.Jobs[job.UID].ScheduleStartTimestamp = metav1.Time{ Time: time.Now(), } } // 获取job queue if queue, found := ssn.Queues[job.Queue]; !found { continue } else if !queueSet.Has(string(queue.UID)) { // queue入堆 queueSet.Insert(string(queue.UID)) queues.Push(queue) } // pg为空/pgPending/pgPhase为空 if job.IsPending() { if _, found := jobsMap[job.Queue]; !found { jobsMap[job.Queue] = util.NewPriorityQueue(ssn.JobOrderFn) } // job入堆 jobsMap[job.Queue].Push(job) } } for { if queues.Empty() { break } // 获取优先级最高的queue queue := queues.Pop().(*api.QueueInfo) // skip the Queue that has no pending job jobs, found := jobsMap[queue.UID] if !found || jobs.Empty() { continue } // 获取优先级最高的job job := jobs.Pop().(*api.JobInfo) // job未设置最小资源限制 || job申请资源通过 if job.PodGroup.Spec.MinResources == nil || ssn.JobEnqueueable(job) { // plugin.jobEnqueued ssn.JobEnqueued(job) // 更新job.pg状态为Inqueue job.PodGroup.Status.Phase = scheduling.PodGroupInqueue // 记录jobs ssn.Jobs[job.UID] = job } // Added Queue back until no job in Queue. queues.Push(queue) } }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
注意
enqueue action基于plugin依次处理queue job至queue为空,相关的资源检查及入队处理均依赖plugin实现
# 4.2.reclaim
reclaim用于根据队列权重回收资源,负责将空闲队列的资源临时复用给超额队列,确保集群任务的高效运行及资源复用,提高资源整体的利用上限。func (ra *Action) Execute(ssn *framework.Session) { ... // 遍历session job for _, job := range ssn.Jobs { if job.IsPending() { continue } // plugin.JobValidFn if vr := ssn.JobValid(job); vr != nil && !vr.Pass { continue } // queue入队 if queue, found := ssn.Queues[job.Queue]; !found { continue } else if _, existed := queueMap[queue.UID]; !existed { queueMap[queue.UID] = queue queues.Push(queue) } // job饥饿检查(plugin.JobStarving) if ssn.JobStarving(job) { if _, found := preemptorsMap[job.Queue]; !found { preemptorsMap[job.Queue] = util.NewPriorityQueue(ssn.JobOrderFn) } // queue相关job preemptorsMap[job.Queue].Push(job) // job相关pendingTask preemptorTasks[job.UID] = util.NewPriorityQueue(ssn.TaskOrderFn) for _, task := range job.TaskStatusIndex[api.Pending] { if task.SchGated { continue } preemptorTasks[job.UID].Push(task) } } } for { // If no queues, break if queues.Empty() { break } ... // 优先级最高的queue queue := queues.Pop().(*api.QueueInfo) // queue超出资源限制(plugin.overusedFn) if ssn.Overused(queue) { continue } // queue优先级最高的job jobs, found := preemptorsMap[queue.UID] if !found || jobs.Empty() { continue } else { job = jobs.Pop().(*api.JobInfo) } // job没有pendingTask||job不饥饿 if tasks, found := preemptorTasks[job.UID]; !found || tasks.Empty() || !ssn.JobStarving(job) { continue // job优先级最高的task } else { task = tasks.Pop().(*api.TaskInfo) } // task不允许抢占 if task.Pod.Spec.PreemptionPolicy != nil && *task.Pod.Spec.PreemptionPolicy == v1.PreemptNever { jobs.Push(job) queues.Push(queue) continue } // queue task不允许抢占(plugin.PreemptiveFn) if !ssn.Preemptive(queue, task) { continue } // task预检查(plugin.PrePredicateFn) if err := ssn.PrePredicateFn(task); err != nil { continue } ... // filter out those nodes that are UnschedulableAndUnresolvable status got in allocate action totalNodes := ssn.GetUnschedulableAndUnresolvableNodesForTask(task) for _, n := range totalNodes { // 节点可行性检查(plugin.PredicateFn) if err := ssn.PredicateForPreemptAction(task, n); err != nil { continue } ... for _, task := range n.Tasks { // Ignore non running task. if task.Status != api.Running { continue } // task不允许抢占 if !task.Preemptable { continue } // task关联job不存在 if j, found := ssn.Jobs[task.Job]; !found { continue // 仅抢占其它queue的task } else if j.Queue != job.Queue { q := ssn.Queues[j.Queue] // queue不允许回收(spec.reclaimable) if !q.Reclaimable() { continue } // Clone task to avoid modify Task's status on node. reclaimees = append(reclaimees, task.Clone()) } } if len(reclaimees) == 0 { continue } // 过滤可驱逐的task——交集(plugin.ReclaimableFn) victims := ssn.Reclaimable(task, reclaimees) // 检查驱逐空出的资源 if err := util.ValidateVictims(task, n, victims); err != nil { continue } // 基于优先级排序回收任务 victimsQueue := ssn.BuildVictimsPriorityQueue(victims, task) ... // Reclaim victims for tasks. for !victimsQueue.Empty() { // 优先级最高的victimTask reclaimee := victimsQueue.Pop().(*api.TaskInfo) // 驱逐 if err := ssn.Evict(reclaimee, "reclaim"); err != nil { continue } // 记录回收的资源 reclaimed.Add(reclaimee.Resreq) // reclaimed enough resources, break loop to avoid Sub panic. if resreq.LessEqual(reclaimed, api.Zero) { break } } // 回收的资源充足 if task.InitResreq.LessEqual(reclaimed, api.Zero) { // 预占用资源 ssn.Pipeline(task, n.Name) ... // Ignore error of pipeline, will be corrected in next scheduling loop. assigned = true break } } // 分配到资源,继续下一个job task if assigned { jobs.Push(job) } queues.Push(queue) } }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注意
reclaim会尝试抢占低优先级任务的资源供高优先级任务复用,确保高优任务完成
# 4.3.shuffle
shuffle的操作简单粗暴,直接将所有运行的job task驱逐重调度,相对来说比较简单,目前未识别出使用场景及用意。// Execute select evictees according given strategies and evict them. func (shuffle *Action) Execute(ssn *framework.Session) { ... // select pods that may be evicted for _, jobInfo := range ssn.Jobs { for _, taskInfo := range jobInfo.Tasks { // running task才值得驱逐 if taskInfo.Status == api.Running { tasks = append(tasks, taskInfo) } } } // 过滤待驱逐任务(plugin.victimTaskFn) victims := ssn.VictimTasks(tasks) for victim := range victims { // 执行驱逐及资源恢复(updateTaskStatus-->node.UpdateTask-->callbacks) if err := ssn.Evict(victim, "shuffle"); err != nil { continue } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24