daemonset
# 1.简介
# 1.1.定义
daemonset类似Linux上的守护进程,会向每个node部署Pod副本,确保部署的Pod与节点数严格一致,一般用于日志组件或网络插件部署。// dscontroller is responsible for synchronizing ds objects stored in the system with actual running pods. type DaemonSetsController struct { ... burstReplicas int // 批操作上限 syncHandler func(ctx context.Context, dsKey string) error // 同步器 enqueueDaemonSet func(ds *apps.DaemonSet) // 入队器 // A TTLCache of pod creates/deletes each ds expects to see expectations controller.ControllerExpectationsInterface // TTLCache(记录期望创建/删除Pod数) dsLister appslisters.DaemonSetLister // ds缓存 ... historyLister appslisters.ControllerRevisionLister // 历史版本缓存 ... podLister corelisters.PodLister // Pod缓存 ... nodeLister corelisters.NodeLister // node缓存 ... queue workqueue.RateLimitingInterface // ds queue failedPodsBackoff *flowcontrol.Backoff // 退避器(流控) }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注意
daemonset会基于node下发Pod,node删除会触发Pod回收,node加入会触发Pod创建
# 1.2.原理
daemonset会监听节点的加入或删除事件,匹配的node会创建Pod,清理不该存在的Pod,根据daemonset策略增删或滚动更新Pod。补充
worker默认激活2个,并发处理不同ds对象的同步,核心逻辑都放在ds.syncHandler()
# 2.分析
# 2.1.start
startDaemonSetController()负责实例化dsc及调用dsc.Run()激活同步,监听node/creversion/ds/pod资源变化调整Pod行为。// NewDaemonSetsController creates a new DaemonSetsController func NewDaemonSetsController(...) (*DaemonSetsController, error) { ... // 实例化ds controller dsc := &DaemonSetsController{ ... burstReplicas: 250, expectations: controller.NewControllerExpectations(), queue: workqueue.NewNamedRateLimitingQueue(DefaultControllerRateLimiter(), "daemonset"), } // dsInformer回调 daemonSetInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { dsc.addDaemonset(logger, obj) }, UpdateFunc: func(oldObj, newObj interface{}) { dsc.updateDaemonset(logger, oldObj, newObj) }, DeleteFunc: func(obj interface{}) { dsc.deleteDaemonset(logger, obj) }, }) dsc.dsLister = daemonSetInformer.Lister() ... // creversionInformer回调 historyInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { dsc.addHistory(logger, obj) }, UpdateFunc: func(oldObj, newObj interface{}) { dsc.updateHistory(logger, oldObj, newObj) }, DeleteFunc: func(obj interface{}) { dsc.deleteHistory(logger, obj) }, }) dsc.historyLister = historyInformer.Lister() ... // podInformer回调 podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { dsc.addPod(logger, obj) }, UpdateFunc: func(oldObj, newObj interface{}) { dsc.updatePod(logger, oldObj, newObj) }, DeleteFunc: func(obj interface{}) { dsc.deletePod(logger, obj) }, }) dsc.podLister = podInformer.Lister() ... // nodeInformer回调 nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { dsc.addNode(logger, obj) }, UpdateFunc: func(oldObj, newObj interface{}) { dsc.updateNode(logger, oldObj, newObj) }, }, ) dsc.nodeStoreSynced = nodeInformer.Informer().HasSynced ... dsc.syncHandler = dsc.syncDaemonSet dsc.enqueueDaemonSet = dsc.enqueue // 退避器[1s,15min] dsc.failedPodsBackoff = failedPodsBackoff return dsc, nil } func startDaemonSetController(...) (controller.Interface, bool, error) { // 实例化dsc dsc, err := daemon.NewDaemonSetsController(...) ... // 默认为2 go dsc.Run(ctx, int(controllerContext.ComponentConfig.DaemonSetController.ConcurrentDaemonSetSyncs)) 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
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注意
daemonset正常同步依赖informer回调推入workqueue的ds对象,下面会分析几个informer如何将关联ds推入workqueue
# 2.2.dsInformer
dsInformer会监听ds对象变化,根据Add/Update/Del事件将关联的ds对象推入workqueue,ds对象重建会触发oldDS回收逻辑。// 创建/刚启动事件 func (dsc *DaemonSetsController) addDaemonset(logger klog.Logger, obj interface{}) { ds := obj.(*apps.DaemonSet) dsc.enqueueDaemonSet(ds) } // 更新事件 func (dsc *DaemonSetsController) updateDaemonset(logger klog.Logger, cur, old interface{}) { oldDS := old.(*apps.DaemonSet) curDS := cur.(*apps.DaemonSet) // ds对象重建 if curDS.UID != oldDS.UID { ... // oldDs对象清理 dsc.deleteDaemonset(logger, cache.DeletedFinalStateUnknown{ Key: key, Obj: oldDS, }) } // curDs对象入队 dsc.enqueueDaemonSet(curDS) } // 删除事件 func (dsc *DaemonSetsController) deleteDaemonset(logger klog.Logger, obj interface{}) { ds, ok := obj.(*apps.DaemonSet) ... key, err := controller.KeyFunc(ds) ... // 清理exp期望状态 dsc.expectations.DeleteExpectations(key) // 入队 dsc.queue.Add(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注意
dsInformer基于ds对象变化触发推入workqueue,由于监听的就是自身对象,相对简单一些
# 2.3.crInformer
crInfomer指的是controllerrevisions对象,记录ds历史版本,类似deployment所属的replicaset对象,用于提供版本回退。// addHistory enqueues the ds that manages a ControllerRevision when the cr is created or restarted. func (dsc *DaemonSetsController) addHistory(logger klog.Logger, obj interface{}) { history := obj.(*apps.ControllerRevision) // cr对象正在删除 if history.DeletionTimestamp != nil { // 走DEL流程 dsc.deleteHistory(logger, history) return } // owner不为空(什么都没做,用意是什么?) if controllerRef := metav1.GetControllerOf(history); controllerRef != nil { ds := dsc.resolveControllerRef(history.Namespace, controllerRef) ... return } // orphan cr对象,基于selector匹配ds daemonSets := dsc.getDaemonSetsForHistory(logger, history) if len(daemonSets) == 0 { return } // 匹配的ds入队 for _, ds := range daemonSets { dsc.enqueueDaemonSet(ds) } } // updateHistory figures out what ds manage a cr when the cr is updated and wake them up. func (dsc *DaemonSetsController) updateHistory(logger klog.Logger, old, cur interface{}) { ... // rv版本无变化 if curHistory.ResourceVersion == oldHistory.ResourceVersion { return } ... // owner变化及oldOwner存在 if controllerRefChanged && oldControllerRef != nil { // 所属owner ds入队 if ds := dsc.resolveControllerRef(oldHistory.Namespace, oldControllerRef); ds != nil { dsc.enqueueDaemonSet(ds) } } // curOwner存在 if curControllerRef != nil { // 解析curOwner ds对象 ds := dsc.resolveControllerRef(curHistory.Namespace, curControllerRef) if ds == nil { return } // 入队 dsc.enqueueDaemonSet(ds) return } ... // label变化或owner变化(owner-->无主) if labelChanged || controllerRefChanged { // 获取curCR匹配的ds daemonSets := dsc.getDaemonSetsForHistory(logger, curHistory) if len(daemonSets) == 0 { return } // 匹配的ds入队 for _, ds := range daemonSets { dsc.enqueueDaemonSet(ds) } } } // deleteHistory enqueues the ds that manages a cr when the cr is deleted. func (dsc *DaemonSetsController) deleteHistory(logger klog.Logger, obj interface{}) { history, ok := obj.(*apps.ControllerRevision) ... // owner为空 if controllerRef == nil { // No controller should care about orphans being deleted. return } // 解析owner ds ds := dsc.resolveControllerRef(history.Namespace, controllerRef) if ds == nil { return } // owner ds入队 dsc.enqueueDaemonSet(ds) }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注意
crInformer监听controllerrevisions对象变化,入队关联的owner/labelMatch ds对象
# 2.4.podInformer
podInformer之前分析的控制器也出现过,还是基于owner/label匹配daemonset对象,将owner/match ds对象推入workqueue。func (dsc *DaemonSetsController) addPod(logger klog.Logger, obj interface{}) { pod := obj.(*v1.Pod) // pod正在删除 if pod.DeletionTimestamp != nil { // 走DEL流程 dsc.deletePod(logger, pod) return } // owner存在 if controllerRef := metav1.GetControllerOf(pod); controllerRef != nil { ds := dsc.resolveControllerRef(pod.Namespace, controllerRef) ... dsKey, err := controller.KeyFunc(ds) ... // exp的add-1 dsc.expectations.CreationObserved(dsKey) // ds入队 dsc.enqueueDaemonSet(ds) return } // orphan pod,基于label selector匹配ds dss := dsc.getDaemonSetsForPod(pod) ... // 匹配的ds入队 for _, ds := range dss { dsc.enqueueDaemonSet(ds) } } // When a pod is updated, figure out what sets manage it and wake them up. func (dsc *DaemonSetsController) updatePod(logger klog.Logger, old, cur interface{}) { ... // rv无变化 if curPod.ResourceVersion == oldPod.ResourceVersion { // Periodic resync will send update events for all known pods. // Two different versions of the same pod will always have different RVs. return } // curPod正在删除 if curPod.DeletionTimestamp != nil { // 走DEL流程 dsc.deletePod(logger, curPod) return } ... // owner变化及oldOwner存在 if controllerRefChanged && oldControllerRef != nil { // old owner ds入队 if ds := dsc.resolveControllerRef(oldPod.Namespace, oldControllerRef); ds != nil { dsc.enqueueDaemonSet(ds) } } // curOwner存在 if curControllerRef != nil { ds := dsc.resolveControllerRef(curPod.Namespace, curControllerRef) ... // cur owner ds入队 dsc.enqueueDaemonSet(ds) ... // Pod刚变为ready及设置MinReadySeconds if changedToReady && ds.Spec.MinReadySeconds > 0 { // 延迟一个MinReadySeconds再入队检查ready dsc.enqueueDaemonSetAfter(ds, (time.Duration(ds.Spec.MinReadySeconds)*time.Second)+time.Second) } return } // orphan Pod基于label selector匹配ds dss := dsc.getDaemonSetsForPod(curPod) ... // label或owner变化 if labelChanged || controllerRefChanged { // 匹配的ds入队 for _, ds := range dss { dsc.enqueueDaemonSet(ds) } } } func (dsc *DaemonSetsController) deletePod(logger klog.Logger, obj interface{}) { pod, ok := obj.(*v1.Pod) ... // owner为空 controllerRef := metav1.GetControllerOf(pod) if controllerRef == nil { // No controller should care about orphans being deleted. return } // 获取owner ds ds := dsc.resolveControllerRef(pod.Namespace, controllerRef) if ds == nil { return } dsKey, err := controller.KeyFunc(ds) ... // exp del-1 dsc.expectations.DeletionObserved(dsKey) // owner ds入队 dsc.enqueueDaemonSet(ds) }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注意
Pod存在ownerRef,将owner ds推入workqueue,部分情况会基于label selector为无主Pod匹配owner ds
# 2.5.nodeInformer
nodeInformer监听node资源变化,基于daemonset检测node是否可以运行ds Pod,满足会将ds推入workqueue,否则不处理。func (dsc *DaemonSetsController) addNode(logger klog.Logger, obj interface{}) { // 获取所有ds对象 dsList, err := dsc.dsLister.List(labels.Everything()) ... node := obj.(*v1.Node) for _, ds := range dsList { // 检测node是否可以运行Pod(nodeName/nodeAffinity/taints) if shouldRun, _ := NodeShouldRunDaemonPod(node, ds); shouldRun { // 满足则入队ds dsc.enqueueDaemonSet(ds) } } } func (dsc *DaemonSetsController) updateNode(logger klog.Logger, old, cur interface{}) { ... // condition无变化或node其它字段无变化 if shouldIgnoreNodeUpdate(*oldNode, *curNode) { return } // 获取所有ds dsList, err := dsc.dsLister.List(labels.Everything()) ... for _, ds := range dsList { // oldNode是否可以允许Pod(nodeName/nodeAffinity/taints) oldShouldRun, oldShouldContinueRunning := NodeShouldRunDaemonPod(oldNode, ds) // curNode是否可以允许Pod(nodeName/nodeAffinity/taints) currentShouldRun, currentShouldContinueRunning := NodeShouldRunDaemonPod(curNode, ds) // 运行条件变化,ds入队 if (oldShouldRun != currentShouldRun) || (oldShouldContinueRunning != currentShouldContinueRunning) { dsc.enqueueDaemonSet(ds) } } }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注意
nodeInformer更关注的是node是否可以运行ds Pod,由nodeName/nodeAffinity/taints三个维度检查