jobcontroller
# 1.executor
# 1.1.state
vcjob本质是一个状态机,供提供10种phase对应8种executor,executor会执行ActionFn/KillActionFn及updateStatuFn。// State interface. type State interface { // Execute executes the actions based on current state. Execute(act Action) error } // NewState gets the state from the volcano job Phase. func NewState(jobInfo *apis.JobInfo) State { job := jobInfo.Job switch job.Status.State.Phase { case vcbatch.Pending: return &pendingState{job: jobInfo} case vcbatch.Running: return &runningState{job: jobInfo} case vcbatch.Restarting: return &restartingState{job: jobInfo} case vcbatch.Terminated, vcbatch.Completed, vcbatch.Failed: return &finishedState{job: jobInfo} case vcbatch.Terminating: return &terminatingState{job: jobInfo} case vcbatch.Aborting: return &abortingState{job: jobInfo} case vcbatch.Aborted: return &abortedState{job: jobInfo} case vcbatch.Completing: return &completingState{job: jobInfo} } // default pending return &pendingState{job: jobInfo} }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
注意
state差异会执行不同的execute逻辑,针对job.status更新策略也不同
# 1.2.pending
pendingState基于action行为触发ActionFn/KillActionFn,不同的action会关联不同的updateStatusFn更新job状态。func (ps *pendingState) Execute(action Action) error { switch action.Action { // restart job case v1alpha1.RestartJobAction: return KillJob(ps.job, PodRetainPhaseNone, func(status *vcbatch.JobStatus) bool { status.RetryCount++ status.State.Phase = vcbatch.Restarting return true }) // restart task/pod case v1alpha1.RestartTaskAction, v1alpha1.RestartPodAction: return KillTarget(ps.job, action.Target, func(status *vcbatch.JobStatus) bool { status.RetryCount++ status.State.Phase = vcbatch.Restarting return true }) // abort job case v1alpha1.AbortJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Aborting return true }) // complete job case v1alpha1.CompleteJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Completing return true }) // terminate job case v1alpha1.TerminateJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Terminating return true }) default: return SyncJob(ps.job, func(status *vcbatch.JobStatus) bool { if ps.job.Job.Spec.MinAvailable <= status.Running+status.Succeeded+status.Failed { status.State.Phase = vcbatch.Running return true } return false }) } }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
补充
pendingState会基于action执行不同actionFn,区别在于回调的updateStatusFn
# 1.3.running
runningState类似于pendingState,相对复杂的是syncJob逻辑,根据job及关联task副本数和最小可用关系更新job.status状态。func (ps *runningState) Execute(action Action) error { switch action.Action { // restart job case v1alpha1.RestartJobAction: return KillJob(ps.job, PodRetainPhaseNone, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Restarting status.RetryCount++ return true }) // restart task/pod case v1alpha1.RestartTaskAction, v1alpha1.RestartPodAction: return KillTarget(ps.job, action.Target, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Restarting status.RetryCount++ return true }) // abort job case v1alpha1.AbortJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Aborting return true }) // terminate job case v1alpha1.TerminateJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Terminating return true }) // complete job case v1alpha1.CompleteJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Completing return true }) default: return SyncJob(ps.job, func(status *vcbatch.JobStatus) bool { ... if jobReplicas == 0 { // when scale down to zero, keep the current job phase return false } minSuccess := ps.job.Job.Spec.MinSuccess if minSuccess != nil && status.Succeeded >= *minSuccess { status.State.Phase = vcbatch.Completed return true } totalTaskMinAvailable := TotalTaskMinAvailable(ps.job.Job) // task结束 if status.Succeeded+status.Failed == jobReplicas { // 达到totalTaskMinAvailable if ps.job.Job.Spec.MinAvailable >= totalTaskMinAvailable { for _, task := range ps.job.Job.Spec.Tasks { if task.MinAvailable == nil { continue } if taskStatus, ok := status.TaskStatusCount[task.Name]; ok { // task未达到MinAvailable if taskStatus.Phase[v1.PodSucceeded] < *task.MinAvailable { status.State.Phase = vcbatch.Failed return true } } } } // job未达到minSuccess if minSuccess != nil && status.Succeeded < *minSuccess { status.State.Phase = vcbatch.Failed // job达到MinAvailable } else if status.Succeeded >= ps.job.Job.Spec.MinAvailable { status.State.Phase = vcbatch.Completed } else { status.State.Phase = vcbatch.Failed } return true } // task未全部完成 if status.Pending > jobReplicas-ps.job.Job.Spec.MinAvailable { status.State.Phase = vcbatch.Pending return true } return false }) } }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
注意
runningState执行syncJob会检查job及关联task副本数与最小可用限制,基于对比状态更新job.status
# 1.4.restarting
restartingState基于不同action执行ActionFn,updateStatusFn相同,会基于重试次数及剩余副本数对比限制条件更新job.status。func (ps *restartingState) Execute(action Action) error { switch action.Action { // sync job case v1alpha1.SyncJobAction: return SyncJob(ps.job, ps.restartingUpdateStatus) // restart task/pod case v1alpha1.RestartTaskAction, v1alpha1.RestartPodAction: return KillTarget(ps.job, action.Target, ps.restartingUpdateStatus) default: return KillJob(ps.job, PodRetainPhaseNone, ps.restartingUpdateStatus) } } func (ps *restartingState) restartingUpdateStatus(status *vcbatch.JobStatus) bool { // maximum number of retries. maxRetry := ps.job.Job.Spec.MaxRetry // 达到maxRetry if status.RetryCount >= maxRetry { status.State.Phase = vcbatch.Failed return true } total := int32(0) for _, task := range ps.job.Job.Spec.Tasks { total += task.Replicas } // remaining>MinAvailable,回退为pending if total-status.Terminating >= status.MinAvailable { status.State.Phase = vcbatch.Pending return true } return false }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
# 1.5.remaining
finishedState/terminatingState/abortingState/abortedState/completingState类似,基于action执行ActionFn及更新状态。func (ps *finishedState) Execute(action Action) error { // In finished state, always kill the whole job. return KillJob(ps.job, PodRetainPhaseSoft, nil) } func (ps *terminatingState) Execute(action Action) error { return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { // any "alive" pods, still in Terminating phase. if status.Terminating != 0 || status.Pending != 0 || status.Running != 0 { return false } status.State.Phase = vcbatch.Terminated return true }) } func (ps *abortingState) Execute(action Action) error { switch action.Action { // resume job case v1alpha1.ResumeJobAction: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Restarting status.RetryCount++ return true }) default: return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { // any "alive" pods, still in Aborting phase if status.Terminating != 0 || status.Pending != 0 || status.Running != 0 { return false } status.State.Phase = vcbatch.Aborted return true }) } } func (as *abortedState) Execute(action Action) error { switch action.Action { // resume job case v1alpha1.ResumeJobAction: return KillJob(as.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { status.State.Phase = vcbatch.Restarting status.RetryCount++ return true }) default: return KillJob(as.job, PodRetainPhaseSoft, nil) } } func (ps *completingState) Execute(action Action) error { return KillJob(ps.job, PodRetainPhaseSoft, func(status *vcbatch.JobStatus) bool { // any "alive" pods, still in Completing phase if status.Terminating != 0 || status.Pending != 0 || status.Running != 0 { return false } status.State.Phase = vcbatch.Completed 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
59
60
61
62
63
64
# 2.synchandler
# 2.1.syncJob
cc.syncJob()会检查job关联的pod,根据job.spec.tasks定义及job状态构建需要创建和删除的pod集合,createPod会设置调度器。func (cc *jobcontroller) syncJob(jobInfo *apis.JobInfo, updateStatus state.UpdateStatusFn) error { job := jobInfo.Job if jobInfo.Job.DeletionTimestamp != nil { return nil } // copy job to prevent mutate it job = job.DeepCopy() // job关联的queue queueInfo, err := cc.GetQueueInfo(job.Spec.Queue) ... var jobForwarding bool // 跨集群标记 if len(queueInfo.Spec.ExtendClusters) != 0 { jobForwarding = true ... // 标记job可能调度到其它集群 job.Annotations[batch.JobForwardingKey] = "true" job = cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).Update(context.TODO(), job, ...) ... } // phase为空或pending if !isInitiated(job) { // 更新job状态,调用add插件更新podgroup job = cc.initiateJob(job) ... } else { // 调用add插件更新podgroup cc.initOnJobUpdate(job) ... } // 尝试更新跨集群标记 if len(queueInfo.Spec.ExtendClusters) != 0 { jobForwarding = true job.Annotations[batch.JobForwardingKey] = "true" cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).Update(context.TODO(), job, ...) ... } ... // 获取job关联podgroup pg := cc.getPodGroupByJob(job) ... // pg设置过状态,未卡在pending(资源足够调度) if pg != nil & pg.Status.Phase != "" && pg.Status.Phase != scheduling.PodGroupPending { syncTask = true ... } ... // pg无法推进 if !syncTask { // 更新job.status if updateStatus != nil { updateStatus(&job.Status) } // 状态无变化 if equality.Semantic.DeepEqual(job.Status, oldStatus) { return nil } // 更新job状态 job.Status.State.LastTransitionTime = metav1.Now() jobCondition = newCondition(job.Status.State.Phase, &job.Status.State.LastTransitionTime) job.Status.Conditions = append(job.Status.Conditions, jobCondition) newJob := cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).UpdateStatus(context.TODO(), job, ...) ... // 更新缓存 cc.cache.Update(newJob) ... return nil } ... // 调度job task for _, ts := range job.Spec.Tasks { ... // 创建副本Pod for i := 0; i < int(ts.Replicas); i++ { podName := fmt.Sprintf("%s-%s-%d", job.Name, ts.Name, i) // 未创建过 if pod, found := pods[podName]; !found { // 生成pod,挂载PVC,设置annotation及label newPod := createJobPod(job, tc, ts.TopologyPolicy, i, jobForwarding) // 基于plugin补充rsa/cm/volumeMount/label/annotation/spec cc.pluginOnPodCreate(job, newPod) ... podToCreateEachTask = append(podToCreateEachTask, newPod) // pod创建过 } else { delete(pods, podName) // pod正在删除 if pod.DeletionTimestamp != nil { atomic.AddInt32(&terminating, 1) continue } // 分类计算pod phase数量 classifyAndAddUpPodBaseOnPhase(pod, &pending, &running, &succeeded, &failed, &unknown) // 计算task pod phase数量 calcPodStatus(pod, taskStatusCount) } } podToCreate[ts.Name] = podToCreateEachTask for _, pod := range pods { podToDelete = append(podToDelete, pod) } } // 待创建pod for taskName, podToCreateEachTask := range podToCreate { if len(podToCreateEachTask) == 0 { continue } go func(taskName string, podToCreateEachTask []*v1.Pod) { ... // depend未就绪 if job.Spec.Tasks[taskIndex].DependsOn != nil && !cc.waitDependsOnTaskMeetCondition(taskIndex,job) { // release wait group ... return } for _, pod := range podToCreateEachTask { go func(pod *v1.Pod) { ... // 创建pod newPod, err := cc.kubeClient.CoreV1().Pods(pod.Namespace).Create(context.TODO(), pod, ...) ... // 计算pod phase状态,统计taskStatusCount classifyAndAddUpPodBaseOnPhase(newPod, &pending, &running, &succeeded, &failed, &unknown) calcPodStatus(newPod, taskStatusCount) }(pod) } }(taskName, podToCreateEachTask) } ... // 创建出错 if len(creationErrs) != 0 { return fmt.Errorf("failed to create %d pods of %d", len(creationErrs), len(podToCreate)) } ... // 多余的pod for _, pod := range podToDelete { go func(pod *v1.Pod) { ... // 删除pod err := cc.deleteJobPod(job.Name, pod) if err != nil { ... // 失败推入errsTask队列 cc.resyncTask(pod) } else { atomic.AddInt32(&terminating, 1) } }(pod) } ... // 删除出错 if len(deletionErrs) != 0 { return fmt.Errorf("failed to delete %d pods of %d", len(deletionErrs), len(podToDelete)) } ... // updateStatusFn if updateStatus != nil { updateStatus(&newStatus) } // 状态无变化 if reflect.DeepEqual(job.Status, newStatus) { return nil } // 更新job.status job.Status = newStatus job.Status.State.LastTransitionTime = metav1.Now() jobCondition = newCondition(job.Status.State.Phase, &job.Status.State.LastTransitionTime) job.Status.Conditions = append(job.Status.Conditions, jobCondition) newJob, err := cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).UpdateStatus(context.TODO(), job, ...) ... // 更新cache job cc.cache.Update(newJob) ... 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
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
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208注意
cc.syncJob会检查关联的podgroup状态,不可推进仅更新job.status,否则尝试创建或回收job.task.pod及更新job.status
# 2.2.initiatejob
cc.initiateJob()会初始化job.status,尝试创建hostsCM/noneSvc/NetworkPolicy/rsaSecret/PVC/PG资源,供pod挂载或调度。func (cc *jobcontroller) initiateJob(job *batch.Job) (*batch.Job, error) { jobInstance := cc.initJobStatus(job) ... cc.pluginOnJobAdd(jobInstance) ... newJob := cc.createJobIOIfNotExist(jobInstance) ... cc.createOrUpdatePodGroup(newJob) ... return newJob, nil } func (cc *jobcontroller) initJobStatus(job *batch.Job) (*batch.Job, error) { if job.Status.State.Phase != "" { return job, nil } job.Status.State.Phase = batch.Pending job.Status.State.LastTransitionTime = metav1.Now() job.Status.MinAvailable = job.Spec.MinAvailable jobCondition := newCondition(job.Status.State.Phase, &job.Status.State.LastTransitionTime) job.Status.Conditions = append(job.Status.Conditions, jobCondition) newJob := cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).UpdateStatus(context.TODO(), job, ...) ... cc.cache.Update(newJob) ... return newJob, nil } func (cc *jobcontroller) pluginOnJobAdd(job *batch.Job) error { client := pluginsinterface.PluginClientset{KubeClients: cc.kubeClient} ... for name, args := range job.Spec.Plugins { pb, found := plugins.GetPluginBuilder(name) ... // 基于plugin更新job pb(client, args).OnJobAdd(job) ... } return nil } // plugin func (sp *servicePlugin) OnJobAdd(job *batch.Job) error { // 执行过 if job.Status.ControlledResources["plugin-"+sp.Name()] == sp.Name() { return nil } // 生成hosts文件内容 hostFile := GenerateHosts(job) // create or update ConfigMap of hosts for Pods to mount. helpers.CreateOrUpdateConfigMap(job, sp.Clientset.KubeClients, hostFile, sp.cmName(job)) ... // create none svc for job pod. sp.createServiceIfNotExist(job) ... // 未禁用网络策略 if !sp.disableNetworkPolicy { // 创建网络策略(NetworkPolicy),限制允许job所属pod互相访问 sp.createNetworkPolicyIfNotExist(job) ... } job.Status.ControlledResources["plugin-"+sp.Name()] = sp.Name() return nil } func (sp *sshPlugin) OnJobAdd(job *batch.Job) error { // 执行过 if job.Status.ControlledResources["plugin-"+sp.Name()] == sp.Name() { return nil } ... // 公私钥 if len(sp.sshPrivateKey) > 0 { data, err = withUserProvidedRsaKey(job, sp.sshPrivateKey, sp.sshPublicKey) } else { data, err = generateRsaKey(job) } ... // 创建或更新rsa secret helpers.CreateOrUpdateSecret(job, sp.client.KubeClients, data, sp.secretName(job)) ... job.Status.ControlledResources["plugin-"+sp.Name()] = sp.Name() return nil } func (cc *jobcontroller) createJobIOIfNotExist(job *batch.Job) (*batch.Job, error) { ... for index, volume := range job.Spec.Volumes { vcName := volume.VolumeClaimName // 未设置pvcName if len(vcName) == 0 { for { // 生成唯一的pvcName vcName = jobhelpers.GenPVCName(job.Name) if cc.checkPVCExist(job, vcName) { continue } ... job.Spec.Volumes[index].VolumeClaimName = vcName needUpdate = true break } // 定义volume claim if volume.VolumeClaim != nil { cc.createPVC(job, vcName, volume.VolumeClaim) ... } // volume关联pvcName } else { // 检查pvc存在 if !cc.checkPVCExist(job, vcName) { return job, fmt.Errorf("pvc not found, job will be in Pending state until PVC created", vcName) } } job.Status.ControlledResources["volume-pvc-"+vcName] = vcName } // 更新job if needUpdate { newJob,err := cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).Update(context.TODO(), job, ...) ... newJob.Status = job.Status return newJob, err } return job, nil } func (cc *jobcontroller) createOrUpdatePodGroup(job *batch.Job) error { // If PodGroup does not exist, create one for Job. pg, err := cc.getPodGroupByJob(job) if err != nil { if !apierrors.IsNotFound(err) { return err } ... // 未找到job pg进行创建 cc.vcClient.SchedulingV1beta1().PodGroups(job.Namespace).Create(context.TODO(), pg, ...) ... return nil } ... // 计算pg最小资源 // 1.job.Spec.MinAvailable < totalTaskMinAvailable: 基于优先级排序task,基于高优task凑够job.Spec.MinAvailable // 2.job.Spec.MinAvailable ≥ totalTaskMinAvailable: 基于优先级排序task分MinAvailable,多出的基于高优凑够 minResources := cc.calcPGMinResources(job) // 更新MinMember和MinResources if pg.Spec.MinMember != job.Spec.MinAvailable || !equality.DeepEqual(pg.Spec.MinResources, minResources) { pg.Spec.MinMember = job.Spec.MinAvailable pg.Spec.MinResources = minResources pgShouldUpdate = true } ... // 更新MinTaskMember for _, task := range job.Spec.Tasks { cnt := task.Replicas if task.MinAvailable != nil { cnt = *task.MinAvailable } if taskMember, ok := pg.Spec.MinTaskMember[task.Name]; !ok { pgShouldUpdate = true pg.Spec.MinTaskMember[task.Name] = cnt } else { if taskMember == cnt { continue } pgShouldUpdate = true pg.Spec.MinTaskMember[task.Name] = cnt } } if !pgShouldUpdate { return nil } cc.vcClient.SchedulingV1beta1().PodGroups(job.Namespace).Update(context.TODO(), pg, metav1.UpdateOptions{}) ... return 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
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
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
注意
pluginInterface主要实现是servicePlugin/sshPlugin,其它几个均基于设置env工作,PG加载会根据Task优先级更新资源
# 2.3.initupdate
cc.initOnJobUpdate()尝试执行plugin.OnJobUpdate()更新hostsCM,基于job最新定义调整关联podgroup的资源信息及更新。func (cc *jobcontroller) initOnJobUpdate(job *batch.Job) error { cc.pluginOnJobUpdate(job) ... cc.createOrUpdatePodGroup(job) ... return nil } func (cc *jobcontroller) pluginOnJobUpdate(job *batch.Job) error { ... for name, args := range job.Spec.Plugins { pb, found := plugins.GetPluginBuilder(name) ... pb(client, args).OnJobUpdate(job) ... } return nil } func (sp *servicePlugin) OnJobUpdate(job *batch.Job) error { hostFile := GenerateHosts(job) // updates ConfigMap of hosts for Pods to mount. return helpers.CreateOrUpdateConfigMap(job, sp.Clientset.KubeClients, hostFile, sp.cmName(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
注意
servicePlugin实现了plugin.OnJobUpdate,其它plugin均未实现
# 2.4.podcreate
cc.pluginOnPodCreate()会基于plugin进一步修饰Pod,加入configmap/rsaSecret挂载、hosts及host:port环境变量配置。func (cc *jobcontroller) pluginOnPodCreate(job *batch.Job, pod *v1.Pod) error { ... for name, args := range job.Spec.Plugins { pb, found := plugins.GetPluginBuilder(name) ... pb(client, args).OnPodCreate(pod, job) ... } return nil } // dns:podName.jobName.default.svc.cluster.local func (sp *servicePlugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { ... // 设置VC_taskName_HOSTS和VC_taskName_NUM for _, ts := range job.Spec.Tasks { formateENVKey := strings.Replace(ts.Name, "-", "_", -1) envNames = append(envNames, fmt.Sprintf(EnvTaskHostFmt, strings.ToUpper(formateENVKey))) envNames = append(envNames, fmt.Sprintf(EnvHostNumFmt, strings.ToUpper(formateENVKey))) } // 引用之前生成的cm内容(hosts/num) for _, name := range envNames { hostEnv = append(hostEnv, v1.EnvVar{ Name: name, ValueFrom: &v1.EnvVarSource{ ConfigMapKeyRef: &v1.ConfigMapKeySelector{ LocalObjectReference: v1.LocalObjectReference{Name: sp.cmName(job)}, Key: name, }}}, ) } // 设置container env for i := range pod.Spec.Containers { pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, hostEnv...) } // 设置initcontainer env for i := range pod.Spec.InitContainers { pod.Spec.InitContainers[i].Env = append(pod.Spec.InitContainers[i].Env, hostEnv...) } // 挂载cm sp.mountConfigmap(pod, job) return nil } func (sp *sshPlugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { // 挂载rsa secret sp.mountRsaKey(pod, job) return nil } func (ep *envPlugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { index := jobhelpers.GetPodIndexUnderTask(pod) // add VK_TASK_INDEX and VC_TASK_INDEX env to each container for i := range pod.Spec.Containers { pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{TaskVkIndex, index}) pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{TaskIndex, index}) } // add VK_TASK_INDEX and VC_TASK_INDEX env to each init container for i := range pod.Spec.InitContainers { pod.Spec.InitContainers[i].Env = append(pod.Spec.InitContainers[i].Env, v1.EnvVar{TaskVkIndex, index}) pod.Spec.InitContainers[i].Env = append(pod.Spec.InitContainers[i].Env, v1.EnvVar{TaskIndex, index}) } return nil } func (tp *tensorflowPlugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { // No need to generate TF_CONFIG for stand-alone tensorflow job if len(job.Spec.Tasks) == 1 && job.Spec.Tasks[0].Replicas == 1 { return nil } // Generate TF_CONFIG value // task pod index index, err := strconv.Atoi(jobhelpers.GetPodIndexUnderTask(pod)) ... // Generate tensorflow task info spec := tfClusterSpec{ Task: taskInfo{ Type: tp.getTaskType(jobhelpers.GetTaskKey(pod)), // pod type from annotation/default Index: index, }, } // Generate tensorflow cluster info for _, ts := range job.Spec.Tasks { ... // 各task pod生成dns:port for i := 0; i < int(ts.Replicas); i++ { hosts = append(hosts, fmt.Sprintf("%s:%d", jobhelpers.MakeDomainName(ts, job, i), tp.port)) } // 分类回填hosts switch ts.Name { case tp.psName: spec.Cluster.PS = hosts case tp.workerName: spec.Cluster.Worker = hosts case tp.chiefName: spec.Cluster.Chief = hosts case tp.evaluatorName: spec.Cluster.Evaluator = hosts } } ... raw, err := json.Marshal(spec) ... // Add TF_CONFIG enviroment variables for i := range pod.Spec.Containers { pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{ Name: TFConfig, Value: string(raw), }) } return nil } func (pp *pytorchPlugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { taskType := helpers.GetTaskKey(pod) masterIndex := helpers.GetTaskIndexUnderJob(pp.masterName, job) ... // master dns域名及端口 masterAddr := pp.generateMasterAddr(job.Spec.Tasks[masterIndex], job.Name) masterEnvVars = append(masterEnvVars, v1.EnvVar{ Name: EnvMasterAddr, Value: masterAddr, }, v1.EnvVar{ Name: EnvMasterPort, Value: fmt.Sprintf("%v", pp.port), }) ... // worker pod if taskType == pp.workerName { index, err := strconv.Atoi(helpers.GetPodIndexUnderTask(pod)) ... workerRank = index + 1 } // job task总副本数 totalReplicas := pp.getTotalReplicas(job) for i, c := range pod.Spec.Containers { // 打开container pytorch依赖端口 pp.openContainerPort(&c, i, pod) // 设置masterAddr及worldSize env pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, masterEnvVars...) pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{ Name: EnvWorldSize, Value: strconv.Itoa(int(totalReplicas)), }) // worker pod,设置workRank if taskType == pp.workerName { pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{ Name: EnvRank, Value: strconv.Itoa(workerRank), }) // master pod,rank=0 } else if taskType == pp.masterName { pod.Spec.Containers[i].Env = append(pod.Spec.Containers[i].Env, v1.EnvVar{ Name: EnvRank, Value: strconv.Itoa(masterRank), }) } } return nil } func (mp *Plugin) OnPodCreate(pod *v1.Pod, job *batch.Job) error { ... // master pod if helpers.GetTaskKey(pod) == mp.masterName { // 生成worker hosts workerHosts = mp.generateTaskHosts(job.Spec.Tasks[GetTaskIndexUnderJob(mp.workerName, job)], job.Name) env = v1.EnvVar{ Name: MPIHost, Value: workerHosts, } isMaster = true } // open port for ssh and add MPI_HOST env for master task for index, ic := range pod.Spec.InitContainers { // 打开ssh port mp.openContainerPort(&ic, index, pod, true) // 设置mpi_host env if isMaster { pod.Spec.InitContainers[index].Env = append(pod.Spec.InitContainers[index].Env, env) } } // open port for ssh and add MPI_HOST env for master task for index, c := range pod.Spec.Containers { mp.openContainerPort(&c, index, pod, false) if isMaster { pod.Spec.Containers[index].Env = append(pod.Spec.Containers[index].Env, env) } } 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
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
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
注意
tensorflowPlugin/pytorchPlugin/mpiPlugin本质做了一些定制化适配,注入打开特定端口及挂载框架运行必须的hosts/host:port
# 3.killhandler
# 3.1.killJob
cc.killJob()会尝试清理job关联的Pod,更新job.status及删除关联的PodGroup,Pod清理后会进一步触发plugin维护资源的回收。func (cc *jobcontroller) killJob(...) error { return cc.killPods(jobInfo, podRetainPhase, nil, updateStatus) } func (cc *jobcontroller) killTarget(...) error { ... return cc.killPods(jobInfo, nil, &target, updateStatus) } func (cc *jobcontroller) killPods(...) error { job := jobInfo.Job if job.DeletionTimestamp != nil { return nil } ... if target != nil { if target.Type == state.TargetTypeTask { podsToKill = jobInfo.Pods[target.TaskName] } else if target.Type == state.TargetTypePod { podsToKill[target.PodName] = jobInfo.Pods[target.TaskName][target.PodName] } total += len(podsToKill) } else { // Job version is bumped only when job is killed job.Status.Version++ for _, pods := range jobInfo.Pods { for _, pod := range pods { total++ // 正在删除 if pod.DeletionTimestamp != nil { terminating++ continue } ... // 检查最后一次重试 if job.Status.RetryCount >= maxRetry-1 { lastRetry = true } // only retain the Failed and Succeeded pods at the last retry. // If it is not the last retry, kill pod as defined in `podRetainPhase`. retainPhase := podRetainPhase if lastRetry { retainPhase = state.PodRetainPhaseSoft } // 保留phase _, retain := retainPhase[pod.Status.Phase] // 删除未保留phase if !retain { podsToKill[pod.Name] = pod } } } } for _, pod := range podsToKill { // 正在删除 if pod.DeletionTimestamp != nil { terminating++ continue } // 删除pod err := cc.deleteJobPod(job.Name, pod) if err == nil { terminating++ continue } ... // 推到errTasks,由process task异步处理 cc.errTasks.AddRateLimited(pod) // 分类统计phase数量 classifyAndAddUpPodBaseOnPhase(pod, &pending, &running, &succeeded, &failed, &unknown) calcPodStatus(pod, taskStatusCount) } ... // update job status with hook if updateStatus != nil & updateStatus(&job.Status) { job.Status.State.LastTransitionTime = metav1.Now() jobCondition := newCondition(job.Status.State.Phase, &job.Status.State.LastTransitionTime) job.Status.Conditions = append(job.Status.Conditions, jobCondition) } ... // must be called before update job status cc.pluginOnJobDelete(job) ... // Update Job status newJob := cc.vcClient.BatchV1alpha1().Jobs(job.Namespace).UpdateStatus(context.TODO(), job, ...) ... cc.cache.Update(newJob) ... // Delete PodGroup pg, err := cc.getPodGroupByJob(job) ... cc.vcClient.SchedulingV1beta1().PodGroups(job.Namespace).Delete(context.TODO(), pg.Name, ...) ... 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
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112注意
pluginOnJobDelete会清理hostsCM/rsaSecret/PVC资源,避免资源泄漏
# 3.2.jobdelete
cc.pluginOnJobDelete()基于job清理创建阶段准备的hosts configmap/headles svc/rsa secret资源及job相关的plugin注解。func (cc *jobcontroller) pluginOnJobDelete(job *batch.Job) error { ... for name, args := range job.Spec.Plugins { pb, found := plugins.GetPluginBuilder(name) ... pb(client, args).OnJobDelete(job) ... } return nil } func (sp *servicePlugin) OnJobDelete(job *batch.Job) error { // 未执行过 if job.Status.ControlledResources["plugin-"+sp.Name()] != sp.Name() { return nil } // 清理job configmap helpers.DeleteConfigmap(job, sp.Clientset.KubeClients, sp.cmName(job)) ... // 清理none svc sp.Clientset.KubeClients.CoreV1().Services(job.Namespace).Delete(context.TODO(), job.Name, ...) ... // 移除注解 delete(job.Status.ControlledResources, "plugin-"+sp.Name()) // 未禁用NetworkPolicy if !sp.disableNetworkPolicy { // 清理NetworkPolicy对象 sp.client.KubeClients.NetworkingV1().NetworkPolicies(job.Namespace).Delete(context.TODO(), job.Name,...) ... } return nil } func (sp *sshPlugin) OnJobDelete(job *batch.Job) error { // 未执行过 if job.Status.ControlledResources["plugin-"+sp.Name()] != sp.Name() { return nil } // 清理rsa secret helpers.DeleteSecret(job, sp.client.KubeClients, sp.secretName(job)) ... delete(job.Status.ControlledResources, "plugin-"+sp.Name()) 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
注意
envPlugin/tensorflowPlugin/pytorchPlugin/mpiPlugin实现均为空,仅清理pluginName注解