client-go SharedInformer
在client-go中,informer是一个核心的概念,用于从Kubernetes API服务器中读取对象。它可以监视一个或多个API对象,并在对象发生变化时自动更新本地缓存。 Informer主要组成部分: sharedIndexInformer sharedIndexInformer是informer的核心组件。它负责从API服务器获取对象并更新本地缓存。在更新缓存后,它会触发事件处理器,以便通知其他组件对象已经更新。 informerSyncHandler informerSyncHandler是sharedIndexInformer的事件处理程序,它将缓存中的对象与API服务器中的对象进行比较,并更新缓存。当处理完所有的更新操作后,informerSyncHandler会将更新后的对象发送到事件队列。 informerEventHandler informerEventHandler是事件处理程序的接口,用于接收事件并执行特定的操作。当sharedIndexInformer接收到更新事件时,它将调用informerEventHandler来处理事件 indexer indexer是informer的本地缓存,用于存储从API服务器获取的对象。indexer使用map数据结构存储对象,其中键是对象的名称,值是对象本身。此外,indexer还使用索引数据结构(例如Set、List)来优化查找和筛选操作。每个indexer都关联一个ObjectStore,ObjectStore是一个更高级别的抽象,用于管理一组API对象 watcher watcher是informer的事件源,用于从API服务器中获取对象并将它们发送到事件队列中。watcher实现了Kubernetes API服务器上的watch机制,它会定期向API服务器发送请求,以获取当前对象的状态。在获取状态后,watcher将对比上一次请求的结果,找到任何发生变化的对象,并将它们发送到事件队列中 在运行时,informer将indexer和watcher组合在一起,以便实现对象的监视和更新。当watcher从API服务器中获取到对象时,它会将它们添加到indexer中。如果对象已经存在于indexer中,则watcher将更新对象的状态。一旦对象被更新,informerSyncHandler将被触发,以便将更新后的对象发送到事件队列中。在事件处理程序中,可以通过调用indexer对象的方法来访问缓存中的对象 画了一个可能不太标准的图: 了解了Informer架构和Reflector、DeltaFIFO、Indexer几个组件后,再回头来看示例中的SharedInformer使用,它如何将前面几个组件关联起来 示例程序 创建Informer对象 注册事件处理程序 启动Informer 以下面代码为例: 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 func main() { var err error var config *rest.Config var kubeconfig *string if home := homedir.HomeDir(); home != "" { kubeconfig = flag.String("kubeconfig", filepath.Join(home, ".kube", "config"), "[可选] kubeconfig 绝对路径") } else { kubeconfig = flag.String("kubeconfig", "", "kubeconfig 绝对路径") } flag.Parse() // 初始化 *rest.Config 对象 if config, err = rest.InClusterConfig(); err != nil { if config, err = clientcmd.BuildConfigFromFlags("", *kubeconfig); err != nil { panic(err.Error()) } } // 创建 *ClientSet 对象 clientSet, err := kubernetes.NewForConfig(config) if err != nil { panic(err.Error()) } // 初始化 SharedInformerFactory,暂且10s relist一次 sharedInformerFactory := informers.NewSharedInformerFactory(clientSet, time.Second*10) // 每种kubernetes内置资源都实现了Informer // 如创建podInformer:podInformer := sharedInformerFactory.Core().V1().Pods() // 下面创建deploymentInformer deployInformer := sharedInformerFactory.Apps().V1().Deployments() // 实际上调用了 InformerFor(&appsv1.Deployment{}, f.defaultInformer) 构造 deployInformer 中的 SharedIndexInformer informer := deployInformer.Informer() // 创建 DeploymentLister,拥有list方法 //type deploymentLister struct { // indexer cache.Indexer //} deployLister := deployInformer.Lister() // 注册事件处理程序 informer.AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: onAdd, UpdateFunc: onUpdate, DeleteFunc: onDelete, }) // 程序进程退出前,通知informer退出,informer以goroutine运行 stopper := make(chan struct{}) defer close(stopper) // 启动 informer 去 list & watch // 即 SharedIndexInformer 的 run方法,new了DeltaFIFO,然后new出sharedIndexInformer.controller,调用Controller的run方法(即Reflector的run方法),run里面创建了Reflector并执行 sharedInformerFactory.Start(stopper) // 等待所有启动的 Informer 的缓存被同步 sharedInformerFactory.WaitForCacheSync(stopper) informer.Run(stopper) // 从本地缓存中获取 default 中的所有 deployment 列表 deployments, err := deployLister.Deployments("default").List(labels.Everything()) if err != nil { panic(err) } for idx, deploy := range deployments { fmt.Printf("%d -> %s\n", idx+1, deploy.Name) } <-stopper } func onAdd(obj interface{}) { deploy := obj.(*v1.Deployment) fmt.Println("add a deployment:", deploy.Name) } func onUpdate(old, new interface{}) { oldDeploy := old.(*v1.Deployment) newDeploy := new.(*v1.Deployment) fmt.Println("update deployment:", oldDeploy.Name, newDeploy.Name) } func onDelete(obj interface{}) { deploy := obj.(*v1.Deployment) fmt.Println("delete a deployment:", deploy.Name) } 执行逻辑 NewSharedInformerFactory创建Informer对象,函数中传入了clientset和defaultResync参数,clientset是访问APIServer的具体实现,resyncPeriod指定了informer在从API服务器中获取数据之间等待的时间,sharedInformerFactory的informers字段是一个map结构,每种资源类型为key,对应value是SharedIndexInformer 1 2 3 4 5 6 7 8 9 10 func NewSharedInformerFactory(client kubernetes.Interface, defaultResync time.Duration) SharedInformerFactory factory := &sharedInformerFactory{ client: client, namespace: v1.NamespaceAll, defaultResync: defaultResync, // 示例中我们传入的是10s informers: make(map[reflect.Type]cache.SharedIndexInformer), startedInformers: make(map[reflect.Type]bool), customResync: make(map[reflect.Type]time.Duration), } 每种kubernetes内置资源都已经实现了资源对应的Informer,通过sharedInformerFactory.Apps().V1().Deployments()链式调用构造出deployInformer,遵循组、版本、资源(即GVK)格式,其他k8s资源的informer创建也是一样的操作 1 2 3 4 5 6 type deploymentInformer struct { factory internalinterfaces.SharedInformerFactory // TweakListOptionsFunc is a function that transforms a v1.ListOptions tweakListOptions internalinterfaces.TweakListOptionsFunc namespace string } deployInformer.Informer() 用deployInformer创建deployment对应的SharedIndexInformer。这里面初始化了重要的两个对象:cache.Indexers和ListerWatcher 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 // 给InformerFor传入两个参数: // 1. Deployment对象空结构体 // 2. Deployment的默认Informer构造函数(它在调用后返回SharedIndexInformer) // 返回Deployment对应的SharedIndexInformer func (f *deploymentInformer) Informer() cache.SharedIndexInformer { return f.factory.InformerFor(&appsv1.Deployment{}, f.defaultInformer) } // InformerFor的参数defaultInformer方法的实现,返回了构造出的SharedIndexInformer func (f *deploymentInformer) defaultInformer(client kubernetes.Interface, resyncPeriod time.Duration) cache.SharedIndexInformer { // 传入了参数cache.Indexers // Indexers的键是索引分类名,为常量namespace,值为IndexFunc,取出对象的namespace值 return NewFilteredDeploymentInformer(client, f.namespace, resyncPeriod, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}, f.tweakListOptions) } func NewFilteredDeploymentInformer(client kubernetes.Interface, namespace string, resyncPeriod time.Duration, indexers cache.Indexers, tweakListOptions internalinterfaces.TweakListOptionsFunc) cache.SharedIndexInformer { // 传入了ListerWatcher对象和Indexer参数来构造SharedIndexInformer return cache.NewSharedIndexInformer( &cache.ListWatch{ // ListFunc/WatchFunc函数就是来list和watch APIServer资源的动作,具体实现在clientset里面 ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { if tweakListOptions != nil { tweakListOptions(&options) } return client.AppsV1().Deployments(namespace).List(context.TODO(), options) }, WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) { if tweakListOptions != nil { tweakListOptions(&options) } return client.AppsV1().Deployments(namespace).Watch(context.TODO(), options) }, }, &appsv1.Deployment{}, resyncPeriod, // NewSharedInformerFactory时传入的10s indexers, ) } func NewSharedIndexInformer(lw ListerWatcher, exampleObject runtime.Object, defaultEventHandlerResyncPeriod time.Duration, indexers Indexers) SharedIndexInformer { realClock := &clock.RealClock{} // 最终返回的informer是这个sharedIndexInformer sharedIndexInformer := &sharedIndexInformer{ processor: &sharedProcessor{clock: realClock}, indexer: NewIndexer(DeletionHandlingMetaNamespaceKeyFunc, indexers), listerWatcher: lw, objectType: exampleObject, resyncCheckPeriod: defaultEventHandlerResyncPeriod, // 10s defaultEventHandlerResyncPeriod: defaultEventHandlerResyncPeriod, // 10s cacheMutationDetector: NewCacheMutationDetector(fmt.Sprintf("%T", exampleObject)), clock: realClock, } return sharedIndexInformer } // 入参: // 1. 通用对象Object,此处即Deployment{} // 2. deployment资源类型对应的sharedIndexInformer的构造函数 func (f *sharedInformerFactory) InformerFor(obj runtime.Object, newFunc internalinterfaces.NewInformerFunc) cache.SharedIndexInformer { f.lock.Lock() defer f.lock.Unlock() // 反射取得对象类型为Deployment informerType := reflect.TypeOf(obj) informer, exists := f.informers[informerType] if exists { return informer } resyncPeriod, exists := f.customResync[informerType] if !exists { resyncPeriod = f.defaultResync // 默认是零值,所以赋予我们传入的10s } // 如果没有此类型的informer的话就新建一个,此例是deployment类型的informer informer = newFunc(f.client, resyncPeriod) // 填充了sharedInformerFactory的informers字段 // informers map[reflect.Type]cache.SharedIndexInformer f.informers[informerType] = informer // 返回sharedIndexInformer实例 return informer } 创建 DeploymentLister,它里面就一个字段indexer,indexer是索引器对象,可以在本地缓存中根据索引取数。如从本地缓存indexer中获取 default 名称空间的所有 deployment 列表:deployments, err := deployLister.Deployments("default").List(labels.Everything()) 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 func (f *deploymentInformer) Lister() v1.DeploymentLister { // 传入一个参数为cache.Indexer,即前面func (f *deploymentInformer) defaultInformer方法中 // 初始化的cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc} return v1.NewDeploymentLister(f.Informer().GetIndexer()) } func (s *sharedIndexInformer) GetIndexer() Indexer { return s.indexer } func NewDeploymentLister(indexer cache.Indexer) DeploymentLister { return &deploymentLister{indexer: indexer} } // deploymentLister implements the DeploymentLister interface. type deploymentLister struct { indexer cache.Indexer } 至此,我们有了sharedIndexInformer,其已拥有重要的字段cache.Indexers和ListerWatcher已初始化,但目前还没有启动ListerWatcher去从APIServer拿数据,没有把这些数据包装为Delta缓存下来,没有存入到Indexer,也没有后续对数据的处理逻辑 ...