daemonset
# 1.调度
# 1.1.runworker
dsc.Run()会启动多个worker执行syncHandler处理workqueue的daemonset任务,此外会启动一个gc worker回收过期的退避对象。// Run begins watching and syncing daemon sets. func (dsc *DaemonSetsController) Run(ctx context.Context, workers int) { ... defer dsc.queue.ShutDown() ... // informer同步完成 if !cache.WaitForNamedCacheSync("daemon sets", ctx.Done(), dsc.podStoreSynced, dsc.nodeStoreSynced, dsc.historyStoreSynced, dsc.dsStoreSynced) { return } // 一般激活2个worker for i := 0; i < workers; i++ { go wait.UntilWithContext(ctx, dsc.runWorker, time.Second) } // gc worker(间隔1min触发) go wait.Until(dsc.failedPodsBackoff.GC, BackoffGCInterval, ctx.Done()) <-ctx.Done() } func (dsc *DaemonSetsController) runWorker(ctx context.Context) { for dsc.processNextWorkItem(ctx) { } } // processNextWorkItem deals with one key off the queue. It returns false when it's time to quit. func (dsc *DaemonSetsController) processNextWorkItem(ctx context.Context) bool { dsKey, quit := dsc.queue.Get() ... defer dsc.queue.Done(dsKey) // 执行同步 err := dsc.syncHandler(ctx, dsKey.(string)) if err == nil { dsc.queue.Forget(dsKey) return true } ... dsc.queue.AddRateLimited(dsKey) return true } // GC records that have aged past maxDuration. Backoff users are expected to invoke this periodically. func (p *Backoff) GC() { p.Lock() defer p.Unlock() now := p.Clock.Now() for id, entry := range p.perItemBackoff { // 退避数据超时 if now.Sub(entry.lastUpdate) > p.maxDuration*2 { // GC when entry has not been updated for 2*maxDuration delete(p.perItemBackoff, id) } } }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注意
GC相对简单,仅定时清理超时的backoff item数据,同步的核心是syncHandler
# 1.2.syncHandler
dsc.syncDaemonSet()会获取node和ds对象,根据ds对象同步条件执行Pod创建或删除,维护node和pod匹配的数量状态。func (dsc *DaemonSetsController) syncDaemonSet(ctx context.Context, key string) error { ... namespace, name, err := cache.SplitMetaNamespaceKey(key) ... // 获取ds对象 ds, err := dsc.dsLister.DaemonSets(namespace).Get(name) if apierrors.IsNotFound(err) { // ds不存在,清理exp期望 dsc.expectations.DeleteExpectations(key) return nil } ... // 获取所有node nodeList, err := dsc.nodeLister.List(labels.Everything()) ... // ds及所属资源清理由GC负责,避免ds无法准确计算状态造成时序问题 if ds.DeletionTimestamp != nil { return nil } // 基于template匹配拆分新旧history对象 cur, old, err := dsc.constructHistory(ctx, ds) ... hash := cur.Labels[apps.DefaultDaemonSetUniqueLabelKey] if !dsc.expectations.SatisfiedExpectations(dsKey) { // 状态更新 return dsc.updateDaemonSetStatus(ctx, ds, nodeList, hash, false) } // ds同步 dsc.updateDaemonSet(ctx, ds, nodeList, hash, dsKey, old) // 状态更新 dsc.updateDaemonSetStatus(ctx, ds, nodeList, hash, true) ... return nil } // constructHistory finds all histories controlled by the given DaemonSet, and // update current history revision number, or create current history if need to. func (dsc *DaemonSetsController) constructHistory(ctx context.Context, ds *apps.DaemonSet) (...) { ... // 获取ds所属的history对象(涉及领养和弃养,参考RS分析) histories, err = dsc.controlledHistories(ctx, ds) ... // 遍历history for _, history := range histories { // history.label未标记hash if _, ok := history.Labels[apps.DefaultDaemonSetUniqueLabelKey]; !ok { toUpdate := history.DeepCopy() // 标记hash label=history.Name toUpdate.Labels[apps.DefaultDaemonSetUniqueLabelKey] = toUpdate.Name // 更新 history, err = dsc.kubeClient.AppsV1().ControllerRevisions(ds.Namespace).Update(ctx, toUpdate, metav1.UpdateOptions{}) ... } ... // 匹配ds和history的template found, err = Match(ds, history) ... // template匹配的history作为最新的 if found { currentHistories = append(currentHistories, history) } else { old = append(old, history) } } // 计算新的RV=maxOld+1 currRevision := maxRevision(old) + 1 switch len(currentHistories) { // currHistory不存在,走创建流程 case 0: // 1.先基于ds对象创建 // 2.已存在,对比template内容是否一致 // 3.不一致,获取底层的ds对象,对比ds盐是否有变化(hash冲突) // 4.无变化,更新下盐值,触发下次重试 cur, err = dsc.snapshot(ctx, ds, currRevision) ... // 清理当前history重复项,保留最新版本 default: // 1.筛选最新的history // 2.重新标记其它history对象的Pod hash label // 3.删除其它history对象 cur, err = dsc.dedupCurHistories(ctx, ds, currentHistories) ... // curHistory版本比较旧 if cur.Revision < currRevision { toUpdate := cur.DeepCopy() // 更新RV版本 toUpdate.Revision = currRevision dsc.kubeClient.AppsV1().ControllerRevisions(ds.Namespace).Update(ctx, toUpdate, UpdateOptions{}) ... } } return cur, old, 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注意
ds处于删除状态dsc不会进行任何处理,避免ds子资源状态无法精确计算造成删除时序问题,ds及子资源清理由GC完成
# 1.3.updateHandler
dsc.updateDaemonSet()调用dsc.manage()同步ds状态,或调用dsc.rollingUpdate()执行滚动更新及清理过期的creversion对象。func (dsc *DaemonSetsController) updateDaemonSet(...) error { // 执行创建 dsc.manage(ctx, ds, nodeList, hash) ... // 执行滚动更新(add/del<=0、距上次执行超出5min、exp对象不存在则触发一次) if dsc.expectations.SatisfiedExpectations(key) { switch ds.Spec.UpdateStrategy.Type { ... case apps.RollingUpdateDaemonSetStrategyType: dsc.rollingUpdate(ctx, ds, nodeList, hash) } ... } // 清理过期controllerrevision对象 dsc.cleanupHistory(ctx, ds, old) ... return nil } func (dsc *DaemonSetsController) cleanupHistory(...) error { // ds Pod基于nodeName分组(nodeName基于已调度Pod字段或Affinity获取) // Pod获取涉及领养及弃养 // 包括正在删除或终止Pod nodesToDaemonPods, err := dsc.getNodesToDaemonPods(ctx, ds, true) ... // 计算应该清理多少history toKeep := int(*ds.Spec.RevisionHistoryLimit) toKill := len(old) - toKeep if toKill <= 0 { return nil } // 收集hash label liveHashes := make(map[string]bool) for _, pods := range nodesToDaemonPods { for _, pod := range pods { if hash := pod.Labels[apps.DefaultDaemonSetUniqueLabelKey]; len(hash) > 0 { liveHashes[hash] = true } } } // 基于RV排序history sort.Sort(historiesByRevision(old)) // 根据排序结果遍历history for _, history := range old { // 清理完成 if toKill <= 0 { break } // 跳过hash label活跃的history if hash := history.Labels[apps.DefaultDaemonSetUniqueLabelKey]; liveHashes[hash] { continue } // 删除history对象 dsc.kubeClient.AppsV1().ControllerRevisions(ds.Namespace).Delete(ctx, history.Name, DeleteOptions{}) ... // 剩余待删除数量 toKill-- } 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注意
dsc.cleanupHistory()负责清理超限的history对象,核心逻辑主要在dsc.manage()和dsc.rollingUpdate()
# 1.4.updateStatus
dsc.updateDaemonSetStatus()负责状态计算及更新,主要根据nodeList及Pod推算ds统计信息,交互APIServer更新ds对象状态。func (dsc *DaemonSetsController) updateDaemonSetStatus(...) error { ... // 根据node分组获取ds activePod nodeToDaemonPods, err := dsc.getNodesToDaemonPods(ctx, ds, false) ... // 遍历node列表 for _, node := range nodeList { // node支持运行ds Pod(nodeName/affinity/taints) shouldRun, _ := NodeShouldRunDaemonPod(node, ds) // node已调度ds Pod scheduled := len(nodeToDaemonPods[node.Name]) > 0 // node支持运行ds Pod if shouldRun { // 更新期望值 desiredNumberScheduled++ // 未调度 if !scheduled { continue } // 统计已调度数量 currentNumberScheduled++ // 排序node调度的ds activePod daemonPods, _ := nodeToDaemonPods[node.Name] sort.Sort(podByCreationTimestampAndPhase(daemonPods)) pod := daemonPods[0] // oldest Pod已经ready,统计ready及available数量 if podutil.IsPodReady(pod) { numberReady++ if podutil.IsPodAvailable(pod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}) { numberAvailable++ } } // 获取ds对象generation generation, err := util.GetTemplateGeneration(ds) ... // generation或hash label匹配,更新updatedNumberScheduled数量 if util.IsPodUpdated(pod, hash, generation) { updatedNumberScheduled++ } } else { // node不支持运行Pod if scheduled { // 更新miss数量 numberMisscheduled++ } } } // 计算不可用数量 numberUnavailable := desiredNumberScheduled - numberAvailable // 重试触发ds更新 storeDaemonSetStatus(ctx, dsc.kubeClient.AppsV1().DaemonSets(ds.Namespace), ds, ...) ... // ready与可用数量不一致 if ds.Spec.MinReadySeconds > 0 && numberReady != numberAvailable { // 延迟一轮minReadySeconds时长入队 dsc.enqueueDaemonSetAfter(ds, time.Duration(ds.Spec.MinReadySeconds)*time.Second) } 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注意
dsc.updateDaemonSetStatus()主要根据node运行的ds activePod状态计算副本信息,相关的状态会patch到ds对象
# 2.同步
# 2.1.manage
dsc.manage()会获取nodeToPods分组数据,根据node支持ds Pod运行条件计算toAdd和toDelete集合,基于两者创建或删除Pod。// manage manages the scheduling and running of Pods of ds on nodes. func (dsc *DaemonSetsController) manage(...) error { // ds activePod基于nodeName分组(nodeName基于已调度/Affinity) // activePod获取涉及领养及弃养 // 不包括正在删除或终止Pod nodeToDaemonPods, err := dsc.getNodesToDaemonPods(ctx, ds, false) ... // 遍历node列表 for _, node := range nodeList { // 获取node需创建或删除的Pod toAdd, toDelete := dsc.podsShouldBeOnNode(logger, node, nodeToDaemonPods, ds, hash) // 记录到待新增集合 nodesNeedingDaemonPods = append(nodesNeedingDaemonPods, toAdd...) // 记录到待删除集合 podsToDelete = append(podsToDelete, toDelete...) } // node不存在且nodeName为空的Pod,追加到待删除列表(node不存在但nodeName不为空由PodGC清理) podsToDelete = append(podsToDelete, getUnscheduledPodsWithoutNode(nodeList, nodeToDaemonPods)...) // 执行创建或删除 dsc.syncNodes(ctx, ds, podsToDelete, nodesNeedingDaemonPods, hash) ... return nil } // podsShouldBeOnNode figures out the DaemonSet pods to be created and deleted on the given node. func (dsc *DaemonSetsController) podsShouldBeOnNode(...) (nodesNeedingDaemonPods, podsToDelete []string) { // node应该允许或继续允许ds Pod shouldRun, shouldContinueRunning := NodeShouldRunDaemonPod(node, ds) // node预期或已调度的Pod daemonPods, exists := nodeToDaemonPods[node.Name] switch { // node应该允许但未允许ds Pod case shouldRun && !exists: // 记录node需创建Pod nodesNeedingDaemonPods = append(nodesNeedingDaemonPods, node.Name) // node应该继续运行Pod case shouldContinueRunning: ... for _, pod := range daemonPods { // 正在删除的Pod跳过 if pod.DeletionTimestamp != nil { continue } // Pod创建失败 if pod.Status.Phase == v1.PodFailed { ... // 已处于退避 inBackoff := dsc.failedPodsBackoff.IsInBackOffSinceUpdate(backoffKey, now) if inBackoff { // 计算延迟时间 delay := dsc.failedPodsBackoff.Get(backoffKey) // 延迟入队 dsc.enqueueDaemonSetAfter(ds, delay) continue } // 未处于退避,设置backoff退避时间 dsc.failedPodsBackoff.Next(backoffKey, now) // 记录到待删除Pod podsToDelete = append(podsToDelete, pod.Name) } else if .Status.Phase == v1.PodSucceeded { // Pod走到succeed状态,记录到待删除Pod podsToDelete = append(podsToDelete, pod.Name) } else { // Pending/Running,暂时视为已成功运行 daemonPodsRunning = append(daemonPodsRunning, pod) } } // 未启用maxSurge if !util.AllowsSurge(ds) { // node应该运行的Pod未超出1个,什么都不做 if len(daemonPodsRunning) <= 1 { // There are no excess pods to be pruned, and no pods to create break } // 根据创建时间排序 sort.Sort(podByCreationTimestampAndPhase(daemonPodsRunning)) // 保留最老的Running Pod,其它记录到待删除列表 for i := 1; i < len(daemonPodsRunning); i++ { podsToDelete = append(podsToDelete, daemonPodsRunning[i].Name) } break } // node应该运行但还未运行 if len(daemonPodsRunning) <= 1 { // 记录node需创建Pod if len(daemonPodsRunning) == 0 && shouldRun { // We are surging so we need to have at least one non-deleted pod on the node nodesNeedingDaemonPods = append(nodesNeedingDaemonPods, node.Name) } break } ... // node应该运行的Pod临时有多个(滚动更新),基于创建时间排序 sort.Sort(podByCreationTimestampAndPhase(daemonPodsRunning)) for _, pod := range daemonPodsRunning { // 当前Pod是新版本 if pod.Labels[apps.ControllerRevisionHashLabelKey] == hash { // 记录最老的新Pod if oldestNewPod == nil { oldestNewPod = pod continue } // 记录最老的旧Pod } else { if oldestOldPod == nil { oldestOldPod = pod continue } } // 其它Pod都删掉 podsToDelete = append(podsToDelete, pod.Name) } // 新旧Pod共存 if oldestNewPod != nil && oldestOldPod != nil { switch { // 旧Pod非ready,记录到待删除列表 case !podutil.IsPodReady(oldestOldPod): podsToDelete = append(podsToDelete, oldestOldPod.Name) // 新Pod可用,旧Pod记录到待删除列表 case IsPodAvailable(oldestNewPod, MinReadySeconds, metav1.Time{Time: dsc.backoff.Clock.Now()}): podsToDelete = append(podsToDelete, oldestOldPod.Name) } } // node不应继续运行ds Pod case !shouldContinueRunning && exists: for _, pod := range daemonPods { // 正在删除的可以跳过 if pod.DeletionTimestamp != nil { continue } // 其它的都追加到待删除列表 podsToDelete = append(podsToDelete, pod.Name) } } return nodesNeedingDaemonPods, podsToDelete }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注意
dsc.podsShouldBeOnNode()主要用于收集哪些node需创建Pod,哪些Pod需要删除
# 2.2.syncNode
dsc.syncNodes()向支持运行的node创建Pod,删除不该存在或多余的Pod,每轮调度最大执行250个Pod,更新exp对象维护的期望状态。// syncNodes deletes given pods and creates new ds pods on the given nodes. func (dsc *DaemonSetsController) syncNodes(...) error { ... createDiff := len(nodesNeedingDaemonPods) deleteDiff := len(podsToDelete) // 每轮创建最大限制250 if createDiff > dsc.burstReplicas { createDiff = dsc.burstReplicas } // 每轮删除最大限制250 if deleteDiff > dsc.burstReplicas { deleteDiff = dsc.burstReplicas } // 设置ds exp的add和del期望 dsc.expectations.SetExpectations(dsKey, createDiff, deleteDiff) ... // 获取ds对象的generation generation, err := util.GetTemplateGeneration(ds) ... // 获取Pod模板(提前设置容忍的taint) template := util.CreatePodTemplate(ds.Spec.Template, generation, hash) // 慢创建初始值 batchSize := integer.IntMin(createDiff, controller.SlowStartInitialBatchSize) // 步长为1,2,4,8,...,left for pos := 0; createDiff > pos; batchSize, pos = integer.IntMin(2*batchSize, createDiff-(pos+batchSize)), pos+batchSize { ... // 以batchSize并发创建 for i := pos; i < pos+batchSize; i++ { go func(ix int) { ... // pod定义 podTemplate := template.DeepCopy() // 替换nodeAffinity为当前nodeName podTemplate.Spec.Affinity = util.ReplaceDaemonSetPodNodeNameNodeAffinity( podTemplate.Spec.Affinity, nodesNeedingDaemonPods[ix]) // 创建Pod dsc.podControl.CreatePods(ctx, ds.Namespace, podTemplate, ds, metav1.NewControllerRef(ds, controllerKind)) ... if err != nil { // 出错需手动更新exp对象add期望 dsc.expectations.CreationObserved(dsKey) ... } }(i) } createWait.Wait() // 出错跳过剩余Pod创建 skippedPods := createDiff - (batchSize + pos) if errorCount < len(errCh) && skippedPods > 0 { // 手动更新exp对象的add期望 dsc.expectations.LowerExpectations(dsKey, skippedPods, 0) break } } ... for i := 0; i < deleteDiff; i++ { go func(ix int) { ... // 尝试删除Pod if err := dsc.podControl.DeletePod(ctx, ds.Namespace, podsToDelete[ix], ds); err != nil { // 出错的话手动更新期望 dsc.expectations.DeletionObserved(dsKey) ... } }(i) } deleteWait.Wait() ... return utilerrors.NewAggregate(errors) }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注意
dsc.syncNodes()创建Pod会基于慢加载的策略,删除则直接批量操作
# 3.滚动
# 3.1.desiredCount
dsc.updatedDesiredNodeCounts()根据ds对象的maxSurge和maxUnavailable计算可多创建数量及最大不可用数量,以滚动增删Pod。// updatedDesiredNodeCounts calculates the true number of allowed surge, unavailable or desired scheduled pods. func (dsc *DaemonSetsController) updatedDesiredNodeCounts(...) (int, int, int, error) { ... // 计算期望创建Pod数量 for i := range nodeList { node := nodeList[i] wantToRun, _ := NodeShouldRunDaemonPod(node, ds) if !wantToRun { continue } desiredNumberScheduled++ // 当前节点没有activePod if _, exists := nodeToDaemonPods[node.Name]; !exists { nodeToDaemonPods[node.Name] = nil } } // maxUnavailable=desired*unavailable maxUnavailable, err := util.UnavailableCount(ds, desiredNumberScheduled) ... // maxSurge=desired*surge maxSurge, err := util.SurgeCount(ds, desiredNumberScheduled) ... // 未设置上下限,maxUnavailable初始化为1 if desiredNumberScheduled > 0 && maxUnavailable == 0 && maxSurge == 0 { maxUnavailable = 1 } return maxSurge, maxUnavailable, desiredNumberScheduled, 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注意
maxSurge可以不设置,maxUnavailable默认为1,打破滚动更新期间由于错误的配置导致每次伸缩总为0的局面
# 3.2.rollUpdate
dsc.rollingUpdate()负责滚动更新nodePod,根据条件计算toAdd和toDelete列表,调用syncNode执行newPod创建及oldPod删除。// rollingUpdate identifies the set of old pods to delete, or additional pods to create on nodes. func (dsc *DaemonSetsController) rollingUpdate(...) error { // 获取activePod nodeToDaemonPods, err := dsc.getNodesToDaemonPods(ctx, ds, false) ... // 计算maxSurge maxUnavailable desired maxSurge, maxUnavailable, desired, err := dsc.updatedDesiredNodeCounts(ctx, ds, nodeList, nodeToDaemonPods) ... // 未设置上限,仅执行删除,创建动作由manage流程补充 if maxSurge == 0 { ... for nodeName, pods := range nodeToDaemonPods { // 获取newPod和oldPod newPod, oldPod, ok := findUpdatedPodsOnNode(ds, pods, hash) // newPod或oldPod存在多个,manage流程会处理 if !ok { // 更新unavailable数量 numUnavailable++ continue } switch { // newPod和oldPod同时不存在或存在,manage流程会处理 case oldPod == nil && newPod == nil, oldPod != nil && newPod != nil: // 更新unavailable数量 numUnavailable++ // 仅有newPod case newPod != nil: // newPod不可用 if !podutil.IsPodAvailable(newPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}) { // 更新unavailable数量 numUnavailable++ } default: // oldPod作为更新候选 switch { // oldPod不是可用 case !podutil.IsPodAvailable(oldPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}): ... // 作为替代候选 allowedReplacementPods = append(allowedReplacementPods, oldPod.Name) // 不可用数量超出maxUnavailable,不再考虑其它候选 case numUnavailable >= maxUnavailable: // no point considering any other candidates continue default: // oldPod作为删除候选 candidatePodsToDelete = append(candidatePodsToDelete, oldPod.Name) } } } // 计算剩余最大不可用数量 remainingUnavailable := maxUnavailable - numUnavailable if remainingUnavailable < 0 { remainingUnavailable = 0 } // remainingUnavailable超出候选删除数量,重置避免越界 if max := len(candidatePodsToDelete); remainingUnavailable > max { remainingUnavailable = max } // 合并到replacePods oldPodsToDelete := append(allowedReplacementPods, candidatePodsToDelete[:remainingUnavailable]...) // 执行syncNodes删除多出Pod return dsc.syncNodes(ctx, ds, oldPodsToDelete, nil, hash) } ... // 设置maxSurge for nodeName, pods := range nodeToDaemonPods { // 获取newPod和oldPod newPod, oldPod, ok := findUpdatedPodsOnNode(ds, pods, hash) // newPod或oldPod超出1个 if !ok { // 更新numSurge numSurge++ continue } // oldPod可用 if oldPod != nil { if podutil.IsPodAvailable(oldPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}) { numAvailable++ } // newPod可用 } else if newPod != nil { if podutil.IsPodAvailable(newPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}) { numAvailable++ } } switch { // oldPod为空什么都不做(没有替换目标) case oldPod == nil: // we don't need to do anything to this node, the manage loop will handle it // newPod为空 case newPod == nil: // this is a surge candidate switch { // oldPod不可用 case !podutil.IsPodAvailable(oldPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}): // 获取node node, err := dsc.nodeLister.Get(nodeName) ... // node已不支持运行Pod if shouldRun, _ := NodeShouldRunDaemonPod(node, ds); !shouldRun { // oldPod作为删除候选 oldPodsToDelete = append(oldPodsToDelete, oldPod.Name) continue } ... // node记录需要newPod allowedNewNodes = append(allowedNewNodes, nodeName) // oldPod可用 default: // 获取node node, err := dsc.nodeLister.Get(nodeName) ... // node已不支持运行Pod if shouldRun, _ := NodeShouldRunDaemonPod(node, ds); !shouldRun { // oldPod作为删除候选 shouldNotRunPodsToDelete = append(shouldNotRunPodsToDelete, oldPod.Name) continue } // 达到最大数量,不再新建 if numSurge >= maxSurge { // no point considering any other candidates continue } ... // node作为新建候选 candidateNewNodes = append(candidateNewNodes, nodeName) } // oldPod和newPod均不为空 default: // newPod不可用 if !podutil.IsPodAvailable(newPod, ds.Spec.MinReadySeconds, metav1.Time{Time: now}) { // 更新numSurge numSurge++ continue } // newPod可用,oldPod作为删除候选 oldPodsToDelete = append(oldPodsToDelete, oldPod.Name) } } // 统计剩余可创建数量 remainingSurge := maxSurge - numSurge // 计算可安全删除的Pod(多出的) if deletable := numAvailable - desiredNumberScheduled; deletablePodsNumber > 0 { // 最多只能删除不该运行Pod的node数量 if shouldNotRun := len(shouldNotRunPodsToDelete); deletable > shouldNotRun { deletablePodsNumber = shouldNotRunPodsToDeleteNumber } // 删除条件检查 for _, podToDeleteName := range shouldNotRunPodsToDelete[:deletablePodsNumber] { // 获取待删除Pod podToDelete, err := dsc.podLister.Pods(ds.Namespace).Get(podToDeleteName) // 跳过不存在的 if err != nil { if errors.IsNotFound(err) { continue } return err } // 存在的作为删除候选 oldPodsToDelete = append(oldPodsToDelete, podToDeleteName) } } // 计算剩余可多创建的Pod数量 if remainingSurge < 0 { remainingSurge = 0 } if max := len(candidateNewNodes); remainingSurge > max { remainingSurge = max } // 最大创建remainingSurge个Pod newNodesToCreate := append(allowedNewNodes, candidateNewNodes[:remainingSurge]...) // 执行创建与删除 return dsc.syncNodes(ctx, ds, oldPodsToDelete, newNodesToCreate, hash) }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注意
dsc.rollingUpdate()会基于ds条件进行滚动更新,maxSurge未设置仅执行缩容,newPod由manage流程补充