informer
Overview
Factory & Informers
Single Informer
Informer monitors the changes of target resource. An informer is created for each of the target resources if you need to handle multiple resources (e.g. podInformer, deploymentInformer).
The snippets and diagrams below target client-go v0.37.1. Structs show implementation details rather than APIs to construct directly.
types
Interface SharedInformerFactory
type SharedInformerFactory interface {
internalinterfaces.SharedInformerFactory
Start(stopCh <-chan struct{})
StartWithContext(ctx context.Context)
Shutdown()
WaitForCacheSync(stopCh <-chan struct{}) map[reflect.Type]bool
WaitForCacheSyncWithContext(ctx context.Context) cache.SyncResult
ForResource(resource schema.GroupVersionResource) (GenericInformer, error)
InformerFor(obj runtime.Object, newFunc internalinterfaces.NewInformerFunc) cache.SharedIndexInformer
Admissionregistration() admissionregistration.Interface
Internal() apiserverinternal.Interface
Apps() apps.Interface
Autoscaling() autoscaling.Interface
Batch() batch.Interface
Certificates() certificates.Interface
Coordination() coordination.Interface
Core() core.Interface
Discovery() discovery.Interface
Events() events.Interface
Extensions() extensions.Interface
Flowcontrol() flowcontrol.Interface
Lifecycle() lifecycle.Interface
Networking() networking.Interface
Node() node.Interface
Policy() policy.Interface
Rbac() rbac.Interface
Resource() resource.Interface
Scheduling() scheduling.Interface
Storage() storage.Interface
Storagemigration() storagemigration.Interface
}
Implementation sharedInformerFactory
type sharedInformerFactory struct {
client kubernetes.Interface
namespace string
tweakListOptions internalinterfaces.TweakListOptionsFunc
lock sync.Mutex
defaultResync time.Duration
customResync map[reflect.Type]time.Duration
transform cache.TransformFunc
informerName *cache.InformerName
informers map[reflect.Type]cache.SharedIndexInformer
startedInformers map[reflect.Type]bool
wg sync.WaitGroup
shuttingDown bool
}
Fields:
client: clientset to interact with API servernamespace: you can specify a namespace or all namespaces (v1.NamespaceAll) by defaultinformers: store created informers to start them whenfactory.Startis called.
The factory exposes API groups, each group exposes versions, and each version exposes resource-specific informer accessors. Accessor construction, shared-informer registration, and starting watches are separate steps.
How a new informer is created with a Factory:
-
Create a factory.
kubeInformerFactory := kubeinformers.NewSharedInformerFactory(kubeClient, time.Second*30) -
Obtain the Deployment informer accessor.
deploymentInformer := kubeInformerFactory.Apps().V1().Deployments()The v0.37.1 call chain is:
Call Declared return type Implementation factory.Apps() apps.InterfaceCalls apps.New(f, f.namespace, f.tweakListOptions)apps.New(...) InterfaceReturns &group{factory: f, namespace: namespace, tweakListOptions: tweakListOptions}group.V1() v1.InterfaceCalls v1.New(g.factory, g.namespace, g.tweakListOptions)v1.New(...) InterfaceReturns &version{factory: f, namespace: namespace, tweakListOptions: tweakListOptions}version.Deployments() TypedDeploymentInformerReturns &deploymentInformer{factory: v.factory, namespace: v.namespace, tweakListOptions: v.tweakListOptions}In particular, Apps is a method on the concrete
sharedInformerFactory;kubeInformerFactoryis the sample variable, not a type or an upstream declaration. The method body is:func (f *sharedInformerFactory) Apps() apps.Interface { return apps.New(f, f.namespace, f.tweakListOptions) }At this point the Deployment accessor exists, but it has not yet registered a shared informer or started API watches.
-
Obtain the shared informer and register handlers, usually while constructing your controller.
The accessor's Informer and TypedInformer methods are:
func (f *deploymentInformer) Informer() cache.SharedIndexInformer { return f.TypedInformer() } func (f *deploymentInformer) TypedInformer() DeploymentIndexInformer { return cache.NewTypedSharedIndexInformer[*apiappsv1.Deployment](f.factory.InformerFor(&apiappsv1.Deployment{}, f.defaultInformer)) }factory.InformerFor locks the factory and looks up
reflect.TypeOf(obj)inf.informers. If registered, it returns the existing informer. Otherwise, it selects the custom or default resync period, callsnewFunc(f.client, resyncPeriod), applies the factory transform if configured, and registers the result inf.informers.Here
newFuncis deploymentInformer.defaultInformer:func (f *deploymentInformer) defaultInformer(client kubernetes.Interface, resyncPeriod time.Duration) cache.SharedIndexInformer { return NewTypedDeploymentInformerWithOptions(client, f.namespace, internalinterfaces.InformerOptions{ResyncPeriod: resyncPeriod, Indexers: cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}, InformerName: f.factory.InformerName(), TweakListOptions: f.tweakListOptions}) }NewTypedDeploymentInformerWithOptions builds the Deployment ListWatch, wraps it with WatchList semantics, and calls
cache.NewSharedIndexInformerWithOptions. The List/Watch callbacks useclient.AppsV1().Deployments(namespace)and apply TweakListOptions. The resulting informer is wrapped withcache.NewTypedSharedIndexInformer[*apiappsv1.Deployment]. Construction sets up the objects; running them starts API requests.Register your handler and check the error:
_, err := deploymentInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: handleAdd, }) if err != nil { return err }handleAddis your callback. Calling Lister also obtains the informer, because Lister needs its Indexer. -
Start the registered informers.
kubeInformerFactory.StartWithContext(ctx)StartWithContext starts each registered informer that has not already started, using
informer.RunWithContext(ctx). Merely calling Apps().V1().Deployments() before Start is insufficient. If an informer is registered after Start, call StartWithContext again to start it.Wait for cache synchronization before reading through its Lister. Cancel the context before calling Shutdown to wait for informer goroutines. The channel-based Start used by the sample delegates to StartWithContext.
The next sections explain the shared informer's internals and RunWithContext lifecycle.
Interface SharedInformer
- Interface:
SharedInformer
SharedIndexInformer
// Selected methods; see the linked interface for all options. type SharedInformer interface { AddEventHandler(handler ResourceEventHandler) (ResourceEventHandlerRegistration, error) AddEventHandlerWithResyncPeriod(handler ResourceEventHandler, resyncPeriod time.Duration) (ResourceEventHandlerRegistration, error) AddEventHandlerWithOptions(handler ResourceEventHandler, options HandlerOptions) (ResourceEventHandlerRegistration, error) RemoveEventHandler(handle ResourceEventHandlerRegistration) error GetStore() Store GetController() Controller RunWithContext(ctx context.Context) HasSynced() bool LastSyncResourceVersion() string SetWatchErrorHandlerWithContext(handler WatchErrorHandlerWithContext) error SetTransform(handler TransformFunc) error }type SharedIndexInformer interface { SharedInformer // AddIndexers add indexers to the informer before it starts. AddIndexers(indexers Indexers) error GetIndexer() Indexer }
Implementation sharedIndexInformer
type sharedIndexInformer struct {
indexer Indexer
controller Controller
synced chan struct{}
processor *sharedProcessor
cacheMutationDetector MutationDetector
listerWatcher ListerWatcher
objectType runtime.Object
objectDescription string
resyncCheckPeriod time.Duration
defaultEventHandlerResyncPeriod time.Duration
clock clock.Clock
started, stopped bool
startedLock sync.Mutex
blockDeltas sync.Mutex
watchErrorHandler WatchErrorHandlerWithContext
transform TransformFunc
identifier InformerNameAndResource
informerMetricsProvider InformerMetricsProvider
keyFunc KeyFunc
}
NewSharedIndexInformerWithOptions initializes the Indexer, processor, mutation detector, and synchronization channels. NewSharedIndexInformer delegates to that constructor. The informer retains its ListerWatcher until RunWithContext creates the low-level controller and Reflector.
Components: - Indexer - controller: explained below - sharedProcessor: explained below - ListerWatcher
- Call
newQueueFIFOto construct the informer queue. Its implementation depends on feature gates (includingInOrderInformers); do not assume that every shared informer always uses DeltaFIFO. The queue accepts changes from the Reflector and supplies them toConfig.ProcessorConfig.ProcessBatch. - Create Controller with New
- Run s.cacheMutationDetector.Run
- Run
s.processor.run<- start all listeners using a separate processor context, stopped after the low-level controller. listeners are added viaAddEventHandler. (usually withcache.ResourceEventHandlerFuncs{AddFunc: xx, UpdateFunc: xx, DeleteFunc: xx}) - Run s.controller.Run <- refer the controller section 1. Create a new Reflector and call r.Run (ListAndWatch is called inside)
NewSharedInformer:
- NewSharedInformer: call NewSharedIndexInformer with
Indexers{}.NewSharedIndexInformer(lw, exampleObject, defaultEventHandlerResyncPeriod, Indexers{}) - NewSharedIndexInformer
func NewSharedIndexInformer(lw ListerWatcher, exampleObject runtime.Object, defaultEventHandlerResyncPeriod time.Duration, indexers Indexers) SharedIndexInformer { return NewSharedIndexInformerWithOptions( lw, exampleObject, SharedIndexInformerOptions{ ResyncPeriod: defaultEventHandlerResyncPeriod, Indexers: indexers, }, ) }
sharedProcessor
Role: hold a collection of listeners and distribute notification objects to them. distribute selects listeners; each listener buffers notifications and invokes the registered handler. Cache synchronization and completion of a handler's initial notifications are distinct; use the registration handle's HasSynced for the latter.
type sharedProcessor struct {
listenersStarted bool
listenersLock sync.RWMutex
listenersRCond *sync.Cond // Caller of Wait must hold a read lock on listenersLock.
listeners map[*processorListener]bool
clock clock.Clock
wg wait.Group
}
Listenersare added for ResourceEventHandler via AddEventHandlerdistribute()callslistener.addto propagate new events to each listener.distribute()is called byinformer.OnAdd,informer.OnUpdate, andinformer.OnDeleterun()callslistener.runandlistener.popfor all listeners.handler.OnAdd,handler.OnUpdate,handler.OnDeletebased on the notification type.
Controller
Role: Run a reflector and enqueue item to Queue from ListerWatcher and process item from the queue with processfunc.
Interface:
type Controller interface {
RunWithContext(ctx context.Context)
Run(stopCh <-chan struct{})
HasSynced() bool
HasSyncedChecker() DoneChecker
LastSyncResourceVersion() string
}
Implementation:
type controller struct {
config Config
reflector *Reflector
reflectorMutex sync.RWMutex
clock clock.Clock
}
- Most things are passed by
Config(ListerWatcher, ObjectType, Queue (informer queue))
Run:
- Create a Reflector with NewReflectorWithOptions
- Run
reflector.RunWithContext(details -> ref reflector)ListAndWatchWithContexthandleWatch:- event.Added -> store.Add
- event.Modified -> store.Update
- event.Deleted -> store.Delete (store = Queue)
- Run processLoop until the context is canceled.
- Pop item from the Queue and process it repeatedly. (Actual process is given by
Config.Process, controller is just a container to executeProcess)Config.Process: handleDeltashandleDeltas(logger, obj, isInInitialList)calls processDeltas(logger, s, s.indexer, deltas, isInInitialList, s.keyFunc)handler: sharedIndexInformerclientState: s.indexer
- Keep indexer up-to-date by calling
indexer.Update(),indexer.Add(),indexer.Delete(). - Distribute notification and add object to cacheMutationDetector by calling
sharedIndexInformer.OnUpdate(),sharedIndexInformer.OnAdd(),sharedIndexInformer.OnDelete()
- Pop item from the Queue and process it repeatedly. (Actual process is given by
MutationDetector
Role: Check if a cached object is mutated. Call failurefunc or panic if mutated.
- By default, mutation detector is not enabled. (You can skip this components)
var mutationDetectionEnabled = false func init() { mutationDetectionEnabled, _ = strconv.ParseBool(os.Getenv("KUBE_CACHE_MUTATION_DETECTOR")) } - Run periodically calls CompareObjects.
- CompareObjects compares
cachedandcopiedofcacheObjind.cachedObjsandd.retainedCachedObjs.type cacheObj struct { cached interface{} copied interface{} } - If any object is altered, call
failureFunc. (if created with NewCacheMutationDetector, it doesn't have failureFunc, the program goespanic) - AddObject adds an object to
d.addedObjs. - You can enable the mutation detector to catch accidental mutation of shared cached objects. The updated sample calls
DeepCopy()before changing labels, so it should not producepanic: cache *v1.Pod modified. Mutating the original object would trigger that failure.KUBE_CACHE_MUTATION_DETECTOR=true go run ./contents/kubernetes-operator/client-go/informer
Example
- Initialize clientset with
.kube/config - Create an informer factory with the following line.
The second argument specifies ResyncPeriod, which defines the interval of resync (The resync operation consists of delivering to the handler an update notification for every object in the informer's local cache). For more detail, please read NewSharedInformer
informerFactory := informers.NewSharedInformerFactory(kubeClient, time.Second*30) -
Obtain a Pod informer accessor. Calling Informer() or Lister() registers its shared informer; starting the factory begins watching Pods.
podInformer := informerFactory.Core().V1().Pods()factory -> group -> version -> resource accessor (
TypedPodInformer, which embedsPodInformer)type PodInformer interface { Informer() cache.SharedIndexInformer Lister() v1.PodLister }Informer()returnsSharedIndexInformer- Delegate to
TypedInformer(), which calls the factory's InformerFor and wraps the result withcache.NewTypedSharedIndexInformer[*apicorev1.Pod]. - On first registration,
defaultInformerconstructs it throughNewTypedPodInformerWithOptions, including the namespace index, informer name, and TweakListOptions. - Reuse the registered shared informer on later calls. See the Pod accessor implementation.
- Delegate to
Lister()returns PodLister- call
v1.NewPodLister(f.Informer().GetIndexer()) - NewPodLister returns podLister with the given indexer.
type podLister struct { listers.ResourceIndexer[*corev1.Pod] }
- call
-
Add event handlers (
AddFunc,UpdateFunc, andDeleteFunc) to the pod informer._, err := podInformer.Informer().AddEventHandler( cache.ResourceEventHandlerFuncs{ AddFunc: handleAdd, UpdateFunc: handleUpdate, DeleteFunc: handleDelete, }, )handleAdd,handleUpdate, andhandleDeletedefine custom logic for each event. In this example, just print"handleXXX is called" -
Create a signal-aware context and start the factory. Check the error returned by
AddEventHandlerfirst.ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() ch := ctx.Done() informerFactory.Start(ch) defer informerFactory.Shutdown() -
Wait until the cache is synced.
cacheSynced := podInformer.Informer().HasSynced if ok := cache.WaitForCacheSync(ch, cacheSynced); !ok { log.Print("cache sync stopped") return } log.Println("cache is synced")func WaitForCacheSync(stopCh <-chan struct{}, cacheSyncs ...InformerSynced) bool { err := wait.PollImmediateUntil(syncedPollPeriod, func() (bool, error) { for _, syncFunc := range cacheSyncs { if !syncFunc() { return false, nil } } return true, nil }, stopCh) if err != nil { return false } return true }The legacy stop-channel wrapper can be adapted to a context with wait.ContextForChannel. New code can retain the context directly:
func ContextForChannel(parentCh <-chan struct{}) context.Context { return channelContext{stopCh: parentCh} } -
Wait for cancellation. The factory watches continuously without a separate polling loop.
<-ctx.Done()
Run and check
Run from the repository root with Pod list/watch permissions. The logs below illustrate the event sequence; timestamps and Pod names depend on the cluster. The current sample also prints labels from a copied Pod. It does not modify the shared cache or API object. Delete keys use DeletionHandlingMetaNamespaceKeyFunc, which handles DeletedFinalStateUnknown tombstones.
1. Run
go run ./contents/kubernetes-operator/client-go/informer
-
All Pods are synced in the cache.
1. Create a2021/12/21 09:05:08 handleAdd is called for Pod (key: local-path-storage/local-path-provisioner-547f784dff-lhwfk) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/kube-scheduler-kind-control-plane) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/etcd-kind-control-plane) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/kube-apiserver-kind-control-plane) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/kindnet-nzc7p) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/coredns-558bd4d5db-b4wjg) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/kube-controller-manager-kind-control-plane) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/kube-proxy-vrcbc) 2021/12/21 09:05:08 handleAdd is called for Pod (key: kube-system/coredns-558bd4d5db-8q78s) 2021/12/21 09:05:08 handleAdd is called for Pod (key: default/foo-sample-688594b488-782kw) 2021/12/21 09:05:08 cache is syncedPodwith namenginx.1. Handlers are called by the events of the createdkubectl run nginx --image=nginxPod.1. Delete the2021/12/21 09:05:20 handleAdd is called for Pod (key: default/nginx) 2021/12/21 09:05:20 handleUpdate is called for Pod (key: default/nginx) 2021/12/21 09:05:20 handleUpdate is called for Pod (key: default/nginx)Pod1. Handlers are called by the events of the Pod deletion.kubectl delete po nginx1. Stop with Ctrl+C. Context cancellation stops watches, and2021/12/21 09:05:29 handleUpdate is called for Pod (key: default/nginx) 2021/12/21 09:05:30 handleUpdate is called for Pod (key: default/nginx) 2021/12/21 09:05:31 handleUpdate is called for Pod (key: default/nginx) 2021/12/21 09:05:31 handleUpdate is called for Pod (key: default/nginx) 2021/12/21 09:05:31 handleDelete is called for Pod (key: default/nginx)factory.Shutdown()waits for informer goroutines. 1. Resync delivers Update notifications for cached objects every 30 seconds. It does not perform a fresh API List every 30 seconds.2021/12/21 09:27:08 handleUpdate is called for Pod (key: local-path-storage/local-path-provisioner-547f784dff-lhwfk) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/kube-apiserver-kind-control-plane) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/coredns-558bd4d5db-b4wjg) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/kube-controller-manager-kind-control-plane) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/coredns-558bd4d5db-8q78s) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: default/foo-sample-688594b488-782kw) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/kube-scheduler-kind-control-plane) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/etcd-kind-control-plane) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/kindnet-nzc7p) 2021/12/21 09:27:08 handleUpdate is called for Pod (key: kube-system/kube-proxy-vrcbc)
reference
- https://adevjoe.com/post/client-go-informer/
- https://www.huweihuang.com/kubernetes-notes/code-analysis/kube-controller-manager/sharedIndexInformer.html
- https://yangxikun.com/kubernetes/2020/03/05/informer-lister.html
Tests
go test ./contents/kubernetes-operator/client-go/informer checks that handleAdd does not mutate a shared cached Pod and that deletion keys work for both objects and DeletedFinalStateUnknown tombstones. The live watch sequence above requires a cluster; these regression tests do not exercise API-server delivery.