schedqueue
# 1.简介
# 1.1.概述
schedulingQueue负责待调度pod的存储,,内部实现依赖三个队列:activeQ、unscheduleQ及backoffQ,以维护pod调度顺序。--- activeQ(heap结构) 基于优先级存放待调度pod,调度器会获取队列中的pod进行调度,调度失败的pod加入unscheduleQ,成功的pod会移除 --- unscheduleQ 存放调度失败的pod --- backoffQ(heap结构) 基于重试时间按序存放待重试的pod,重试算法利用指数退避机制,重试时间为[1s,10s]1
2
3
4
5
6
7
8注意
unscheduleQ和backoffQ暂存待调度pod以限流及延迟,activeQ作为scheduler调度数据来源入口
# 1.2.队列数据
调度队列的
Pod结构是QueuedPodInfo,由PodInfo加上部分队列属性组成,队列属性用于退避重试,以限制重试流量。// QueuedPodInfo is a Pod wrapper with additional information related to // the pod's status in the scheduling queue, such as the timestamp when // it's added to the queue. type QueuedPodInfo struct { *PodInfo // 入队时间戳 Timestamp time.Time // 失败次数 Attempts int // 首次入队时间 InitialAttemptTimestamp time.Time // 上一次造成调度失败的插件名称集合(NodeResourcesFit/PodAffinity/TaintToleration) UnschedulablePlugins sets.String } // PodInfo is a wrapper to a Pod with additional pre-computed information to // accelerate processing. type PodInfo struct { Pod *v1.Pod // 硬性亲和性(过滤阶段) RequiredAffinityTerms []AffinityTerm // 硬性反亲和性(过滤阶段) RequiredAntiAffinityTerms []AffinityTerm // 软性亲和性(打分阶段) PreferredAffinityTerms []WeightedAffinityTerm // 软性反亲和性(打分阶段) PreferredAntiAffinityTerms []WeightedAffinityTerm ParseError error }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注意
queuedPodInfo基于退避调度考量对pod进行再包装,以限制重试频率
# 1.3.调度队列
scheduler的调度队列schedulingQueue是行为限制接口,定义入队、更新及出队方法,以维护及暂存待调度pod供调度器使用。// SchedulingQueue is an interface for a queue to store pods waiting to be scheduled. type SchedulingQueue interface { // pod状态同步器,更新提名器信息 framework.PodNominator // 注册待调度pod Add(pod *v1.Pod) error // 注册pod到activeQ Activate(pods map[string]*v1.Pod) // 将无法调度pod注册到unscheduleQ AddUnschedulableIfNotPresent(pod *framework.QueuedPodInfo, podSchedulingCycle int64) error // 调度周期,pod对应一次周期 SchedulingCycle() int64 // 队头弹出pod Pop() (*framework.QueuedPodInfo, error) // 更新pod Update(oldPod, newPod *v1.Pod) error // 删除pod Delete(pod *v1.Pod) error // unscheduleQ的pod移到activeQ或backoffQ MoveAllToActiveOrBackoffQueue(event framework.ClusterEvent, preCheck PreEnqueueCheck) // 关联pod被添加 AssignedPodAdded(pod *v1.Pod) // 关联pod被更新 AssignedPodUpdated(pod *v1.Pod) // 获取队列中的pod PendingPods() []*v1.Pod // 关闭调度队列 Close() // 启动调度队列 Run() }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
注意
schedulingQueue是调度队列接口,具体实现由PriorityQueue对象完成
# 1.4.优先队列
priorityQueue实现schedulingQueue提供pod调度数据维护,加入priorityQueue的pod会进一步包装,注入部分管理属性。type PriorityQueue struct { // 提名器 *nominator stop chan struct{} clock util.Clock // 初始退避时长(1s) podInitialBackoffDuration time.Duration // 最大退避时长(10s) podMaxBackoffDuration time.Duration // pod处于unscheduleQ最大时长 podMaxInUnschedulablePodsDuration time.Duration cond sync.Cond // 可立即调度队列 activeQ *heap.Heap // 退避调度队列 podBackoffQ *heap.Heap // 无法调度队列 unschedulablePods *UnschedulablePods // 调度周期编号 schedulingCycle int64 // 收到移动请求周期编号 moveRequestCycle int64 clusterEventMap map[framework.ClusterEvent]sets.String closed bool // namespace缓存器 nsLister listersv1.NamespaceLister } // Heap is a producer/consumer queue that implements a heap data structure. // It can be used to implement priority queues and similar data structures. type Heap struct { // data stores objects and has a queue that keeps their ordering according // to the heap invariant. data *data ... } // data is an internal struct that implements the standard heap interface and keeps the data stored in the heap. type data struct { // key-->pod映射 items map[string]*heapItem // key顺序 queue []string // pod唯一标识生成器(name/namespace) keyFunc KeyFunc // 比较器(优先级/退避时间) lessFunc lessFunc } // UnschedulablePods holds pods that cannot be scheduled. type UnschedulablePods struct { // key-->pod映射 podInfoMap map[string]*framework.QueuedPodInfo // pod唯一标识生成器 keyFunc func(*v1.Pod) string ... }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注意
priorityQueue的核心部分就是activeQ、backoffQ及unscheduleQ
# 2.源码
# 2.1.add
p.add()用于将pod转换为queuedPodInfo加入activeQ队列,后续scheduler会获取activeQ的pod进行调度。// Add adds a pod to the active queue. It should be called only when a new pod is added. func (p *PriorityQueue) Add(pod *v1.Pod) error { p.lock.Lock() defer p.lock.Unlock() pInfo := p.newQueuedPodInfo(pod) // 加入activeQ p.activeQ.Add(pInfo) ... // 由unscheduleQ删除 if p.unschedulablePods.get(pod) != nil { p.unschedulablePods.delete(pod) } // 由backoffQ删除 p.podBackoffQ.Delete(pInfo) ... // pod与提名node映射 p.addNominatedPodUnlocked(pInfo.PodInfo, nil) // 广播唤醒调度器 p.cond.Broadcast() return nil } func (npm *nominator) addNominatedPodUnlocked(pi *framework.PodInfo, nominatingInfo *framework.NominatingInfo) { // 清理旧记录(pod-->node,node-->pod) npm.delete(pi.Pod) var nodeName string // 沿用pod已提名node if nominatingInfo.Mode() == framework.ModeOverride { nodeName = nominatingInfo.NominatedNodeName // 沿用pod关联的提名node } else if nominatingInfo.Mode() == framework.ModeNoop { if pi.Pod.Status.NominatedNodeName == "" { return } nodeName = pi.Pod.Status.NominatedNodeName } // 幽灵pod检查 if npm.podLister != nil { // 确认pod存在 updatedPod, err := npm.podLister.Pods(pi.Pod.Namespace).Get(pi.Pod.Name) ... // pod已调度 if updatedPod.Spec.NodeName != "" { return } } // 记录pod调度到的node npm.nominatedPodToNode[pi.Pod.UID] = nodeName // node已关联pod则不重新记录 for _, npi := range npm.nominatedPods[nodeName] { if npi.Pod.UID == pi.Pod.UID { return } } // 记录node允许容纳的pod npm.nominatedPods[nodeName] = append(npm.nominatedPods[nodeName], pi) }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注意
注册
pod至activeQ会清理backoffQ和unscheduleQ数据避免重复调度,同时通知调度器触发流程
# 2.2.update
pod属性发生变化会同步到schedulingQueue,依次更新activeQ、backoffQ及unscheduleQ的pod信息供scheduler调度。// Update updates a pod in the active or backoff queue if present. func (p *PriorityQueue) Update(oldPod, newPod *v1.Pod) error { p.lock.Lock() defer p.lock.Unlock() // activeQ和backoffQ依赖管理数据,已加入的pod才能触发更新 if oldPod != nil { oldPodInfo := newQueuedPodInfoForLookup(oldPod) // 更新activeQ队列的pod if oldPodInfo, exists, _ := p.activeQ.Get(oldPodInfo); exists { // pod元数据 pInfo := updatePod(oldPodInfo, newPod) // 调整node提名 p.updateNominatedPodUnlocked(oldPod, pInfo.PodInfo) // 更新 return p.activeQ.Update(pInfo) } // 更新backoffQ队列的pod if oldPodInfo, exists, _ := p.podBackoffQ.Get(oldPodInfo); exists { // pod元数据 pInfo := updatePod(oldPodInfo, newPod) // 调整node提名 p.updateNominatedPodUnlocked(oldPod, pInfo.PodInfo) // 更新 return p.podBackoffQ.Update(pInfo) } } // unscheduleQ仅缓存无法调度pod,不依赖管理数据,可直接基于pod唯一标识更新 if usPodInfo := p.unschedulablePods.get(newPod); usPodInfo != nil { // pod元数据 pInfo := updatePod(usPodInfo, newPod) // 调整node提名 p.updateNominatedPodUnlocked(oldPod, pInfo.PodInfo) // pod更新 if isPodUpdated(oldPod, newPod) { // 退避周期内 if p.isPodBackingoff(usPodInfo) { // 加入backoffQ p.podBackoffQ.Add(pInfo) ... // 由unscheduleQ删除 p.unschedulablePods.delete(usPodInfo.Pod) // 退避结束 } else { // 加入activeQ p.activeQ.Add(pInfo) ... // 由unscheduleQ删除 p.unschedulablePods.delete(usPodInfo.Pod) // 通知调度器有可调度pod p.cond.Broadcast() } // pod未更新 } else { // 仅更新unscheduleQ的pod数据,由异步循环周期加入其它队列 p.unschedulablePods.addOrUpdate(pInfo) } return nil } // 未加入任何队列pod放到activeQ pInfo := p.newQueuedPodInfo(newPod) p.activeQ.Add(pInfo) ... // 注册node提名 p.addNominatedPodUnlocked(pInfo.PodInfo, nil) // 通知调度器 p.cond.Broadcast() return nil } // 调整node提名 func (npm *nominator) updateNominatedPodUnlocked(oldPod *v1.Pod, newPodInfo *framework.PodInfo) { ... // 检查沿用提名node情况(apiserver同步延迟) if NominatedNodeName(oldPod) == "" && NominatedNodeName(newPodInfo.Pod) == "" { // 内存维护提名节点 if nnn, ok := npm.nominatedPodToNode[oldPod.UID]; ok { // 沿用内存已提名节点 nominatingInfo = &framework.NominatingInfo{ NominatingMode: framework.ModeOverride, NominatedNodeName: nnn, } } } // 清理旧pod提名节点 npm.delete(oldPod) // 重新注册提名节点 npm.addNominatedPodUnlocked(newPodInfo, nominatingInfo) }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注意
activeQ和backoffQ均维护类似优先级、延迟时间的管理数据,要求已加入队列的pod才触发更新,unscheduleQ仅作为无法调度暂存区
# 2.3.pop
q.pop()负责弹出activeQ队列的堆顶pop,activeQ没有可调度pod会执行cond.Wait()阻塞,其它协程向activeQ转移数据会通知唤醒。// Pop removes the head of the active queue and returns it. func (p *PriorityQueue) Pop() (*framework.QueuedPodInfo, error) { p.lock.Lock() defer p.lock.Unlock() // activeQ为空 for p.activeQ.Len() == 0 { // 队列关闭 if p.closed { return nil, fmt.Errorf(queueClosed) } // 阻塞待唤醒 p.cond.Wait() } // 获取activeQ堆顶pod obj, err := p.activeQ.Pop() ... pInfo := obj.(*framework.QueuedPodInfo) // 更新调度属性 pInfo.Attempts++ p.schedulingCycle++ return pInfo, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22注意
p.pop()会弹出activeQ堆顶pod,没有可调度pod会阻塞直到activeQ加入pod才被唤醒
# 2.4.delete
p.delete()用于将指定的pod由队列中删除,内部会依次尝试清理activeQ、backoffQ和unscheduleQ队列的数据,取消pod调度。// Delete deletes the item from either of the two queues. It assumes the pod is only in one queue. func (p *PriorityQueue) Delete(pod *v1.Pod) error { p.lock.Lock() defer p.lock.Unlock() // 清理node提名 p.deleteNominatedPodIfExistsUnlocked(pod) // 清理activeQ数据 if err := p.activeQ.Delete(newQueuedPodInfoForLookup(pod)); err != nil { //activeQ没找到,清理backoffQ数据 p.podBackoffQ.Delete(newQueuedPodInfoForLookup(pod)) // 清理unscheduleQ数据 p.unschedulablePods.delete(pod) } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15注意
p.delete()会依次清理activeQ、backoffQ和unscheduleQ队列pod,同时清理nominated提名的node
# 2.5.activate
p.activate()用于激活pod调度,将给定pod由backoffQ和unscheduleQ转移到activeQ队列,以重新将未调度的pod纳入流程。// Activate moves the given pods to activeQ iff they're in unschedulablePods or backoffQ. func (p *PriorityQueue) Activate(pods map[string]*v1.Pod) { p.lock.Lock() defer p.lock.Unlock() activated := false for _, pod := range pods { // 将pod加入activeQ if p.activate(pod) { activated = true } } // activeQ加入新pod,唤醒调度 if activated { p.cond.Broadcast() } } func (p *PriorityQueue) activate(pod *v1.Pod) bool { // activeQ已存在pod if _, exists, _ := p.activeQ.Get(newQueuedPodInfoForLookup(pod)); exists { return false } ... // 检查unscheduleQ队列 if pInfo = p.unschedulablePods.get(pod); pInfo == nil { // 检查backoffQ队列 if obj, exists, _ := p.podBackoffQ.Get(newQueuedPodInfoForLookup(pod)); !exists { return false } else { pInfo = obj.(*framework.QueuedPodInfo) } } // pod未加入backoffQ和unscheduleQ队列,无需激活 if pInfo == nil { return false } // pod加入过backoffQ和unscheduleQ队列,转移到activeQ队列 if err := p.activeQ.Add(pInfo); err != nil { return false } // 清理backoffQ和unscheduleQ队列数据 p.unschedulablePods.delete(pod) p.podBackoffQ.Delete(pInfo) ... // 注册pod提名node p.addNominatedPodUnlocked(pInfo.PodInfo, nil) return true }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注意
activate()仅用于将可调度pod由backoffQ和unscheduleQ转移到activeQ队列,不会更新队列已存在pod
# 2.6.addIfNotPresent
p.AddUnschedulableIfNotPresent()用于将调度失败的pod重新放回backoffQ或unscheduleQ队列,以待后续重新尝试调度。// AddUnschedulableIfNotPresent inserts a pod that cannot be scheduled into // the queue, unless it is already in the queue. func (p *PriorityQueue) AddUnschedulableIfNotPresent(...) error { p.lock.Lock() defer p.lock.Unlock() ... // 检查pod未加入任何队列 ... // 更新时间戳 pInfo.Timestamp = p.clock.Now() ... // 已经唤醒过unscheduleQ if p.moveRequestCycle >= podSchedulingCycle { // pod加入backoffQ p.podBackoffQ.Add(pInfo) ... // 发生调度后未唤醒过unscheduleQ } else { // pod加入或更新unscheduleQ p.unschedulablePods.addOrUpdate(pInfo) ... } // 重新注册提名node p.addNominatedPodUnlocked(pInfo.PodInfo, nil) 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注意
调度失败的
pod会根据唤醒触发情况加入backoffQ或unscheduleQ,以及重新触发提名节点注册
# 2.7.movePod
p.movePodsToActiveOrBackoffQueue()用于将一批不可调度pod转移到activeQ或backoffQ,以将搁置的pod重新投入调度循环。// MoveAllToActiveOrBackoffQueue moves all pods from unschedulablePods to activeQ or backoffQ. func (p *PriorityQueue) MoveAllToActiveOrBackoffQueue(event framework.ClusterEvent, preCheck PreEnqueueCheck) { p.lock.Lock() defer p.lock.Unlock() ... // 基于检查条件加入待调度pod for _, pInfo := range p.unschedulablePods.podInfoMap { if preCheck == nil || preCheck(pInfo.Pod) { unschedulablePods = append(unschedulablePods, pInfo) } } // 待调度pod加入activeQ或backoffQ p.movePodsToActiveOrBackoffQueue(unschedulablePods, event) } // AssignedPodAdded is called when a bound pod is added. func (p *PriorityQueue) AssignedPodAdded(pod *v1.Pod) { p.lock.Lock() // // 将满足当前绑定成功pod亲和的pod加入activeQ或backoffQ p.movePodsToActiveOrBackoffQueue(p.getUnschedulablePodsWithMatchingAffinityTerm(pod), AssignedPodAdd) p.lock.Unlock() } // AssignedPodUpdated is called when a bound pod is updated. func (p *PriorityQueue) AssignedPodUpdated(pod *v1.Pod) { p.lock.Lock() // 将满足当前绑定成功pod亲和的pod加入activeQ或backoffQ p.movePodsToActiveOrBackoffQueue(p.getUnschedulablePodsWithMatchingAffinityTerm(pod), AssignedPodUpdate) p.lock.Unlock() } // NOTE: this function assumes lock has been acquired in caller func (p *PriorityQueue) movePodsToActiveOrBackoffQueue(...) { ... // 遍历待调度pod for _, pInfo := range podInfoList { // 检查集群事件与调度失败的插件事件重合(有重合说明失败可能被解决,才有价值调度) if len(pInfo.UnschedulablePlugins) != 0 && !p.podMatchesEvent(pInfo, event) { continue } moved = true ... // pod处于退避期 if p.isPodBackingoff(pInfo) { // 加入backoffQ队列 p.podBackoffQ.Add(pInfo) ... // 由不可调度队列删除 p.unschedulablePods.delete(pod) // 未处于退避期 } else { // 加入activeQ队列 p.activeQ.Add(pInfo) ... // 由不可调度队列删除 p.unschedulablePods.delete(pod) } } // 记录调度周期(用于调度失败重入队检查) p.moveRequestCycle = p.schedulingCycle // 通知给调度器 if moved { p.cond.Broadcast() } }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注意
未调度
pod激活主要应对亲和性匹配和集群资源事件变更场景,将有价值触发调度的pod重新纳入流程
# 2.8.run
priorityQueue.Run()会启动两个协程不停扫描backoffQ和unscheduleQ,将未调度的pod重新加入activeQ,以重新纳入调度流程。// Run starts the goroutine to pump from podBackoffQ to activeQ func (p *PriorityQueue) Run() { // 间隔1s检查backoffQ go wait.Until(p.flushBackoffQCompleted, 1.0*time.Second, p.stop) // 间隔30s检查unscheduleQ go wait.Until(p.flushUnschedulablePodsLeftover, 30*time.Second, p.stop) } // flushBackoffQCompleted Moves all pods from backoffQ which have completed backoff in to activeQ func (p *PriorityQueue) flushBackoffQCompleted() { p.lock.Lock() defer p.lock.Unlock() broadcast := false for { // backoffQ堆顶pod rawPodInfo := p.podBackoffQ.Peek() if rawPodInfo == nil { break } pod := rawPodInfo.(*framework.QueuedPodInfo).Pod // 计算退避时间 boTime := p.getBackoffTime(rawPodInfo.(*framework.QueuedPodInfo)) if boTime.After(p.clock.Now()) { break } // 退避时间结束,弹出堆顶pod _, err := p.podBackoffQ.Pop() ... // 加入activeQ p.activeQ.Add(rawPodInfo) ... broadcast = true } // 通知调度器 if broadcast { p.cond.Broadcast() } } // flushUnschedulablePodsLeftover moves pods which stay in unschedulablePods // longer than podMaxInUnschedulablePodsDuration to backoffQ or activeQ. func (p *PriorityQueue) flushUnschedulablePodsLeftover() { p.lock.Lock() defer p.lock.Unlock() ... // 遍历unscheduleQ队列pod for _, pInfo := range p.unschedulablePods.podInfoMap { // pod加入unscheduleQ队列时间 lastScheduleTime := pInfo.Timestamp // 加入时间超出最大保留时间则记录 if currentTime.Sub(lastScheduleTime) > p.podMaxInUnschedulablePodsDuration { podsToMove = append(podsToMove, pInfo) } } if len(podsToMove) > 0 { // 开始转换队列,退避期转入backoffQ,非退避期转入activeQ p.movePodsToActiveOrBackoffQueue(podsToMove, UnschedulableTimeout) } }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注意
run()其实就是周期扫描backoffQ和unscheduleQ队列的pod,将有调度可能的pod重新纳入调度流程