qcontroller
# 1.qcontroller
# 1.1.initialize
queue controller监听Queue/PodGroup/Command资源对象,基于对象状态更新Queue资源,后续根据Queue资源触发调度。// Initialize creates QueueController from option. func (c *queuecontroller) Initialize(opt *framework.ControllerOption) error { ... // queue watcher queueInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: c.addQueue, UpdateFunc: c.updateQueue, DeleteFunc: c.deleteQueue, }) // podgroup watcher pgInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: c.addPodGroup, UpdateFunc: c.updatePodGroup, DeleteFunc: c.deletePodGroup, }) // enable commandSync if utilfeature.DefaultFeatureGate.Enabled(features.QueueCommandSync) { c.cmdInformer = factory.Bus().V1alpha1().Commands() c.cmdInformer.Informer().AddEventHandler(cache.FilteringResourceEventHandler{ FilterFunc: func(obj interface{}) bool { switch v := obj.(type) { // command ref queue case *busv1alpha1.Command: return IsQueueReference(v.TargetObject) default: return false } }, Handler: cache.ResourceEventHandlerFuncs{ AddFunc: c.addCommand, }, }) c.cmdLister = c.cmdInformer.Lister() c.cmdSynced = c.cmdInformer.Informer().HasSynced } // global handleHook queuestate.SyncQueue = c.syncQueue queuestate.OpenQueue = c.openQueue queuestate.CloseQueue = c.closeQueue // queue controller handleHook c.syncHandler = c.handleQueue c.syncCommandHandler = c.handleCommand c.enqueueQueue = c.enqueue 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
注意
queue controller注册eventHandler,还会初始化一些handleHook用于后续逻辑处理
# 1.2.vcstart
controller.Run()会激活queueInformer/pgInformer/cmdInformer,同步完成后异步执行worker和commandworker协调处理队列项。// Run starts QueueController. func (c *queuecontroller) Run(stopCh <-chan struct{}) { ... // wait cached c.vcInformerFactory.Start(stopCh) for informerType, ok := range c.vcInformerFactory.WaitForCacheSync(stopCh) { if !ok { return } } // default 5 worker for i := 0; i < int(c.workers); i++ { go wait.Until(c.worker, 0, stopCh) go wait.Until(c.commandWorker, 0, stopCh) } <-stopCh } func (c *queuecontroller) commandWorker() { for c.processNextCommand() { } } func (c *queuecontroller) processNextCommand() bool { cmd, shutdown := c.commandQueue.Get() ... defer c.commandQueue.Done(cmd) // process and handle err c.handleCommandErr(c.syncCommandHandler(cmd), cmd) return true } func (c *queuecontroller) handleCommand(cmd *busv1alpha1.Command) error { ... // process and delete command c.vcClient.BusV1alpha1().Commands(cmd.Namespace).Delete(context.TODO(), cmd.Name, metav1.DeleteOptions{}) ... req := &apis.Request{ QueueName: cmd.TargetObject.Name, Event: busv1alpha1.CommandIssuedEvent, Action: busv1alpha1.Action(cmd.Action), } c.enqueueQueue(req) return nil } func (c *queuecontroller) handleCommandErr(err error, cmd *busv1alpha1.Command) { if err == nil { c.commandQueue.Forget(cmd) return } // 还能重试 if c.maxRequeueNum == -1 || c.commandQueue.NumRequeues(cmd) < c.maxRequeueNum { c.commandQueue.AddRateLimited(cmd) return } c.commandQueue.Forget(cmd) }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
注意
command对象是一次性的,入队处理后会转为request推到queue,相应的command对象会删除
# 1.3.worker
controller.worker()持续消费队列,基于queue对象状态初始化state执行queueState.Execute(),队列数据由queueInformer产生。// worker runs a worker thread that just dequeues items, processes them, and marks them done. func (c *queuecontroller) worker() { for c.processNextWorkItem() { } } func (c *queuecontroller) processNextWorkItem() bool { req, shutdown := c.queue.Get() ... defer c.queue.Done(req) c.handleQueueErr(c.syncHandler(req), req) return true } func (c *queuecontroller) handleQueue(req *apis.Request) error { ... queue, err := c.queueLister.Get(req.QueueName) ... // 初始化queueState(基于queue.status.state) queueState := queuestate.NewState(queue) ... // 执行stateHook // openState.Execute // closedState.Execute // closingState.Execute // unknownState.Execute queueState.Execute(req.Action) return nil } func (c *queuecontroller) handleQueueErr(err error, req *apis.Request) { if err == nil { c.queue.Forget(req) return } if c.maxRequeueNum == -1 || c.queue.NumRequeues(req) < c.maxRequeueNum { c.queue.AddRateLimited(req) return } c.queue.Forget(req) }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
注意
queuestate.NewState会基于queue状态实例化不同state对象,相应触发不同Execute实现
# 2.state
# 2.1.openState
openState意为Queue是新建的或打开的,匹配预期状态或状态未知会执行SyncQueue同步,未匹配执行OpenQueue/CloseQueue更新状态。func (os *openState) Execute(action v1alpha1.Action) error { switch action { case v1alpha1.OpenQueueAction: // openAction,进行syncQueue调和 return SyncQueue(os.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { status.State = v1beta1.QueueStateOpen }) case v1alpha1.CloseQueueAction: // closeAction,进行queue收尾 return CloseQueue(os.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing }) default: // 重新同步 return SyncQueue(os.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { specState := os.queue.Status.State if len(specState) == 0 || specState == v1beta1.QueueStateOpen { status.State = v1beta1.QueueStateOpen return } if specState == v1beta1.QueueStateClosed { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing return } status.State = v1beta1.QueueStateUnknown }) } }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
# 2.2.closingState
closingState代表Queue正在关闭,OpenQueueAction会重新打开Queue,否则基于queue.status和关联的PodGroup重置状态。func (cs *closingState) Execute(action v1alpha1.Action) error { switch action { // 重新打开 case v1alpha1.OpenQueueAction: return OpenQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { status.State = v1beta1.QueueStateOpen }) // 关闭Queue case v1alpha1.CloseQueueAction: return SyncQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing }) // 重新同步 default: return SyncQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { specState := cs.queue.Status.State if specState == v1beta1.QueueStateOpen { status.State = v1beta1.QueueStateOpen return } if specState == v1beta1.QueueStateClosing { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing return } status.State = v1beta1.QueueStateUnknown }) } }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
# 2.3.closedState
closedState基于OpenQueueAction重新打开Queue,CloseQueueAction关闭Queue,未知状态基于queue status重置状态。func (cs *closedState) Execute(action v1alpha1.Action) error { switch action { // 打开Queue case v1alpha1.OpenQueueAction: return OpenQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { status.State = v1beta1.QueueStateOpen }) // 关闭Queue case v1alpha1.CloseQueueAction: return SyncQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { status.State = v1beta1.QueueStateClosed }) // 重新同步 default: return SyncQueue(cs.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { specState := cs.queue.Status.State if specState == v1beta1.QueueStateOpen { status.State = v1beta1.QueueStateOpen return } if specState == v1beta1.QueueStateClosed { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateUnknown }) } }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
# 2.4.unknownState
unknownState基于OpenQueueAction重新打开Queue,CloseQueueAction关闭Queue,未知状态基于status和PodGroup重置状态。func (us *unknownState) Execute(action v1alpha1.Action) error { switch action { // 打开Queue case v1alpha1.OpenQueueAction: return OpenQueue(us.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { status.State = v1beta1.QueueStateOpen }) // 关闭Queue case v1alpha1.CloseQueueAction: return CloseQueue(us.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing }) // 重新同步 default: return SyncQueue(us.queue, func(status *v1beta1.QueueStatus, podGroupList []string) { specState := us.queue.Status.State if specState == v1beta1.QueueStateOpen { status.State = v1beta1.QueueStateOpen return } if specState == v1beta1.QueueStateClosed { if len(podGroupList) == 0 { status.State = v1beta1.QueueStateClosed return } status.State = v1beta1.QueueStateClosing return } status.State = v1beta1.QueueStateUnknown }) } }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
# 3.action
# 3.1.syncQueue
c.syncQueue会尝试补充queue parent,根据关联的podgroup统计计数器,基于updateStatusHook重置及更新queue.status.state。func (c *queuecontroller) syncQueue(queue *schedulingv1beta1.Queue, stateFn state.UpdateQueueStatusFn) error { // add parent queue if parent not specified(root) queue, err := c.updateQueueParent(queue) ... // 关联的podgroups podGroups := c.getPodGroups(queue.Name) ... // 重置queue状态 if stateFn != nil { stateFn(&queueStatus, podGroups) } else { queueStatus.State = queue.Status.State } newQueue := queue.DeepCopy() // ignore update when state does not change if queueStatus.State != queue.Status.State { ... // update queue state newQueue = c.vcClient.SchedulingV1beta1().Queues().ApplyStatus(context.TODO(), queueApply, ...) ... } // sync the state between parent and child queues return c.syncHierarchicalQueue(newQueue) } // sync the state between parent and child queues func (c *queuecontroller) syncHierarchicalQueue(queue *schedulingv1beta1.Queue) error { if queue.Name == "root" { return nil } // 获取parent queue parentQueue, err := c.queueLister.Get(queue.Spec.Parent) ... switch parentQueue.Status.State { case schedulingv1beta1.QueueStateClosed, schedulingv1beta1.QueueStateClosing: // open queue is updated when parent queue is in the closed/closing state. if queue.Status.State != QueueStateClosed && queue.Status.State != QueueStateClosing { // volcano.sh/closed-by-parent=true c.updateQueueAnnotation(queue, ClosedByParentAnnotationKey, ClosedByParentAnnotationTrueValue) ... req := &apis.Request{ QueueName: queue.Name, Action: busv1alpha1.CloseQueueAction, } c.enqueue(req) } case schedulingv1beta1.QueueStateOpen: if queue.Status.State == QueueStateClosed || queue.Status.State == QueueStateClosing { // close queue is updated when parent queue is open state and queue close by parent. if queue.Annotations[ClosedByParentAnnotationKey] == ClosedByParentAnnotationTrueValue { req := &apis.Request{ QueueName: queue.Name, Action: busv1alpha1.OpenQueueAction, } c.enqueue(req) } } } 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
注意
syncQueue主要根据updateStateFn更新queue.state.state状态,同时基于queue.parent修正queue.status.state
# 3.2.openQueue
c.queuecontroller()用于重置queue.status.state=open,重置期间会检查queue parent状态及触发child queue入队协调。func (c *queuecontroller) openQueue(queue *schedulingv1beta1.Queue, stateFn state.UpdateQueueStatusFn) error { // queue not open if queue.Status.State != schedulingv1beta1.QueueStateOpen { c.openHierarchicalQueue(queue) ... } newQueue := queue.DeepCopy() if stateFn != nil { stateFn(&newQueue.Status, nil) } // update queue state if queue.Status.State != newQueue.Status.State { ... c.vcClient.SchedulingV1beta1().Queues().ApplyStatus(context.TODO(), queueApply, ...) } // update annotation with queue parent(volcano.sh/closed-by-parent=false) _, err := c.updateQueueAnnotation(queue, ClosedByParentAnnotationKey, ClosedByParentAnnotationFalseValue) return err } func (c *queuecontroller) openHierarchicalQueue(queue *schedulingv1beta1.Queue) error { if queue.Spec.Parent != "" && queue.Spec.Parent != "root" { parentQueue, err := c.queueLister.Get(queue.Spec.Parent) ... // parent queue close if parentQueue.Status.State == QueueStateClosing || parentQueue.Status.State == QueueStateClosed { // the parent queue may be being opened, and it may take a few attempts to open the child queue. return fmt.Errorf("Failed to open queue %s because its parent queue %s is closing or closed.") } } // queues queueList, err := c.queueLister.List(labels.Everything()) ... // update child queue for _, childQueue := range queueList { // child queue closed by queue if childQueue.Spec.Parent == queue.Name&&childQueue.Annotations[ClosedByParentAnnotationKey] == "true" { // open child queue req := &apis.Request{ QueueName: childQueue.Name, Action: busv1alpha1.OpenQueueAction, } c.enqueue(req) } } 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注意
openQueue会尝试修改queue.status.state为Open,同时尝试open child Queue
# 3.3.closeQueue
c.closeQueue()先基于CloseQueueAction入队child queue触发child close,再尝试更新queue state及更新到queue对象。func (c *queuecontroller) closeQueue(queue *schedulingv1beta1.Queue, stateFn state.UpdateQueueStatusFn) error { // update child queue if queue.Status.State != QueueStateClosed && queue.Status.State != QueueStateClosing { continued, err := c.closeHierarchicalQueue(queue) if !continued { return err } } podGroups := c.getPodGroups(queue.Name) newQueue := queue.DeepCopy() if updateStateFn != nil { updateStateFn(&newQueue.Status, podGroups) } // update queue state if queue.Status.State != newQueue.Status.State { c.vcClient.SchedulingV1beta1().Queues().ApplyStatus(context.TODO(), queueApply, ...) ... } return nil } func (c *queuecontroller) closeHierarchicalQueue(queue *schedulingv1beta1.Queue) (bool, error) { if queue.Name == "root" { return false, nil } queueList, err := c.queueLister.List(labels.Everything()) ... for _, childQueue := range queueList { if childQueue.Spec.Parent != queue.Name { continue } // child queue not close if childQueue.Status.State != QueueStateClosed && childQueue.Status.State != QueueStateClosing { c.updateQueueAnnotation(childQueue, ClosedByParentAnnotationKey, true) ... req := &apis.Request{ QueueName: childQueue.Name, Action: busv1alpha1.CloseQueueAction, } c.enqueue(req) } } return true, 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注意
closeQueue会尝试修改queue.status.state为closed/closing,必要时触发close child Queue