deployment
# 1.简介
# 1.1.定义
deployment controller是kube-controller-manager组件负责deployment资源对象的控制器,监听deploy/rs/pod触发协同操作。// DeploymentController is responsible for synchronizing Deployment objects stored // in the system with actual running replica sets and pods. type DeploymentController struct { rsControl controller.RSControlInterface // 认领/释放rs的控制器 ... syncHandler func(ctx context.Context, dKey string) error // deployment同步回调 enqueueDeployment func(deployment *apps.Deployment) // deployment入队回调 dLister appslisters.DeploymentLister // deployment缓存 rsLister appslisters.ReplicaSetLister // rs缓存 podLister corelisters.PodLister // pod缓存 ... queue workqueue.RateLimitingInterface // deployment调度队列 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
注意
deployment本质是控制replicaset,间接影响pod,pod具体的调度由kubelet完成
# 1.2.原理
deploy controller利用informer监听deploy/rs/pod,将相关的deploy推入workqueue,由syncHandler处理deployment任务。--- 执行过程 1.基于informer监听执行事件回调,将deployment推入workqueue 2.异步worker执行syncHandler,实时获取workqueue的item及处理 3.根据item执行rs领养/弃养 4.根据item的升级策略处理下级资源1
2
3
4
5注意
workqueue是一个先入先出的队列,由item切片、dirty map和processing map构成
# 2.分析
# 2.1.start
controller-manager基于注册的initFunc管理不同的controller,startDeploymentController就负责实例化及启动deployment。// NewDeploymentController creates a new DeploymentController. func NewDeploymentController(...) (*DeploymentController, error) { ... // 创建dc实例 dc := &DeploymentController{ client: client, ... queue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "deployment"), } // rs领养/弃养控制 dc.rsControl = controller.RealRSControl{ client, dc.eventRecorder } // deploy入队回调 dInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ ... }) // rs变化deploy入队回调 rsInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ ... }) // pod变化deploy入队回调 podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ ... }) // item处理回调 dc.syncHandler = dc.syncDeployment // deploy入队回调 dc.enqueueDeployment = dc.enqueue // 缓存 dc.dLister = dInformer.Lister() dc.rsLister = rsInformer.Lister() dc.podLister = podInformer.Lister() ... return dc, nil } func startDeploymentController(...) (controller.Interface, bool, error) { // 初始化 dc, err := deployment.NewDeploymentController(ctx, dInformer, rsInformer, podInfomer, client) ... // 执行 go dc.Run(ctx, int(controllerContext.ComponentConfig.DeploymentController.ConcurrentDeploymentSyncs)) return nil, 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注意
deployment controller会基于informer监听deploy/rs/pod三类资源,利用回调直接或间接入队deploy,由syncHandler处理
# 2.2.informer
deployment基于informer监听deploy/rs/pod资源变化,将相关的deploy推入workqueue,dinformer相对简单,这里分析rs/pod。// addReplicaSet enqueues the deployment that manages a ReplicaSet when the ReplicaSet is created. func (dc *DeploymentController) addReplicaSet(logger klog.Logger, obj interface{}) { rs := obj.(*apps.ReplicaSet) if rs.DeletionTimestamp != nil { // 执行delete流程 dc.deleteReplicaSet(logger, rs) return } // If it has a ControllerRef, that's all that matters. if controllerRef := metav1.GetControllerOf(rs); controllerRef != nil { // 解析owner d := dc.resolveControllerRef(rs.Namespace, controllerRef) ... // owner deploy推入workqueue dc.enqueueDeployment(d) return } // 所有selector匹配的deploy ds := dc.getDeploymentsForReplicaSet(logger, rs) ... // 推入workqueue for _, d := range ds { dc.enqueueDeployment(d) } } // updateReplicaSet figures out what deployment(s) manage a ReplicaSet. func (dc *DeploymentController) updateReplicaSet(logger klog.Logger, old, cur interface{}) { ... // rs无变化 if curRS.ResourceVersion == oldRS.ResourceVersion { return } ... // ownerRef变化,旧的rs ownerRef不为空 if controllerRefChanged && oldControllerRef != nil { // 旧的owner deploy推入workqueue if d := dc.resolveControllerRef(oldRS.Namespace, oldControllerRef); d != nil { dc.enqueueDeployment(d) } } // 新的rs ownerRef不为空 if curControllerRef != nil { // 解析旧的owner deploy d := dc.resolveControllerRef(curRS.Namespace, curControllerRef) ... // 新的owner deploy推入workqueue dc.enqueueDeployment(d) return } ... // rs labels变化或owner变化(无符合的owner deploy) if labelChanged || controllerRefChanged { // 基于selector匹配deploy ds := dc.getDeploymentsForReplicaSet(logger, curRS) ... // 匹配的deploy推入workqueue for _, d := range ds { dc.enqueueDeployment(d) } } } // deleteReplicaSet enqueues the deployment that manages a ReplicaSet. func (dc *DeploymentController) deleteReplicaSet(logger klog.Logger, obj interface{}) { rs, ok := obj.(*apps.ReplicaSet) if !ok { tombstone, ok := obj.(cache.DeletedFinalStateUnknown) ... rs, ok = tombstone.Obj.(*apps.ReplicaSet) ... } // 获取rs owner controllerRef := metav1.GetControllerOf(rs) ... // 解析owner deploy d := dc.resolveControllerRef(rs.Namespace, controllerRef) ... // owner deploy入队 dc.enqueueDeployment(d) } // deletePod will enqueue a Recreate Deployment once all of its pods have stopped running. func (dc *DeploymentController) deletePod(logger klog.Logger, obj interface{}) { pod, ok := obj.(*v1.Pod) if !ok { tombstone, ok := obj.(cache.DeletedFinalStateUnknown) ... pod, ok = tombstone.Obj.(*v1.Pod) ... } // pod隶属recreate类型的deploy if d := dc.getDeploymentForPod(logger, pod); d != nil && d.Spec.Strategy.Type == apps.RecreateDeploymentStrategyType { // 获取deploy关联的rs rsList, err := util.ListReplicaSets(d, util.RsListFromClient(dc.client.AppsV1())) ... // 根据rs认领pod podMap, err := dc.getPodMapForDeployment(d, rsList) ... // 旧pod全部删除,关联的deploy推入workqueue if numPods == 0 { dc.enqueueDeployment(d) } } }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注意
rs和pod基于owner引用及label匹配关联deploy,将相关的deploy推入workqueue处理
# 2.3.run
dc.run()会启动controller,阻塞至informer的缓存同步完成,激活一定数量(默认5)的worker协程消费queue及执行syncHandler。// Run begins watching and syncing. func (dc *DeploymentController) Run(ctx context.Context, workers int) { ... defer dc.queue.ShutDown() // 阻塞至缓存同步完毕 if !cache.WaitForNamedCacheSync("deployment", ctx.Done(), dListerSynced, rsListerSynced, podListerSynced) { return } // 激活worker for i := 0; i < workers; i++ { go wait.UntilWithContext(ctx, dc.worker, time.Second) } <-ctx.Done() } // worker runs a worker thread that just dequeues items, processes them, and marks them done. func (dc *DeploymentController) worker(ctx context.Context) { for dc.processNextWorkItem(ctx) { } } func (dc *DeploymentController) processNextWorkItem(ctx context.Context) bool { // 推出item key, quit := dc.queue.Get() ... // 由processing map删除 defer dc.queue.Done(key) // 执行syncHandler err := dc.syncHandler(ctx, key.(string)) dc.handleErr(ctx, err, key) return true } func (dc *DeploymentController) handleErr(ctx context.Context, err error, key interface{}) { ... // 成功/ns终止 if err == nil || errors.HasStatusCause(err, v1.NamespaceTerminatingCause) { // 清理限流状态,避免满足条件再次加入queue dc.queue.Forget(key) return } ... // 未达到最大重试 if dc.queue.NumRequeues(key) < maxRetries { // 设置限流策略,满足条件加入queue dc.queue.AddRateLimited(key) return } ... // 其它情况,清理限流状态,避免满足条件再次加入queue dc.queue.Forget(key) }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注意
worker本质上是无限循环轮询器,不断取出queue item交给syncHandler处理
# 2.4.syncHandler
syncHandler不同阶段会进行相应处理,syncStatusOnly处理删除,sync处理pause状态,rollback处理回滚,rollout处理升级。// syncDeployment will sync the deployment with the given key. // This function is not meant to be invoked concurrently with the same key. func (dc *DeploymentController) syncDeployment(ctx context.Context, key string) error { // name/namespace namespace, name, err := cache.SplitMetaNamespaceKey(key) ... // 获取deployment deployment, err := dc.dLister.Deployments(namespace).Get(name) ... // 拷贝deployment,避免修改影响缓存 d := deployment.DeepCopy() ... // selector为空 if reflect.DeepEqual(d.Spec.Selector, &everything) { // observeGeneration更新(无selector,没有可筛选的rs和pod) if d.Status.ObservedGeneration < d.Generation { d.Status.ObservedGeneration = d.Generation dc.client.AppsV1().Deployments(d.Namespace).UpdateStatus(ctx, d, metav1.UpdateOptions{}) } return nil } // 获取关联的rs rsList, err := dc.getReplicaSetsForDeployment(ctx, d) ... // 获取关联的pod podMap, err := dc.getPodMapForDeployment(d, rsList) ... // 删除状态处理 if d.DeletionTimestamp != nil { return dc.syncStatusOnly(ctx, d, rsList) } // 检查pause状态及更新pause condition dc.checkPausedConditions(ctx, d) ... // pause状态处理 if d.Spec.Paused { return dc.sync(ctx, d, rsList) } // 回滚处理 // 标记deprecated.deployment.rollback.to注解,该注解后续会废弃 if getRollbackTo(d) != nil { return dc.rollback(ctx, d, rsList) } // 检查期望副本数与rs期望一致性 scalingEvent, err := dc.isScalingEvent(ctx, d, rsList) ... // 扩容 if scalingEvent { return dc.sync(ctx, d, rsList) } // 更新操作 switch d.Spec.Strategy.Type { // 重建模式 case apps.RecreateDeploymentStrategyType: return dc.rolloutRecreate(ctx, d, rsList, podMap) // 滚动更新模式 case apps.RollingUpdateDeploymentStrategyType: return dc.rolloutRolling(ctx, d, rsList) } return fmt.Errorf("unexpected deployment strategy type: %s", d.Spec.Strategy.Type) }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注意
syncHandler执行deployment同步有一定优先级,根据delete-->pause-->rollback-->scale-->rollout顺序分发
# 2.5.rsList
dc.getReplicaSetsForDeployment()会获取namespace所有的replicaset,基于ownerRef和label进行领养、弃养及认领。// getReplicaSetsForDeployment uses ControllerRefManager to reconcile ControllerRef by adopting and orphaning. func (dc *DeploymentController) getReplicaSetsForDeployment(ctx context.Context, d *apps.Deployment) (...) { // 获取命名空间的所有rs rsList, err := dc.rsLister.ReplicaSets(d.Namespace).List(labels.Everything()) ... // 获取deploy的selector deploymentSelector, err := metav1.LabelSelectorAsSelector(d.Spec.Selector) ... // 外部包装RecheckDeletionTimestamp,检查获取的deploy是否出于删除状态 canAdoptFunc := controller.RecheckDeletionTimestamp(func(ctx context.Context) (metav1.Object, error) { // 获取底层的deploy fresh, err := dc.client.AppsV1().Deployments(d.Namespace).Get(ctx, d.Name, metav1.GetOptions{}) ... // deploy删除重建,当前deploy无资格认领rs if fresh.UID != d.UID { return nil, fmt.Errorf("original Deployment %v/%v is gone: got uid %v, wanted %v", d.Namespace, d.Name, fresh.UID, d.UID) } return fresh, nil }) // 初始化认领对象 cm := controller.NewReplicaSetControllerRefManager(dc.rsControl, d, deploymentSelector, controllerKind, canAdoptFunc) // 认领deploy相关rs return cm.ClaimReplicaSets(ctx, rsList) } // ClaimReplicaSets tries to take ownership of a list of ReplicaSets. func (m *ReplicaSetControllerRefManager) ClaimReplicaSets(ctx context.Context, sets []*apps.ReplicaSet) (...) { ... for _, rs := range sets { // 认领rs ok, err := m.ClaimObject(ctx, rs, match, adopt, release) ... if ok { claimed = append(claimed, rs) } } return claimed, utilerrors.NewAggregate(errlist) } func (m *BaseControllerRefManager) ClaimObject(...) error) (bool, error) { // 获取rs的所有者 controllerRef := metav1.GetControllerOfNoCopy(obj) if controllerRef != nil { // 所有者不是当前deploy,无法认领 if controllerRef.UID != m.Controller.GetUID() { // Owned by someone else. Ignore. return false, nil } // selector匹配,认领 if match(obj) { return true, nil } // deploy正在删除,无法认领 if m.Controller.GetDeletionTimestamp() != nil { return false, nil } // deploy作为所有者健康,未匹配selector,弃养(rs成为孤儿) release(ctx, obj) ... // Successfully released. return false, nil } // deploy正在删除或selector未匹配rs,无法认领 if m.Controller.GetDeletionTimestamp() != nil || !match(obj) { // Ignore if we're being deleted or selector doesn't match. return false, nil } // rs正在删除,无法认领 if obj.GetDeletionTimestamp() != nil { // Ignore if the object is being deleted return false, nil } // deploy和rs命名空间不一致,无法认领 if len(m.Controller.GetNamespace()) > 0 && m.Controller.GetNamespace() != obj.GetNamespace() { // Ignore if namespace not match return false, nil } // rs无所有者,deploy健康且匹配selector,领养rs adopt(ctx, obj) ... // Successfully adopted. 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
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注意
deployment会领养孤儿replicaset,也会弃养label selector不匹配的replicaset,领养和弃养都是设置ownerRef
# 2.6.podList
dc.getPodMapForDeployment()会获取// getPodMapForDeployment returns the Pods managed by a Deployment. func (dc *DeploymentController) getPodMapForDeployment(d *apps.Deployment, rsList []*apps.ReplicaSet) (...) { // Get all Pods that potentially belong to this Deployment. selector, err := metav1.LabelSelectorAsSelector(d.Spec.Selector) ... // 获取匹配的pod pods, err := dc.podLister.Pods(d.Namespace).List(selector) ... for _, pod := range pods { // 获取pod ownerRef controllerRef := metav1.GetControllerOf(pod) if controllerRef == nil { continue } // 归属某个认领的rs if _, ok := podMap[controllerRef.UID]; ok { podMap[controllerRef.UID] = append(podMap[controllerRef.UID], pod) } } return podMap, 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注意
pod认领相对简单,基于selector获取匹配pod,基于已经认领的rs分组认领pod