schedpriority
# 1.入口
# 1.1.schedulingCycle
schedulingCycle()是调度周期核心实现,会依次进行节点预选、优选、提名及许可检查,以筛选承载pod的最佳node完成调度。
# 1.2.schedulePod
sched.SchedulePod()尝试由nodes列表为pod选择一个合适的node,分为预选(多个)和打分两个阶段,打分最高的节点就是调度目标。// schedulePod tries to schedule the given pod to one of the nodes in the node list. func (sched *Scheduler) schedulePod(...) (result ScheduleResult, err error) { ... // 更新快照(未删除节点缓存的备份) if err := sched.Cache.UpdateSnapshot(sched.nodeInfoSnapshot); err != nil { return result, err } ... // 没有节点可用 if sched.nodeInfoSnapshot.NumNodes() == 0 { return result, ErrNoNodesAvailable } // 节点预选 feasibleNodes, diagnosis, err := sched.findNodesThatFitPod(ctx, fwk, state, pod) ... // 节点优选 priorityList, err := prioritizeNodes(ctx, sched.Extenders, fwk, state, pod, feasibleNodes) ... // 获取分值最高的节点 host, err := selectHost(priorityList) ... return ScheduleResult{ // 建议节点 SuggestedHost: host, // 评估的节点 EvaluatedNodes: len(feasibleNodes) + len(diagnosis.NodeToStatusMap), // 预选的节点 FeasibleNodes: len(feasibleNodes), }, err }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注意
schedulePod()最重要的两步就是预选和优选阶段,前者匹配满足条件节点,后者匹配最优节点
# 2.预选
# 2.1.findNodesThatFitPod
sched.findNodesThatFitPod()用于实现调度器的预选过程,内部调用插件的preFilter和filter方法筛选满足pod要求的节点。// Filters the nodes to find the ones that fit the pod based on the framework filter plugins and filter extender func (sched *Scheduler) findNodesThatFitPod(...) ([]*v1.Node, framework.Diagnosis, error) { ... // 获取可用node allNodes, err := sched.nodeInfoSnapshot.NodeInfos().List() ... // 执行调度框架的preFilter插件 preRes, s := fwk.RunPreFilterPlugins(ctx, state, pod) ... // 优先尝试提名节点 if len(pod.Status.NominatedNodeName) > 0 { // 评估提名节点可否调度 feasibleNodes, err := sched.evaluateNominatedNode(ctx, pod, fwk, state, diagnosis) ... // 提名节点可行,提前结束 if len(feasibleNodes) != 0 { return feasibleNodes, diagnosis, nil } } nodes := allNodes if !preRes.AllNodes() { nodes = make([]*framework.NodeInfo, 0, len(preRes.NodeNames)) // 根据preFilter结果收缩节点列表 for nodeName := range preRes.NodeNames { // 记录合法的节点 if nodeInfo, err := sched.nodeInfoSnapshot.Get(nodeName); err == nil { nodes = append(nodes, nodeInfo) } } } // 执行调度框架的filter插件限制可选node feasibleNodes, err := sched.findNodesThatPassFilters(ctx, fwk, state, pod, diagnosis, nodes) ... // 更新下一次的起始索引 sched.nextStartNodeIndex = (sched.nextStartNodeIndex + processedNodes) % len(allNodes) ... // 执行调度器的extender扩展进行预选 feasibleNodesAfterExtender, err := findNodesThatPassExtenders(sched.Extenders, pod, feasibleNodes, diagnosis.NodeToStatusMap) ... return feasibleNodesAfterExtender, diagnosis, nil } // RunPreFilterPlugins runs the set of configured PreFilter plugins. func (f *frameworkImpl) RunPreFilterPlugins(...) (_ *framework.PreFilterResult, status *framework.Status) { ... // 遍历preFilter插件 for _, pl := range f.preFilterPlugins { // 执行preFilter插件 r, s := f.runPreFilterPlugin(ctx, pl, state, pod) // 跳过 if s.IsSkip() { skipPlugins.Insert(pl.Name()) continue } ... // 计算插件节点结果交集 result = result.Merge(r) // 交集为空,不可调度 if !result.AllNodes() && len(result.NodeNames) == 0 { ... return nil, framework.NewStatus(framework.Unschedulable, msg) } } return result, 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注意
findNodesThatFitPod()是调度器的预选入口,内部会根据可用节点依次执行preFilter、Filter和Extender插件筛选满足的节点集合
# 2.2.findNodesThatPassExtenders
findNodesThatPassExtenders()是调度器的扩展函数,用于扩展调度器的核心功能,基于http交互外部调度器以影响调度行为。func findNodesThatPassExtenders(...) ([]*v1.Node, error) { // 遍历extender扩展 for _, extender := range extenders { // 无可用节点 if len(feasibleNodes) == 0 { break } // extender未关心该类pod if !extender.IsInterested(pod) { continue } // 执行外部过滤调用 feasibleList, failedMap, failedAndUnresolvableMap, err := extender.Filter(pod, feasibleNodes) ... feasibleNodes = feasibleList } return feasibleNodes, nil } // Filter based on extender implemented predicate functions. The filtered list is // expected to be a subset of the supplied list; otherwise the function returns an error. func (h *HTTPExtender) Filter(...) ([]*v1.Node,failedNodes,failedAndUnresolvableNodes extenderv1.FailedNodesMap,err error) { ... // 建立节点名到节点对象映射 for _, n := range nodes { fromNodeName[n.Name] = n } // 未实现filter接口 if h.filterVerb == "" { return nodes, extenderv1.FailedNodesMap{}, extenderv1.FailedNodesMap{}, nil } ... // 请求参数 args = &extenderv1.ExtenderArgs{ Pod: pod, Nodes: nodeList, NodeNames: nodeNames, } // 发送http请求 h.send(h.filterVerb, args, &result) ... // 轻量模式 if h.nodeCacheCapable && result.NodeNames != nil { nodeResult = make([]*v1.Node, len(*result.NodeNames)) for i, nodeName := range *result.NodeNames { if n, ok := fromNodeName[nodeName]; ok { nodeResult[i] = n } else { return nil, nil, nil, fmt.Errorf("extender %q claims a filtered node %q which is not found") } } // 完整模式 } else if result.Nodes != nil { nodeResult = make([]*v1.Node, len(result.Nodes.Items)) for i := range result.Nodes.Items { nodeResult[i] = &result.Nodes.Items[i] } } return nodeResult, result.FailedNodes, result.FailedAndUnresolvableNodes, 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注意
extender提供非侵入scheduler core的方式扩展功能,相对来说性能较差且无法共享scheduler维护的cache,已不推荐使用
# 2.3.findNodesThatPassFilters
findNodesThatPassFilters()用于遍历执行调度框架的Filter插件,为减少预选阶段的时延,调度框架内部利用Parallelizer并发控制器启动16个协程并发调用插件,通过numFeasibleNodesToFind()方法减少扫描计算的nodes数量。// findNodesThatPassFilters finds the nodes that fit the filter plugins. func (sched *Scheduler) findNodesThatPassFilters(...) ([]*v1.Node, error) { numAllNodes := len(nodes) // 计算需要扫描的node数量,避免超大集群节点扫描 // 节点数低于100直接取出扫描 // 节点数超出100,根据公式 numAllNodes*(50-numAllNodes/125)/100 计算 numNodesToFind := sched.numFeasibleNodesToFind(fwk.PercentageOfNodesToScore(), int32(numAllNodes)) ... // 未注册过滤插件 if !fwk.HasFilterPlugins() { // 由调度索引取出具体数量node for i := range feasibleNodes { feasibleNodes[i] = nodes[(sched.nextStartNodeIndex+i)%numAllNodes].Node() } return feasibleNodes, nil } ... checkNode := func(i int) { // 获取nodeInfo对象 nodeInfo := nodes[(sched.nextStartNodeIndex+i)%numAllNodes] // 遍历执行调度框架的Filter插件 status := fwk.RunFilterPluginsWithNominatedPods(ctx, state, pod, nodeInfo) ... // 执行成功 if status.IsSuccess() { // 累加成功执行插件的节点数量 length := atomic.AddInt32(&feasibleNodesLen, 1) // 结束后续任务 if length > numNodesToFind { cancel() atomic.AddInt32(&feasibleNodesLen, -1) // 记录可提名节点 } else { feasibleNodes[length-1] = nodeInfo.Node() } } ... } ... // 并发执行调度框架的Filter插件 fwk.Parallelizer().Until(ctx, numAllNodes, checkNode, metrics.Filter) // 可提名节点 feasibleNodes = feasibleNodes[:feasibleNodesLen] ... return feasibleNodes, 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注意
findNodesThatPassFilters()本质上并发执行Filter插件检查node可调度pod,取前N个节点作为推荐
# 2.4.numFeasibleNodesToFind
numFeasibleNodesToFind()根据集群可选节点数计算出参与预选的节点数量以缩小范围,避免集群规模膨胀造成无效的资源计算开销,计算结果代表的是取执行Filter插件成功的节点数,前面的节点执行插件失败还是会继续扫描。// numFeasibleNodesToFind returns the number of feasible nodes that once found, the scheduler stops // its search for more feasible nodes. func (sched *Scheduler) numFeasibleNodesToFind(percentageOfNodesToScore *int32, numAllNodes int32) int32 { // 节点数低于100取完整 if numAllNodes < minFeasibleNodesToFind { return numAllNodes } ... // 设置调度框架百分比 if percentageOfNodesToScore != nil { percentage = *percentageOfNodesToScore // 否则设置全局调度器百分比 } else { percentage = sched.percentageOfNodesToScore } // 未设置,计算通用百分比(5%~50%) if percentage == 0 { // 计算百分比 percentage = int32(50) - numAllNodes/125 // 基于最小值修正百分比(5%) if percentage < minFeasibleNodesPercentageToFind { percentage = minFeasibleNodesPercentageToFind } } // 计算参与预选节点 numNodes = numAllNodes * percentage / 100 // 计算出的预选节点低于100进行修正(100) if numNodes < minFeasibleNodesToFind { return minFeasibleNodesToFind } return numNodes }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注意
numFeasibleNodesToFind()的计算策略会将预选节点计算范围控制在5%~50%之间,预选节点数太少则取最小值100进行修正
# 2.5.runFilterPluginsWithNominatedPods
runFilterPluginsWithNominatedPods()用于执行Filter插件对各node进行检查,确定node能否承载pod运行。// RunFilterPluginsWithNominatedPods runs the set of configured filter plugins for nominated pod on the node. func (f *frameworkImpl) RunFilterPluginsWithNominatedPods(...) *framework.Status { var status *framework.Status podsAdded := false // 根据提名状态执行插件 for i := 0; i < 2; i++ { stateToUse := state nodeInfoToUse := info if i == 0 { // 第一轮尝试将node提名的高优先级pod加入,以确定可否承载当前pod 第一轮添加被提名pod到周期状态和节点信息 podsAdded, stateToUse, nodeInfoToUse, err = addNominatedPods(ctx, f, pod, state, info) ... // 上一轮未增加高优先级提名pod或执行失败不再进行下一轮 } else if !podsAdded || !status.IsSuccess() { break } // 执行Filter插件 status = f.RunFilterPlugins(ctx, stateToUse, pod, nodeInfoToUse) // 失败及非不可调度结束检查 if !status.IsSuccess() && !status.IsUnschedulable() { return status } } return status } // 临时将已提名的高优先级Pod绑定到Node再考虑 func addNominatedPods(...) (bool, *framework.CycleState, *framework.NodeInfo, error) { ... // 获取节点的提名pod列表,提名pod通常抢占阶段产生 nominatedPodInfos := fh.NominatedPodsForNode(nodeInfo.Node().Name) ... // 遍历提名pod for _, pi := range nominatedPodInfos { // 仅关心优先级比当前pod高的提名pod if corev1.PodPriority(pi.Pod) >= corev1.PodPriority(pod) && pi.Pod.UID != pod.UID { // 加入克隆的节点 nodeInfoOut.AddPodInfo(pi) // 通知preFilter扩展同步状态 status := fh.RunPreFilterExtensionAddPod(ctx, stateOut, pod, pi, nodeInfoOut) // 失败提前结束 if !status.IsSuccess() { return false, state, nodeInfo, status.AsError() } podsAdded = true } } return podsAdded, stateOut, nodeInfoOut, nil } // RunFilterPlugins runs the set of configured Filter plugins for pod on the given node. func (f *frameworkImpl) RunFilterPlugins(...) *framework.Status { // 遍历filter插件 for _, pl := range f.filterPlugins { // 状态里标记跳过 if state.SkipFilterPlugins.Has(pl.Name()) { continue } // 执行Filter插件 if status := f.runFilterPlugin(ctx, pl, state, pod, nodeInfo); !status.IsSuccess() { ... return status } } 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注意
runFilterPluginsWithNominatedPods()会基于两轮检查筛选node,第一轮假设已提名高优先级pod加入node检查能否承载当前pod,第二轮假设这些pod未加入node能否承载当前pod,两次检查均通过才会采取当前node,这是一种很保守的策略
# 3.优选
# 3.1.prioritizeNodes
prioritizeNodes()是调度器优选阶段的实现,内部遍历调用framework的preScore插件及score插件对各node打分,以推举最优节点。// prioritizeNodes prioritizes the nodes by running the score plugins, // which return a score for each node from the call to RunScorePlugins(). func prioritizeNodes(...) ([]framework.NodePluginScores, error) { // extender为空及score插件为空 if len(extenders) == 0 && !fwk.HasScorePlugins() { ... // 给出默认分值 for i := range nodes { result = append(result, framework.NodePluginScores{ Name: nodes[i].Name, TotalScore: 1, }) } return result, nil } // 执行preScore插件 preScoreStatus := fwk.RunPreScorePlugins(ctx, state, pod, nodes) ... // 执行score插件 nodesScores, scoreStatus := fwk.RunScorePlugins(ctx, state, pod, nodes) ... // extender扩展打分 if len(extenders) != 0 && nodes != nil { ... // 遍历extender for i := range extenders { // extender不关心当前pod类型 if !extenders[i].IsInterested(pod) { continue } wg.Add(1) go func(extIndex int) { ... // http调用extender扩展 prioritizedList, weight, err := extenders[extIndex].Prioritize(pod, nodes) ... // 分值合并 for i := range *prioritizedList { nodename := (*prioritizedList)[i].Host score := (*prioritizedList)[i].Score ... // 缩放到调度器的分值范围 finalscore := score * weight * (framework.MaxNodeScore / extenderv1.MaxExtenderPriority) ... // 记录各extender对节点的打分 allNodeExtendersScores[nodename].Scores = append(allNodeExtendersScores[nodename].Scores, framework.PluginScore{ Name: extenders[extIndex].Name(), Score: finalscore, }) // 记录各extender对节点的打分总和 allNodeExtendersScores[nodename].TotalScore += finalscore } }(i) } // wait for all go routines to finish wg.Wait() // 合并结果 for i := range nodesScores { if score, ok := allNodeExtendersScores[nodes[i].Name]; ok { nodesScores[i].Scores = append(nodesScores[i].Scores, score.Scores...) nodesScores[i].TotalScore += score.TotalScore } } } ... return nodesScores, 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注意
prioritizeNodes()会依次执行preScore、score及extender进行节点打分检查及打分,根据nodeName组织各插件打分结果
# 3.2.runScorePlugins
runScorePlugins()是调度器的核心打分引擎,负责执行所有注册的score插件以向各节点打分,作为后续优选的评分来源。// RunScorePlugins runs the set of configured scoring plugins. // It returns a list that stores scores from each plugin and total score for each Node. func (f *frameworkImpl) RunScorePlugins(...) (ns []framework.NodePluginScores, status *framework.Status) { ... // pluginToNodeScores = { // "PodAffinity": [ {node1, 30}, {node2, 10}, ... ], // "ResourceFit": [ {node1, 60}, {node2, 80}, ... ], // } if len(plugins) > 0 { // 并行执行各插件的score方法 f.Parallelizer().Until(ctx, len(nodes), func(index int) { nodeName := nodes[index].Name for _, pl := range plugins { s, status := f.runScorePlugin(ctx, pl, state, pod, nodeName) ... pluginToNodeScores[pl.Name()][index] = framework.NodeScore{ Name: nodeName, Score: s, } } }, metrics.Score) ... } // 归一化打分 f.Parallelizer().Until(ctx, len(plugins), func(index int) { pl := plugins[index] //未实现打分归一化 if pl.ScoreExtensions() == nil { return } nodeScoreList := pluginToNodeScores[pl.Name()] // 执行打分归一化 status := f.runScoreExtension(ctx, pl, state, pod, nodeScoreList) ... }, metrics.Score) ... // 各插件的分数混入配置的权重 f.Parallelizer().Until(ctx, len(nodes), func(index int) { ... // 根据插件及节点索引获取分数及权重 for i, pl := range plugins { weight := f.scorePluginWeight[pl.Name()] nodeScoreList := pluginToNodeScores[pl.Name()] score := nodeScoreList[index].Score ... // 计算权重分数 weightedScore := score * int64(weight) // 记录节点关联当前插件的分数 nodePluginScores.Scores[i] = framework.PluginScore{ Name: pl.Name(), Score: weightedScore, } // 记录总分数 nodePluginScores.TotalScore += weightedScore } // 记录节点对应插件分数 allNodePluginScores[index] = nodePluginScores }, metrics.Score) ... return allNodePluginScores, 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注意
runScorePlugins()分为打分、归一化及权重调整几个阶段,最终计算出各score插件对node的分值
# 3.3.selectHost
selectHost()基于打分的nodes列表优选分值最高的node,分值相同的两个node会随机选择以实现不同node负载的再均衡。// selectHost takes a prioritized list of nodes and then picks one // in a reservoir sampling manner from the nodes that had the highest score. func selectHost(nodeScores []framework.NodePluginScores) (string, error) { ... maxScore := nodeScores[0].TotalScore selected := nodeScores[0].Name cntOfMaxScore := 1 // 顺序遍历nodeScore列表 for _, ns := range nodeScores[1:] { // 找到分数更大的更新选中node if ns.TotalScore > maxScore { maxScore = ns.TotalScore selected = ns.Name cntOfMaxScore = 1 // 分值相同的节点随机选择 } else if ns.TotalScore == maxScore { cntOfMaxScore++ if rand.Intn(cntOfMaxScore) == 0 { // Replace the candidate with probability of 1/cntOfMaxScore selected = ns.Name } } } return selected, 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注意
selectHost()会根据分值获取运行pod的最优node