statefulset
# 1.简介
# 1.1.定义
statefulset用于管理有状态应用的工作负载,向Pod提供稳定的身份及存储,限制创建、删除及更新执行顺序,确保有状态应用的稳定性。
注意
statefulset用于满足有状态的服务或中间件,诸如mysql、zookeeper或kafka,一般期望存储不希望随服务重启而消失
# 1.2.原理
statefulset基于informer监听statefulset/pod资源变化,缓存controllerRevision/pvc资源,消费workqueue执行sync同步。注意
workqueue是一个先入先出的队列,由item切片、dirty map和processing map构成
# 2.入口
# 2.1.start
startStatefulSetController()是ssc的启动入口,调用NewStatefulSetController()实例化ssc对象及执行Run()激活调谐。// NewStatefulSetController creates a new statefulset controller. func NewStatefulSetController(...) *StatefulSetController { ... // 实例化ssc ssc := &StatefulSetController{ kubeClient: kubeClient, control: NewDefaultStatefulSetControl( NewStatefulPodControl(kubeClient, podInformer.Lister(), pvcInformer.Lister(), recorder), NewRealStatefulSetStatusUpdater(kubeClient, setInformer.Lister()), history.NewHistory(kubeClient, revInformer.Lister()), recorder, ), ... queue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "ss"), podControl: controller.RealPodControl{KubeClient: kubeClient, Recorder: recorder}, ... } // pod informer回调及缓存 podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ // lookup the statefulset and enqueue AddFunc: func(obj interface{}) { ssc.addPod(logger, obj) }, // lookup current and old statefulset if labels changed UpdateFunc: func(oldObj, newObj interface{}) { ssc.updatePod(logger, oldObj, newObj) }, // lookup statefulset accounting for deletion tombstones DeleteFunc: func(obj interface{}) { ssc.deletePod(logger, obj) }, }) ssc.podLister = podInformer.Lister() ... // stateful informer回调及缓存 setInformer.Informer().AddEventHandler( cache.ResourceEventHandlerFuncs{ AddFunc: ssc.enqueueStatefulSet, UpdateFunc: func(old, cur interface{}) { ... ssc.enqueueStatefulSet(cur) }, DeleteFunc: ssc.enqueueStatefulSet, }, ) ssc.setLister = setInformer.Lister() ... return ssc } func startStatefulSetController(...) (controller.Interface, bool, error) { go statefulset.NewStatefulSetController( ctx, controllerContext.InformerFactory.Core().V1().Pods(), controllerContext.InformerFactory.Apps().V1().StatefulSets(), controllerContext.InformerFactory.Core().V1().PersistentVolumeClaims(), controllerContext.InformerFactory.Apps().V1().ControllerRevisions(), controllerContext.ClientBuilder.ClientOrDie("statefulset-controller"), ).Run(ctx, int(controllerContext.ComponentConfig.StatefulSetController.ConcurrentStatefulSetSyncs)) 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
注意
podInformer回调会基于Pod匹配所属的statefulset,基于owner/label selector匹配
# 2.2.podInformer
podInformer会监听Add/Update/DEL事件,基于owner/label selector匹配所属的statefulset入队任务,供sync同步处理。// addPod adds the statefulset for the pod to the sync queue func (ssc *StatefulSetController) addPod(logger klog.Logger, obj interface{}) { pod := obj.(*v1.Pod) // Pod删除分支 if pod.DeletionTimestamp != nil { ssc.deletePod(logger, pod) return } // owner存在 if controllerRef := metav1.GetControllerOf(pod); controllerRef != nil { // 解析owner及入队 set := ssc.resolveControllerRef(pod.Namespace, controllerRef) ... ssc.enqueueStatefulSet(set) return } // 基于label selector匹配 sets := ssc.getStatefulSetsForPod(pod) ... // 入队所属的statefulset for _, set := range sets { ssc.enqueueStatefulSet(set) } } // updatePod adds the statefulset for the current and old pods to the sync queue. func (ssc *StatefulSetController) updatePod(logger klog.Logger, old, cur interface{}) { ... // resourceVersion无变化 if curPod.ResourceVersion == oldPod.ResourceVersion { return } // label差异 labelChanged := !reflect.DeepEqual(curPod.Labels, oldPod.Labels) ... // ownerRef差异 controllerRefChanged := !reflect.DeepEqual(curControllerRef, oldControllerRef) // ownerRef变化及oldOwner存在 if controllerRefChanged && oldControllerRef != nil { // 入队oldOwner if set := ssc.resolveControllerRef(oldPod.Namespace, oldControllerRef); set != nil { ssc.enqueueStatefulSet(set) } } // curOwner存在 if curControllerRef != nil { // 入队curOwner set := ssc.resolveControllerRef(curPod.Namespace, curControllerRef) ... ssc.enqueueStatefulSet(set) // unready-->ready且设置就绪时间 if !podutil.IsPodReady(oldPod) && podutil.IsPodReady(curPod) && set.Spec.MinReadySeconds > 0 { // curOwner延迟入队 ssc.enqueueSSAfter(set, (time.Duration(set.Spec.MinReadySeconds)*time.Second)+time.Second) } return } // label差异或ownerRef差异 if labelChanged || controllerRefChanged { // 基于label selector匹配 sets := ssc.getStatefulSetsForPod(curPod) ... // 入队statefulset for _, set := range sets { ssc.enqueueStatefulSet(set) } } } // deletePod enqueues the statefulset for the pod accounting for deletion tombstones. func (ssc *StatefulSetController) deletePod(logger klog.Logger, obj interface{}) { pod, ok := obj.(*v1.Pod) ... // 解析owner controllerRef := metav1.GetControllerOf(pod) ... set := ssc.resolveControllerRef(pod.Namespace, controllerRef) ... // 入队statefulset ssc.enqueueStatefulSet(set) }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
注意
podInformer监听的任务推入workqueue由sync流程消费
# 2.3.runworker
ssc.Run()会阻塞至informer同步完成,激活5个worker执行sync()同步,sync()推出workqueue任务进行statefulset同步。// Run runs the statefulset controller. func (ssc *StatefulSetController) Run(ctx context.Context, workers int) { ... defer ssc.queue.ShutDown() // informer同步阻塞 if !cache.WaitForNamedCacheSync("stateful set", ctx.Done(), ssc.podListerSynced, ssc.setListerSynced, ssc.pvcListerSynced, ssc.revListerSynced) { return } // 默认激活5个worker for i := 0; i < workers; i++ { go wait.UntilWithContext(ctx, ssc.worker, time.Second) } <-ctx.Done() } // worker runs a worker goroutine that invokes processNextWorkItem until the controller's queue is closed func (ssc *StatefulSetController) worker(ctx context.Context) { for ssc.processNextWorkItem(ctx) { } } // processNextWorkItem dequeues items, processes them, and marks them done. func (ssc *StatefulSetController) processNextWorkItem(ctx context.Context) bool { key, quit := ssc.queue.Get() ... defer ssc.queue.Done(key) // 执行sync同步 if err := ssc.sync(ctx, key.(string)); err != nil { ... ssc.queue.AddRateLimited(key) } else { ssc.queue.Forget(key) } 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
注意
worker执行sync()同步,sync()本质上调用syncStatefulset()进行同步调谐
# 2.4.sync
ssc.sync()会领养或弃养controllerReversion/Pod对象,注入或释放子资源的ownerRef,执行syncStatefulSet()同步调谐状态。// sync syncs the given statefulset. func (ssc *StatefulSetController) sync(ctx context.Context, key string) error { ... set, err := ssc.setLister.StatefulSets(namespace).Get(name) ... // 获取label selector selector, err := metav1.LabelSelectorAsSelector(set.Spec.Selector) ... // 领养cr ssc.adoptOrphanRevisions(ctx, set) ... // 领养或弃养Pod pods, err := ssc.getPodsForStatefulSet(ctx, set, selector) ... // 状态同步 return ssc.syncStatefulSet(ctx, set, pods) } // getPodsForStatefulSet returns the Pods that a given StatefulSet should manage. func (ssc *StatefulSetController) getPodsForStatefulSet(...) ([]*v1.Pod, error) { // 获取Pods pods, err := ssc.podLister.Pods(set.Namespace).List(labels.Everything()) ... // Pod截取的setName与set一致 filter := func(pod *v1.Pod) bool { // Only claim if it matches our StatefulSet name. Otherwise release/ignore. return isMemberOf(set, pod) } ... // 执行领养 return cm.ClaimPods(ctx, pods, filter) } // ClaimPods tries to take ownership of a list of Pods. func (m *PodControllerRefManager) ClaimPods(...) bool) ([]*v1.Pod, error) { ... // label selector匹配及name一致 match := func(obj metav1.Object) bool { pod := obj.(*v1.Pod) // Check selector first so filters only run on potentially matching Pods. if !m.Selector.Matches(labels.Set(pod.Labels)) { return false } for _, filter := range filters { if !filter(pod) { return false } } return true } // 注入ownerRef领养 adopt := func(ctx context.Context, obj metav1.Object) error { return m.AdoptPod(ctx, obj.(*v1.Pod)) } // 摘除ownerRef弃养 release := func(ctx context.Context, obj metav1.Object) error { return m.ReleasePod(ctx, obj.(*v1.Pod)) } // 基于match+adopt领养或match+release弃养 for _, pod := range pods { ok, err := m.ClaimObject(ctx, pod, match, adopt, release) ... if ok { claimed = append(claimed, pod) } } return claimed, utilerrors.NewAggregate(errlist) }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
注意
sync()主要做一些准备工作,涉及controllerReversion/Pod领养,真正的同步由syncStatefulset()实现
# 3.同步
# 3.1.syncState
ssc.syncStatefulset()会触发statefulset同步,更新status状态,基于revisionHistoryLimit限制清理controllerrevision。// syncStatefulSet syncs a tuple of (statefulset, []*v1.Pod). func (ssc *StatefulSetController) syncStatefulSet(...) error { ... // 执行同步 status, err = ssc.control.UpdateStatefulSet(ctx, set, pods) ... // 设置minReady且副本还未全部可用,延迟入队 if set.Spec.MinReadySeconds > 0 && status != nil && status.AvailableReplicas != *set.Spec.Replicas { ssc.enqueueSSAfter(set, time.Duration(set.Spec.MinReadySeconds)*time.Second) } return nil } // UpdateStatefulSet executes the core logic loop for a stateful set. func (ssc *defaultStatefulSetControl) UpdateStatefulSet(...) (*apps.StatefulSetStatus, error) { ... // 基于selector+owner获取匹配的reversion revisions, err := ssc.ListRevisions(set) ... // 由小到大排序(reversion小的-->新建的-->name字典序靠前的) history.SortControllerRevisions(revisions) // 核心同步逻辑 currentRevision, updateRevision, status, err := ssc.performUpdate(ctx, set, pods, revisions) ... // 清理过期reversion return status, ssc.truncateHistory(set, pods, revisions, currentRevision, updateRevision) } // truncateHistory truncates any non-live ControllerRevisions in revisions from set's history. func (ssc *defaultStatefulSetControl) truncateHistory(...) error { ... // 标记curReversion活跃 if current != nil { live[current.Name] = true } // 标记updateReversion活跃 if update != nil { live[update.Name] = true } // 标记Pod使用的reversion活跃 for i := range pods { live[getPodRevision(pods[i])] = true } // 遍历排序后的reversion for i := range revisions { // 不活跃的作为history if !live[revisions[i].Name] { history = append(history, revisions[i]) } } ... // 未超出限制 if historyLen <= historyLimit { return nil } // 清理超出限制的reversion history = history[:(historyLen - historyLimit)] for i := 0; i < len(history); i++ { ssc.controllerHistory.DeleteControllerRevision(history[i]) ... } 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
补充
ssc.performUpdate()是真正的同步模块,会执行同步及状态更新
# 3.2.performUpdate
ssc.performUpdate()获取statefulset reversion情况,执行具体的更新同步,statefulset运行状态会Patch到status更新。func (ssc *defaultStatefulSetControl) performUpdate(...) (...) { ... // 获取当前版本、目标版本及冲突计数 currentRevision, updateRevision, collisionCount, err := ssc.getStatefulSetRevisions(set, revisions) ... // 执行同步 currentStatus, _ = ssc.updateStatefulSet(ctx, set, currentRevision, updateRevision, collisionCount, pods) ... // 更新状态 ssc.updateStatefulSetStatus(ctx, set, currentStatus) ... return currentRevision, updateRevision, currentStatus, nil } // getStatefulSetRevisions returns the current and update ControllerRevisions for set. func (ssc *defaultStatefulSetControl) getStatefulSetRevisions(...) (...) { ... // 排序(reversion↓-->createTime↑-->name↓) history.SortControllerRevisions(revisions) ... // 生成newReversion草稿(基于内容+冲突计数的hash生成名称) updateRevision, err := newRevision(set, nextRevision(revisions), &collisionCount) ... // 基于hash+内容查找相同的历史reversion equalRevisions := history.FindEqualRevisions(revisions, updateRevision) ... // 已有完全相同的reversion且最后一个 if equalCount > 0 && history.EqualRevision(revisions[revisionCount-1], equalRevisions[equalCount-1]) { // 无需创建,直接用最后一个历史版本 updateRevision = revisions[revisionCount-1] // 已有完全相同的reversion且不是最后一个 } else if equalCount > 0 { // 回滚行为,直接更新oldReversion对象,reversionNum会增加 updateRevision, err = ssc.controllerHistory.UpdateControllerRevision(equalRevisions[equalCount-1], updateRevision.Revision) ... // 没有相同版本 } else { // 创建newReversion updateRevision, err = ssc.controllerHistory.CreateControllerRevision(set, updateRevision, &collisionCount) ... } // 查找当前正在用的reversion for i := range revisions { if revisions[i].Name == set.Status.CurrentRevision { currentRevision = revisions[i] break } } // 没有curReversion,取updateReversion if currentRevision == nil { currentRevision = updateRevision } return currentRevision, updateRevision, collisionCount, 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
补充
ssc.performUpdate()会生成newReversion对象,基于curReversion和newReversion执行调谐及更新statefulset状态