vcontroller
# 1.pgcontroller
# 1.1.initialize
podgroup controller监听Pod/PodGroup/Replica资源对象,基于Pod/Replica生成或删除关联的PodGroup资源及入队Req。// Initialize create new Podgroup Controller. func (pg *pgcontroller) Initialize(opt *framework.ControllerOption) error { ... // podLister pg.podInformer = opt.SharedInformerFactory.Core().V1().Pods() pg.podLister = pg.podInformer.Lister() pg.podSynced = pg.podInformer.Informer().HasSynced pg.podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: pg.addPod, }) ... // pgLister pg.vcInformerFactory = factory pg.pgInformer = factory.Scheduling().V1beta1().PodGroups() pg.pgLister = pg.pgInformer.Lister() pg.pgSynced = pg.pgInformer.Informer().HasSynced // workLoad if utilfeature.DefaultFeatureGate.Enabled(features.WorkLoadSupport) { pg.rsInformer = pg.informerFactory.Apps().V1().ReplicaSets() pg.rsSynced = pg.rsInformer.Informer().HasSynced pg.rsInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: pg.addReplicaSet, UpdateFunc: pg.updateReplicaSet, }) } return nil } func (pg *pgcontroller) addPod(obj interface{}) { pod, ok := obj.(*v1.Pod) ... req := podRequest{ podName: pod.Name, podNamespace: pod.Namespace, } pg.queue.Add(req) } func (pg *pgcontroller) updateReplicaSet(oldObj, newObj interface{}) { pg.addReplicaSet(newObj) } func (pg *pgcontroller) addReplicaSet(obj interface{}) { rs, ok := obj.(*appsv1.ReplicaSet) ... if *rs.Spec.Replicas == 0 { pgName := batchv1alpha1.PodgroupNamePrefix + string(rs.UID) ... pg.vcClient.SchedulingV1beta1().PodGroups(rs.Namespace).Delete(context.TODO(), pgName, ...) ... } // In the rolling upgrade scenario, the addReplicasSet(replicas=0) event may be received before // the updateReplicaSet(replicas=1) event. In this event, need to create PodGroup for the pod. if *rs.Spec.Replicas > 0 { selector := metav1.LabelSelector{MatchLabels: rs.Spec.Selector.MatchLabels} podList, err := pg.kubeClient.CoreV1().Pods(rs.Namespace).List(context.TODO(), selector) ... if podList != nil && len(podList.Items) > 0 { pod := podList.Items[0] // scheduler not match if !slices.Contains(pg.schedulerNames, pod.Spec.SchedulerName) { return } // 创建podgroup&更新pod pggroup annotation pg.createNormalPodPGIfNotExist(&pod) ... } } }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
注意
initialize主要初始化informer及eventHandler,维护Pod/PG/Rs Lister缓存
# 1.2.pgstart
controller.Run()会激活podInformer/pgInformer/rsInformer,同步完成后监听资源变化及触发入队,异步执行worker协调处理队列。// Run start NewPodgroupController. func (pg *pgcontroller) Run(stopCh <-chan struct{}) { pg.informerFactory.Start(stopCh) pg.vcInformerFactory.Start(stopCh) for informerType, ok := range pg.informerFactory.WaitForCacheSync(stopCh) { if !ok { klog.Errorf("caches failed to sync: %v", informerType) } } for informerType, ok := range pg.vcInformerFactory.WaitForCacheSync(stopCh) { if !ok { klog.Errorf("caches failed to sync: %v", informerType) return } } // default 5 worker for i := 0; i < int(pg.workers); i++ { go wait.Until(pg.worker, 0, stopCh) } } func (pg *pgcontroller) worker() { for pg.processNextReq() { } } func (pg *pgcontroller) processNextReq() bool { req, shutdown := pg.queue.Get() ... defer pg.queue.Done(req) pod, err := pg.podLister.Pods(req.podNamespace).Get(req.podName) ... // not match scheduler if !slices.Contains(pg.schedulerNames, pod.Spec.SchedulerName) { return true } // 已关联podgroup if pod.Annotations != nil && pod.Annotations[scheduling.KubeGroupNameAnnotationKey] != "" { return true } // 尝试创建podgroup及更新pod annotation if err := pg.createNormalPodPGIfNotExist(pod); err != nil { pg.queue.AddRateLimited(req) return true } // If no error, forget it. pg.queue.Forget(req) 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
56
57
58
注意
pgcontroller较为简单,会根据pod/rs创建podgroup及更新相关pod pgannotation
# 2.tpcontroller
# 2.1.initialize
jobtemplate controller监听job/jobtemplate资源对象,基于eventHandler触发jobtemplate入队及初始化cacheLister。func (jt *jobtemplatecontroller) Initialize(opt *framework.ControllerOption) error { ... // job template informer jt.jobTemplateInformer = factory.Flow().V1alpha1().JobTemplates() jt.jobTemplateSynced = jt.jobTemplateInformer.Informer().HasSynced jt.jobTemplateLister = jt.jobTemplateInformer.Lister() jt.jobTemplateInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: jt.addJobTemplate, }) // job informer jt.jobInformer = factory.Batch().V1alpha1().Jobs() jt.jobSynced = jt.jobInformer.Informer().HasSynced jt.jobLister = jt.jobInformer.Lister() jt.jobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: jt.addJob, }) // default 15qps jt.maxRequeueNum = opt.MaxRequeueNum if jt.maxRequeueNum < 0 { jt.maxRequeueNum = -1 } ... jt.queue = workqueue.NewTypedRateLimitingQueue(DefaultTypedControllerRateLimiter[apis.FlowRequest]()) // handler jt.enqueueJobTemplate = jt.enqueue jt.syncHandler = jt.handleJobTemplate return nil } func (jt *jobtemplatecontroller) addJobTemplate(obj interface{}) { jobTemplate, ok := obj.(*v1alpha1.JobTemplate) ... req := apis.FlowRequest{ Namespace: jobTemplate.Namespace, JobTemplateName: jobTemplate.Name, } jt.enqueueJobTemplate(req) } func (jt *jobtemplatecontroller) addJob(obj interface{}) { job, ok := obj.(*batch.Job) ... if job.Labels[CreatedByJobTemplate] == "" { return } //Filter vcjobs created by JobFlow namespaceName := strings.Split(job.Labels[CreatedByJobTemplate], ".") if len(namespaceName) != CreateByJobTemplateValueNum { return } namespace, name := namespaceName[0], namespaceName[1] req := apis.FlowRequest{ Namespace: namespace, JobTemplateName: name, } jt.enqueueJobTemplate(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
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
注意
initialize会初始化eventHandler监听及入队相关的jobtemplate,同时构造cacheLister缓存job/jobtemplate
# 2.2.tpstart
controller.Run()会激活jobInformer/jobtemplateInformer,同步完成后监听资源变化及触发入队,异步执行worker协调处理队列。func (jt *jobtemplatecontroller) Run(stopCh <-chan struct{}) { defer jt.queue.ShutDown() jt.vcInformerFactory.Start(stopCh) for informerType, ok := range jt.vcInformerFactory.WaitForCacheSync(stopCh) { if !ok { return } } // 间隔1s执行一次 go wait.Until(jt.worker, time.Second, stopCh) <-stopCh } func (jt *jobtemplatecontroller) worker() { for jt.processNextWorkItem() { } } func (jt *jobtemplatecontroller) processNextWorkItem() bool { req, shutdown := jt.queue.Get() ... defer jt.queue.Done(req) err := jt.syncHandler(&req) jt.handleJobTemplateErr(err, req) return true } func (jt *jobtemplatecontroller) handleJobTemplate(req *apis.FlowRequest) error { jobTemplate, err := jt.jobTemplateLister.JobTemplates(req.Namespace).Get(req.JobTemplateName) ... // 同步 jt.syncJobTemplate(jobTemplate) ... return nil } func (jt *jobtemplatecontroller) syncJobTemplate(jobTemplate *v1alpha1flow.JobTemplate) error { // search the jobs created by JobTemplate selector := labels.NewSelector() r, err := labels.NewRequirement(CreatedByJobTemplate, selection.Equals, GetTemplateKey(jobTemplate)) ... // 匹配template创建的job selector = selector.Add(*r) jobList, err := jt.jobLister.Jobs(jobTemplate.Namespace).List(selector) ... if len(jobList) == 0 { return nil } jobListName := make([]string, 0) for _, job := range jobList { jobListName = append(jobListName, job.Name) } jobTemplate.Status.JobDependsOnList = jobListName //update jobTemplate status jt.vcClient.JobTemplates(jobTemplate.Namespace).UpdateStatus(context.Background(), jobTemplate, ...) ... 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
注意
jobtemplate协调主要获取基于模板创建的job更新到jobtemplate.status
# 3.flowcontroller
# 3.1.initialize
jobflow controller监听job/jobtemplate/jobFlow资源对象,基于eventHandler触发jobflow入队及cacheLister初始化。func (jf *jobflowcontroller) Initialize(opt *framework.ControllerOption) error { ... // jobflow informer jf.jobFlowInformer = factory.Flow().V1alpha1().JobFlows() jf.jobFlowSynced = jf.jobFlowInformer.Informer().HasSynced jf.jobFlowLister = jf.jobFlowInformer.Lister() jf.jobFlowInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: jf.addJobFlow, UpdateFunc: jf.updateJobFlow, }) // jobtemplate informer jf.jobTemplateInformer = factory.Flow().V1alpha1().JobTemplates() jf.jobTemplateSynced = jf.jobTemplateInformer.Informer().HasSynced jf.jobTemplateLister = jf.jobTemplateInformer.Lister() // job informer jf.jobInformer = factory.Batch().V1alpha1().Jobs() jf.jobSynced = jf.jobInformer.Informer().HasSynced jf.jobLister = jf.jobInformer.Lister() jf.jobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ UpdateFunc: jf.updateJob, }) ... jf.queue = workqueue.NewTypedRateLimitingQueue(DefaultTypedControllerRateLimiter[apis.FlowRequest]()) // handler jf.enqueueJobFlow = jf.enqueue jf.syncHandler = jf.handleJobFlow jobflowstate.SyncJobFlow = jf.syncJobFlow return nil } func (jf *jobflowcontroller) addJobFlow(obj interface{}) { jobFlow, ok := obj.(*jobflowv1alpha1.JobFlow) ... // use struct instead of pointer req := apis.FlowRequest{ Namespace: jobFlow.Namespace, JobFlowName: jobFlow.Name, Action: jobflowv1alpha1.SyncJobFlowAction, Event: jobflowv1alpha1.OutOfSyncEvent, } jf.enqueueJobFlow(req) } func (jf *jobflowcontroller) updateJobFlow(oldObj, newObj interface{}) { oldJobFlow, ok := oldObj.(*jobflowv1alpha1.JobFlow) ... newJobFlow, ok := newObj.(*jobflowv1alpha1.JobFlow) ... if newJobFlow.ResourceVersion == oldJobFlow.ResourceVersion { return } // jobFlow跑完且保留策略为delete才继续删 if newJobFlow.Status.State.Phase != Succeed || newJobFlow.Spec.JobRetainPolicy != Delete { return } req := apis.FlowRequest{ Namespace: newJobFlow.Namespace, JobFlowName: newJobFlow.Name, Action: jobflowv1alpha1.SyncJobFlowAction, Event: jobflowv1alpha1.OutOfSyncEvent, } jf.enqueueJobFlow(req) } func (jf *jobflowcontroller) updateJob(oldObj, newObj interface{}) { oldJob, ok := oldObj.(*batch.Job) ... newJob, ok := newObj.(*batch.Job) ... // Filter out jobs that are not created from volcano jobflow if !isControlledBy(newJob, helpers.JobFlowKind) { return } if newJob.ResourceVersion == oldJob.ResourceVersion { return } // 由job ownerRef获取jobFlow名称 jobFlowName := getJobFlowNameByJob(newJob) if jobFlowName == "" { return } req := apis.FlowRequest{ Namespace: newJob.Namespace, JobFlowName: jobFlowName, Action: jobflowv1alpha1.SyncJobFlowAction, Event: jobflowv1alpha1.OutOfSyncEvent, } jf.enqueueJobFlow(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
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
注意
JobFlow配合JobTemplate使用,用于vcjob任务的编排
# 3.2.flowstart
controller.Run()用于启动informer,初始化cacheLister及执行worker处理队列的apis.FlowRequest对象,基于状态加载执行器。func (jf *jobflowcontroller) Run(stopCh <-chan struct{}) { defer jf.queue.ShutDown() jf.vcInformerFactory.Start(stopCh) for informerType, ok := range jf.vcInformerFactory.WaitForCacheSync(stopCh) { if !ok { return } } go wait.Until(jf.worker, time.Second, stopCh) <-stopCh } func (jf *jobflowcontroller) worker() { for jf.processNextWorkItem() { } } func (jf *jobflowcontroller) processNextWorkItem() bool { req, shutdown := jf.queue.Get() ... defer jf.queue.Done(req) err := jf.syncHandler(&req) jf.handleJobFlowErr(err, req) return true } func (jf *jobflowcontroller) handleJobFlow(req *apis.FlowRequest) error { ... jobflow, err := jf.jobFlowLister.JobFlows(req.Namespace).Get(req.JobFlowName) ... // Pending/Running/Failed/Terminating/Succeed jobFlowState := jobflowstate.NewState(jobflow) ... // pendingState.Execute // runningState.Execute // failedState.Execute(fake) // terminatingState.Execute(fake) // succeedState.Execute jobFlowState.Execute(req.Action) ... 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
注意
jobFlow有5种状态执行器,terminatingState和failedState都是空实现
# 3.3.state
jobFlow仅实现了三种状态执行器,分别为pendingState/runningState/succeedState,用于同步及修改jobFlow对象状态。func (p *pendingState) Execute(action jobflowv1alpha1.Action) error { switch action { case jobflowv1alpha1.SyncJobFlowAction: return SyncJobFlow(p.jobFlow, func(status *jobflowv1alpha1.JobFlowStatus, allJobList int) { if (len(status.RunningJobs) > 0 || len(status.CompletedJobs) > 0) && len(status.FailedJobs) <= 0 { status.State.Phase = jobflowv1alpha1.Running } else if len(status.FailedJobs) > 0 { status.State.Phase = jobflowv1alpha1.Failed } else { status.State.Phase = jobflowv1alpha1.Pending } }) } return nil } func (p *runningState) Execute(action v1alpha1.Action) error { switch action { case v1alpha1.SyncJobFlowAction: return SyncJobFlow(p.jobFlow, func(status *v1alpha1.JobFlowStatus, allJobList int) { if len(status.CompletedJobs) == allJobList { status.State.Phase = v1alpha1.Succeed } }) } return nil } func (p *succeedState) Execute(action v1alpha1.Action) error { switch action { case v1alpha1.SyncJobFlowAction: return SyncJobFlow(p.jobFlow, func(status *v1alpha1.JobFlowStatus, allJobList int) {}) } return nil } func (jf *jobflowcontroller) syncJobFlow(...) error { // jobFlow完成&保留策略为delete if jobFlow.Spec.JobRetainPolicy == Delete && jobFlow.Status.State.Phase == Succeed { jf.deleteAllJobsCreatedByJobFlow(jobFlow) ... return nil } // 根据jobFlow设置的jobTemplate创建job,声明顺序即为创建顺序 jf.deployJob(jobFlow) ... // all job status jobFlowStatus, err := jf.getAllJobStatus(jobFlow) ... // update jobFlow status jobFlow.Status = *jobFlowStatus updateStateFn(&jobFlow.Status, len(jobFlow.Spec.Flows)) jf.vcClient.FlowV1alpha1().JobFlows(jobFlow.Namespace).UpdateStatus(context.Background(), jobFlow, ...) ... return nil } func (jf *jobflowcontroller) deployJob(jobFlow *v1alpha1flow.JobFlow) error { // load jobTemplate by flow and deploy it for _, flow := range jobFlow.Spec.Flows { jobName := getJobName(jobFlow.Name, flow.Name) if _, err := jf.jobLister.Jobs(jobFlow.Namespace).Get(jobName); err != nil { // job未创建 if errors.IsNotFound(err) { // 无依赖 if flow.DependsOn == nil || flow.DependsOn.Targets == nil { jf.createJob(jobFlow, flow) ... // 有依赖 } else { // 检查依赖项 flag, err := jf.judge(jobFlow, flow) ... // 依赖完成 if flag { jf.createJob(jobFlow, flow) ... } } continue } return err } } 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
注意
state.Execute()主要检查jobFlow状态,回收相关Job或下发未完成的Job,根据下发Job状态回填更新JobFlow状态
# 4.gccontroller
# 4.1.initialize
job运行结束会保留在集群种,大规模任务场景会造成etcd作业压力,为了降低集群性能压力,job允许设置TTL,由gccontroller回收job。// Initialize creates an instance of gccontroller. func (gc *gccontroller) Initialize(opt *framework.ControllerOption) error { ... gc.jobInformer = jobInformer gc.jobLister = jobInformer.Lister() ... gc.queue = workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[string]()) gc.workers = opt.WorkerThreadsForGC jobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: gc.addJob, UpdateFunc: gc.updateJob, }) return nil } func (gc *gccontroller) addJob(obj interface{}) { job := obj.(*v1alpha1.Job) // job completed/failed/terminated&TTL if job.DeletionTimestamp == nil && needsCleanup(job) { gc.enqueue(job) } } func (gc *gccontroller) updateJob(old, cur interface{}) { job := cur.(*v1alpha1.Job) if job.DeletionTimestamp == nil && needsCleanup(job) { gc.enqueue(job) } }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
注意
initialize主要监听job对象,基于DeletionTimestamp和TTL字段入队job触发回收
# 4.2.gcstart
controller.Run()会执行worker协调,检测走到终态的job.spec.TTL过期状态,过期的job会进行删除,未过期的会重新触发入队。// Run starts the worker to clean up Jobs. func (gc *gccontroller) Run(stopCh <-chan struct{}) { defer gc.queue.ShutDown() gc.vcInformerFactory.Start(stopCh) for informerType, ok := range gc.vcInformerFactory.WaitForCacheSync(stopCh) { if !ok { return } } // default 1 worker for i := 0; i < int(gc.workers); i++ { // 间隔1s触发一次 go wait.Until(gc.worker, time.Second, stopCh) } <-stopCh } func (gc *gccontroller) worker() { for gc.processNextWorkItem() { } } func (gc *gccontroller) processNextWorkItem() bool { key, quit := gc.queue.Get() ... defer gc.queue.Done(key) err := gc.processJob(key) gc.handleErr(err, key) return true } // check the Job's state and TTL and delete the Job when it finishes and its TTL after finished has expired. func (gc *gccontroller) processJob(key string) error { namespace, name, err := cache.SplitMetaNamespaceKey(key) ... // Ignore the Jobs that are already deleted or being deleted, or the ones that don't need clean up. job, err := gc.jobLister.Jobs(namespace).Get(name) ... // 未过期 if expired := gc.processTTL(job); !expired { return nil } // 再次获取job fresh, err := gc.vcClient.BatchV1alpha1().Jobs(namespace).Get(context.TODO(), name, metav1.GetOptions{}) ... // 未过期 if expired := gc.processTTL(fresh); !expired { return nil } ... // 删除job gc.vcClient.BatchV1alpha1().Jobs(fresh.Namespace).Delete(context.TODO(), fresh.Name, options) ... return err } // checks whether a given Job's TTL has expired, and add it to the queue after the TTL is expected to expire // if the TTL will expire later. func (gc *gccontroller) processTTL(job *v1alpha1.Job) (expired bool, err error) { // We don't care about the Jobs that are going to be deleted, or the ones that don't need clean up. if job.DeletionTimestamp != nil || !needsCleanup(job) { return false, nil } // 计算距当前的过期时间 now := time.Now() t, err := timeLeft(job, &now) ... // TTL has expired if *t <= 0 { return true, nil } // 未过期重入队 gc.enqueueAfter(job, *t) return false, 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
注意
gc controller监听的是设置TTL及走到终态的pod