endpoint
# 1.简介
# 1.1.service
service是相同功能的pod抽象,向pod提供统一的访问入口。借助service,应用可以实现服务发现与负载均衡,做到应用的零宕机升级。pod与service基于标签匹配组成endpoints,由kubeproxy生成iptables/ipvs规则,劫持及将service流量负载均衡到endpoint。
注意
service controller已迁移到cloud-provider项目,仅监听处理loadbalancer类型service,其它service仅由apiserver回填部分数据
# 1.2.endpoints
endpoints定义service的backend后端,用于描述可供service访问的pod地址,可以实现后端应用的动态发现及注册。本质上,访问svc的流量实际会负载到endpoints定义的某个后端,endpoints关联的后端基于标签选择进行匹配。
注意
service关联服务膨胀时,endpoint对象可能非常巨大,最终超出etcd存储的对象大小限制1.5M,同时频繁的扩容与更新服务会导致会造成完整的endpoint资源多次分发,造成网络性能开销。
# 1.3.syncLoop
endpoints controller持续监听service/pod资源的变更事件及加入队列后供主循环处理,用以组织endpoints后端池,触发kubeproxy规则同步。
# 2.epcontroller
# 2.1.newEndpointController
endpoint controller是kube-controller-manager控制器之一,用于管理endpoints资源对象的生命周期。service变化时,将watch到的事件放入queue,供syncService()取出查询关联的pod列表,完成endpoints对象的生命周期管理。// Controller manages selector-based service endpoints. type Controller struct { client clientset.Interface ... // service资源列表 serviceLister corelisters.ServiceLister ... // pod资源列表 podLister corelisters.PodLister ... // endpoint资源列表 endpointsLister corelisters.EndpointsLister ... // 暂存变化的service queue workqueue.RateLimitingInterface // worker协程循环周期 workerLoopPeriod time.Duration // service标签选择器缓存,避免频繁构造销毁造成cpu消耗 serviceSelectorCache *endpointutil.ServiceSelectorCache } // NewEndpointController returns a new *Controller. func NewEndpointController(...) *Controller { ... e := &Controller{ client: client, queue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "endpoint"), // 1000ms workerLoopPeriod: time.Second, } // service监听 serviceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: e.onServiceUpdate, UpdateFunc: func(old, cur interface{}) { e.onServiceUpdate(cur) }, DeleteFunc: e.onServiceDelete, }) e.serviceLister = serviceInformer.Lister() e.servicesSynced = serviceInformer.Informer().HasSynced // pod监听 podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: e.addPod, UpdateFunc: e.updatePod, DeleteFunc: e.deletePod, }) e.podLister = podInformer.Lister() e.podsSynced = podInformer.Informer().HasSynced // endpoint监听 endpointsInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ DeleteFunc: e.onEndpointsDelete, }) e.endpointsLister = endpointsInformer.Lister() e.endpointsSynced = endpointsInformer.Informer().HasSynced ... // endpoint批量更新时间 e.endpointUpdatesBatchPeriod = endpointUpdatesBatchPeriod // 缓存service标签选择器 e.serviceSelectorCache = endpointutil.NewServiceSelectorCache() return e }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
注意
service/pod监听的事件均会将关联service放入queue供worker创建或更新endpoints,endpoint监听的事件会触发删除
# 2.2.startEndpointController
startEndpointController()会启动协程初始化endpoint controller,调用Run()创建一定数量worker处理queue缓存的service,创建或更新endpoints。func startEndpointController(ctx context.Context, controllerCtx ControllerContext) (controller.Interface, bool, error) { go endpointcontroller.NewEndpointController( ... ).Run(ctx, int(controllerCtx.ComponentConfig.EndpointController.ConcurrentEndpointSyncs)) return nil, true, nil } // Run will not return until stopCh is closed.workers determines how many endpoints will be handled in parallel. func (e *Controller) Run(ctx context.Context, workers int) { ... defer e.queue.ShutDown() // 等待资源同步完成 if !cache.WaitForNamedCacheSync("endpoint", ctx.Done(), e.podsSynced, e.servicesSynced, e.endpointsSynced) { return } // 启动多个worker处理同步(10) for i := 0; i < workers; i++ { // 间隔1s周期调度 go wait.UntilWithContext(ctx, e.worker, e.workerLoopPeriod) } go func() { ... // 清理孤儿endpoints对象 e.checkLeftoverEndpoints() }() <-ctx.Done() } // worker runs a worker thread that just dequeues items, processes them, and marks them done. func (e *Controller) worker(ctx context.Context) { // 主循环调度 for e.processNextWorkItem(ctx) { } } func (e *Controller) processNextWorkItem(ctx context.Context) bool { // 获取队列service(queue-->processing,<-dirty) eKey, quit := e.queue.Get() if quit { return false } // 标记key已处理(移除processing队列的item) defer e.queue.Done(eKey) // 执行真正的同步 err := e.syncService(ctx, eKey.(string)) // 失败尝试重新入队 e.handleErr(err, eKey) 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
注意
endpoint controller启动后,会调度queue里面的service进行处理,调度成功会移除queue item,调度失败则进行重入队处理
# 2.3.syncService
syncService()是endpoint核心同步方法,基于service匹配关联pod,计算最新的endpoints对象,根据已有的endpoints进行对比以更新endpoints状态。// syncService update endpoints base service with pods. func (e *Controller) syncService(ctx context.Context, key string) error { ... // 1.获取service service, err := e.serviceLister.Services(namespace).Get(name) if errors.IsNotFound(err) { ... // service not found,delete endpoints err = e.client.CoreV1().Endpoints(namespace).Delete(ctx, name, metav1.DeleteOptions{}) return nil } ... // 2.获取service关联pod pods, err := e.podLister.Pods(service.Namespace).List(labels.Set(service.Spec.Selector).AsSelectorPreValidated()) ... // 3.初始化端点 for _, pod := range pods { // 跳过未分配地址Pod/正在删除Pod/走到终态Pod if !endpointutil.ShouldPodBeInEndpoints(pod, service.Spec.PublishNotReadyAddresses) { continue } // 实例化endpointAddress对象 ep, err := podToEndpointAddressForService(service, pod) if err != nil { continue } epa := *ep if endpointutil.ShouldSetHostname(pod, service) { epa.Hostname = pod.Spec.Hostname } // 补充endpointPort对象 if len(service.Spec.Ports) == 0 && service.Spec.ClusterIP == api.ClusterIPNone { // headless service允许无端口 subsets, totalReadyEps, totalNotReadyEps = addEndpointSubset(subsets, pod, epa, nil, service.Spec.PublishNotReadyAddresses) } else { // service定义端口 for i := range service.Spec.Ports { servicePort := &service.Spec.Ports[i] // 配合container定义端口确认使用端口号 portNum, err := podutil.FindPort(pod, servicePort) if err != nil { continue } // 生成endpointPort epp := endpointPortFromServicePort(servicePort, portNum) ... // 更新subsets subsets, readyEps, notReadyEps = addEndpointSubset(subsets, pod, epa, epp, service.Spec.PublishNotReadyAddresses) } } } // 计算最终的subsets,去重修正 // subsets: // - Addresses: [{ip: 10.244.1.10}] // Ports: [{port: 80}] // - Addresses: [{ip: 10.244.1.10}] // Ports: [{port: 443}] subsets = endpoints.RepackSubsets(subsets) // 4.获取当前endpoints currentEndpoints, err := e.endpointsLister.Endpoints(service.Namespace).Get(service.Name) ... // 5.对比endpoints是否需要更新 if !createEndpoints && // subsets相同 endpointutil.EndpointSubsetsEqualIgnoreResourceVersion(currentEndpoints.Subsets, subsets) && // endpoints label与service label一致(排除headless label) apiequality.Semantic.DeepEqual(compareLabels, service.Labels) && // subsets未超出容量(1000),无endpoints.kubernetes.io/over-capacity注解 capacityAnnotationSetCorrectly(currentEndpoints.Annotations, currentEndpoints.Subsets) { return nil } // 6.构造新的endpoints newEndpoints := currentEndpoints.DeepCopy() newEndpoints.Subsets = subsets newEndpoints.Labels = service.Labels ... // 7.endpoints subsets容量超出(触发截断,endpoints不再保留完整的后端,endpointslice切分保存) if truncateEndpoints(newEndpoints) { // 优先保留ready address,其次考虑截断ready address newEndpoints.Annotations[v1.EndpointsOverCapacity] = truncated } else { // 容量未超出,不需要标识注解 delete(newEndpoints.Annotations, v1.EndpointsOverCapacity) } // 8.更新endpoints label(headless类型标注label) if !helper.IsServiceIPSet(service) { newEndpoints.Labels = utillabels.CloneAndAddLabel(newEndpoints.Labels, v1.IsHeadlessService, "") } else { newEndpoints.Labels = utillabels.CloneAndRemoveLabel(newEndpoints.Labels, v1.IsHeadlessService) } // 9.创建endpoints if createEndpoints { _, err = e.client.CoreV1().Endpoints(service.Namespace).Create(ctx, newEndpoints, metav1.CreateOptions{}) // 10.更新endpoints } else { _, err = e.client.CoreV1().Endpoints(service.Namespace).Update(ctx, newEndpoints, metav1.UpdateOptions{}) } ... 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
注意
1.
endpoint controller核心作用是订阅service/pod资源变更事件,计算两者映射关系集合记录到endpoints,供kubeproxy进行nat2.
endpoints的映射集合最大容量1000,超出后将进行裁剪及标注truncate annotation,此时endpoints后端池不是完整的