1. 为什么每个写K8s控制器的人都绕不开Informer
我先说个很实在的场景。假设你现在接到一个任务:监控集群里所有Deployment的副本数变化,一旦发现期望副本数和实际副本数不一致,就自动扩缩容。你翻开client-go的文档,第一反应是"这还不简单,我写个for循环,每5秒List一次Deployment,对比一下不就行了?"
这种做法当然能跑,但等到你的控制面组件要同时监听几十种资源、每个资源几万个对象、还要在事件发生后几百毫秒内做出响应的时候,轮询方案就直接崩了——API Server被你打爆不说,数据延迟也完全不可控。这时候再去看社区里那些成熟的控制器代码,你会发现它们清一色在用client-go提供的informers包。Informer机制就是Kubernetes控制器模型的数据底座,理解了它,你看任何控制器代码都会轻松一大截。
这篇文章我准备把自己在项目里用Informer的完整思路、踩过的坑、以及源码层面的关键细节一次讲清楚。适合刚接触client-go的开发者,也适合已经写了几个控制器但还停留在"照着抄"阶段的同学。
简单说,Informer干的事情有三件:监听资源变化、维护本地缓存、触发事件回调。这三件事听起来都不复杂,但合在一起就产生了一个非常优雅的效果——你的控制器永远不用直接去查API Server,它只管从本地缓存里读数据,然后等着被事件通知就行了。
这套机制在Kubernetes生态里几乎是所有controller、operator、scheduler的公共基石。我后面会一层一层拆开讲。
2. 从ListAndWatch说起:Informer之类最底层的运作逻辑
2.1 为什么"只Watch"不够,还要先List一次
Informer最底层依赖的是一个叫ListAndWatch的函数,字面意思就是"先拉全量,再持续监听增量"。很多人不理解为什么非得List一次不可,不能直接Watch吗?
这里有个分布式系统的经典问题:你不能保证Watch连接建立的那一刻,API Server上发生过的历史事件你都能收到。万一网络抖动、连接断开重连,中间丢失的事件怎么补回来?所以client-go的做法很直接——先List一次拿到当前全量状态,然后在此基础上建立Watch。这样即使之前丢失过事件,也以全量数据为准刷新了本地状态,后续增量事件基于这个快照继续叠加,数据就是完整的。
List和Watch的分工是这样的:
// 简化版伪代码,展示核心逻辑 func ListAndWatch(client kubernetes.Interface, resource string) { // 1. List:全量拉取 list, _ := client.CoreV1().Pods("").List(context.TODO(), metav1.ListOptions{}) // 2. 把全量数据交给后续处理(DeltaFIFO) for _, item := range list.Items { store.Add(item) } // 3. Watch:增量监听 watch, _ := client.CoreV1().Pods("").Watch(context.TODO(), metav1.ListOptions{ ResourceVersion: list.ResourceVersion, }) for event := range watch.ResultChan() { // 处理 Added / Modified / Deleted 事件 } }关键就藏在ResourceVersion这里。List返回的时候会带一个资源版本号,这个版本号是API Server上的逻辑时钟。你把这个版本号传给Watch,API Server就知道"从我这个版本之后开始推送变更给我"。这一下就把全量和增量衔接起来了,不会漏事件,也不会重复处理。
2.2 断线重连时为什么不会炸
网络没有百分之百可靠的,Watch连接随时可能断开。Informer怎么处理重连的?它有一套"重新List+重新Watch"的兜底逻辑。
具体来说,Reflector(我后面会详细讲)内部维护了一个lastSyncResourceVersion,每次收到Watch事件都会更新它。当Watch连接断开时,Reflector会退避重试,重试成功后不是直接从断点续传,而是再走一遍"List全量 -> 拿最新ResourceVersion -> 重新Watch"的流程。
我第一次看到这个设计的时候心里打了个问号:断线重连就把全量数据重新拉一遍,那如果集群里有十万个Pod,不就白白浪费大量带宽吗?
后来我想明白了,这其实是在"简单可靠"和"极致的增量效率"之间做了一个务实的选择。直接从断点续传Watch虽然理论上可行,但这要求API Server必须把历史事件全量保留。实际上API Server的etcd里的事件数据是有保留期限的(默认一小时,可以通过--etcd-events-ttl调节),如果断线时间太长,断点早就被清理了,根本续不上。所以重新List就成了最稳妥的方案——它不依赖任何历史状态,天然自愈。
而且这里还有个很巧的优化:Informer的List+Watch只发生在连接建立初期和断线重连时,正常运行期间走的是长连接推送,流量开销并不大。尤其配合后面的本地缓存,95%以上的读操作根本不会碰到API Server。
3. 组件拆解:Reflector、DeltaFIFO、Indexer各管哪一段
3.1 Reflector:那个默默盯着API Server的哨兵
Reflector直译过来是"反射器",但在我的理解里,它就是驻守在API Server旁边的一个哨兵。它负责执行前面说的ListAndWatch循环,把API Server上的资源变化转变成一个个事件。
Reflector内部维护了几个关键信息:
expectedType:要监听的对象类型,比如*v1.Podstore:事件要写入的存储,实际就是DeltaFIFOlastSyncResourceVersion:最后一次同步的版本号,断线重连时用resyncPeriod:周期性重新同步的间隔(这个后面单独讲)
它每次Watch到事件,不是直接分发,而是包装成Delta(变更记录)塞进DeltaFIFO。一个Delta包含两部分:变更类型和变更后的对象。变更类型有几种:
type DeltaType string const ( Added DeltaType = "Added" Updated DeltaType = "Updated" Deleted DeltaType = "Deleted" Sync DeltaType = "Sync" )注意有个Sync类型,这个比较特殊,它不是API Server主动推送的,而是Informer自己定时生成的,用于周期性地把本地缓存里的对象再"重新同步"一遍。目的和作用我放在后面第4节专门说。
3.2 DeltaFIFO:事件的"待办队列"
DeltaFIFO是client-go里非常核心的一个数据结构,名字拆开看就很好理解:
Delta:一个变更记录,包含变更类型和对象FIFO:先进先出队列
也就是说,它是"变更记录的先进先出队列"。Reflector把事件写进来,消费者(processLoop)从队列头部取出去处理。之所以要用队列而不是直接处理,是因为生产和消费的速度不匹配——API Server的推送是突发的,可能瞬间产生几百个事件,而处理程序可能还在处理上一个事件。中间加一个队列做缓冲,能很好地消峰。
DeltaFIFO的精妙之处在于它的去重和合并逻辑。队列里的对象是按键(namespace/name)去重的,同一个对象的新事件不会无限堆积,而是会合并。比如一个Pod在短时间内连续被更新了三次,队列里不会攒三个Updated事件,而是把最新的对象状态直接覆盖旧的,队列里最多只保留一条有效的待处理记录。这样消费端永远处理的是"最新状态",而不是历史事件的堆叠。
这个设计背后的哲学值得品味:Kubernetes的控制器模型是最终一致的,中间过程根本不重要,重要的是最终状态。你去问一个控制器"这个Pod有几个副本?",它不需要知道Pod历史上经历过几次变更,只需要知道当前期望是几个。DeltaFIFO的合并逻辑天然符合这个哲学。
3.3 Indexer:读多写少场景下的本地缓存
Indexer是Informer对外提供的只读缓存,也是我之前说的"控制器不需要直接访问API Server"的关键支撑。
Indexer的内部实现是thread-safe map + 索引。默认的索引是namespace/name -> 对象,但你可以自定义索引。比如按Label分组、按NodeName分组,查询起来就会非常快。
// 自定义索引示例:按Pod的NodeName索引 indexers := cache.Indexers{ "nodeName": func(obj interface{}) ([]string, error) { pod := obj.(*v1.Pod) return []string{pod.Spec.NodeName}, nil }, } // 初始化Informer时传入 sharedInformerFactory := informers.NewSharedInformerFactory(client, 10*time.Minute) podInformer := sharedInformerFactory.Core().V1().Pods() podInformer.Informer().AddIndexers(indexers) // 使用索引查询 podsOnNode, _ := podInformer.Informer().GetIndexer().ByIndex("nodeName", "node-1")缓存的意义在控制器场景下怎么强调都不过分。一个集群有几千个Pod,每个Pod每秒上报一次状态,如果控制器每次处理事件都重新去API Server查一遍Pod详情,API Server的QPS会高得离谱。有了Indexer,控制器直接从本地缓存读,读操作开销几乎为零。这也是为什么我在写所有operator时都会要求"读写分离"——写操作才走client,读操作一律走Informer缓存。
3.4 全流程串起来看
我把Informer的完整数据流整理一下,你跟着走一遍就全通了:
Reflector启动,List全量数据写入DeltaFIFO,然后Watch增量事件Reflector把收到的事件包装成Delta,写入DeltaFIFO- DeltaFIFO的
processLoop持续从队列弹出变更记录,调用HandleDeltas HandleDeltas做两件事:先更新Indexer缓存,再调用用户注册的EventHandler回调- 你的控制器在回调里拿到变更后的对象,开始业务逻辑
这个流程里有一个细节很多人会忽略:先更新缓存,再触发回调。这意味着你在回调里通过Indexer查到的对象必然是最新的,不会出现"回调里读到旧数据"的尴尬。
4. 那些年我踩过的Informer实战坑
4.1 事件处理函数阻塞导致的连锁问题
这是我踩过最深的坑,没有之一。
刚开始写控制器的时候,我在AddFunc里直接调用了业务逻辑,这个业务逻辑里有访问数据库和调用外部API的操作。单个事件处理通常要200毫秒左右。当时测试环境的事件量不大,没发现问题。等上了预发环境,事件量一上来,我发现API Server的Watch连接频繁断开重连,日志里全是reflector: watch of *v1.Deployment closed。
排查下来根因是这样的:processLoop是单goroutine消费DeltaFIFO的,如果我在回调里做耗时操作,整个队列的消费就被卡住了。事件处理不过来,DeltaFIFO就会越堆越长,Reflector往队列里写数据的Add操作也会阻塞(DeltaFIFO是有界队列)。写不进去之后,Reflector认为Watch异常,主动断开重连。最终导致事件积压、重复处理、缓存不一致,整套机制直接崩溃。
正确的做法一定是:回调里只做轻量级操作,耗时逻辑放到独立goroutine或workqueue里异步处理。
// 错误示范:直接在回调里做耗时操作 podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { doExpensiveWork(obj) // 阻塞了processLoop }, }) // 正确做法:回调只负责入队 podWorkQueue := workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter()) podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { key, _ := cache.MetaNamespaceKeyFunc(obj) podWorkQueue.Add(key) // 入队,立即返回 }, UpdateFunc: func(oldObj, newObj interface{}) { key, _ := cache.MetaNamespaceKeyFunc(newObj) podWorkQueue.Add(key) }, }) // 然后启动worker goroutine从队列里取key处理client-go自带的workqueue包就是专门配合Informer使用的。它的核心价值有两个:一是去重,同一个key不需要重复入队多次;二是延迟和限速,失败的任务可以按指数退避的节奏重试,不会因为程序崩溃反复执行。
4.2 事件回调里修改了缓存对象导致的守卫事故
第二个坑更隐蔽。
我在一个项目里写了这么一段代码,UpdateFunc里拿到新对象后想改一下它的Annotation再存储:
UpdateFunc: func(oldObj, newObj interface{}) { pod := newObj.(*v1.Pod) pod.Annotations["processed"] = "true" // 直接修改了缓存里的对象! processPod(pod) }这段代码跑了一段时间,出现了一个诡异的现象:某些Pod明明已经被处理过了,但过一会又会触发一次Update事件,而且业务逻辑执行了两遍。
查了很久才发现,我修改的不是API Server上的对象,而是Indexer缓存里保存的同一份指针引用。缓存里的对象被改了,但API Server并不知道。等到下次有Pod真正被更新时,Reflector收到API Server推送的旧版本对象,放进DeltaFIFO,再更新Indexer时才发现"缓存里的对象和API Server的不一致",于是又触发了一次Update回调。
这种bug最难查的地方在于:它的表现不是直接报错,而是微妙的状态漂移。事件被重复触发,业务逻辑重复执行,可能带来重复计费、重复发送通知等严重后果。
这个事的教训就是:永远不要把缓存里的对象当作自己的私有数据。需要修改对象内容时,必须用DeepCopy:
UpdateFunc: func(oldObj, newObj interface{}) { pod := newObj.(*v1.Pod).DeepCopy() // 深拷贝之后再改 pod.Annotations["processed"] = "true" processPod(pod) }4.3 UpdateFunc触发的乱序处理问题
还有一个让我挠头的问题,是关于UpdateFunc只给新对象不给旧对象使用场景的。
有一次我需要统计Pod重启次数的变化趋势。一开始我在UpdateFunc里拿到新旧两个对象,比较Status.ContainerStatuses里的RestartCount,有变化就上报指标。看起来逻辑没问题,但实际跑起来指标总是对不上,偶尔还会出现负数增量。
问题出在事件不是保序的。Informer不保证回调的执行顺序和API Server产生事件的顺序完全一致,尤其是在watch重连之后,先收到的事件可能是后发生的。所以如果你的事务逻辑依赖"旧对象必然是新对象的前一个状态",那就会出错。
解决思路很直接:不要在回调里做"两个状态之间的差值计算",而是把两个状态都上报,由下游来做聚合计算。或者干脆只用最新状态做幂等处理,不依赖历史。
4.4 Resync的真正目的是什么
我第一次看到NewSharedInformerFactory(client, 10*time.Minute)这个参数的时候,以为是让Informer每10分钟重新拉一遍全量数据。这个理解是错的。
resyncPeriod做的事情是:每过这个时间间隔,Informer会把Indexer缓存里的所有对象重新包装成Sync类型的Delta,再一次投递给DeltaFIFO。这样做的目的不是"拉新数据",而是让你的EventHandler有机会周期性地看到所有对象的最新状态,即使这些对象在API Server上没有任何变化。
这有什么用?两个典型场景:
- 漏事件兜底:如果你的控制器因为某种原因丢失了某些事件(比如回调panic了),resync能让你周期性地补偿处理一遍全量对象。
- 定期核对:如果你需要在业务上周期性扫描所有对象的状态(比如每天检查一次证书是否即将过期),resync就是现成的定时器。
但要注意,resync只触发UpdateFunc(内部会把Sync转成Updated事件),不会触发AddFunc和DeleteFunc。而且resync的单位是"整个Informer工厂",不是"单个Informer",所以生产环境通常给一个适中值比如10分钟或15分钟,不要设置得太短,否则会产生大量的无效回调。
4.5 多Informer实例的事件漂移与一致性选择
当你能熟练使用单个Informer之后,下一个自然的问题就是:如果我的控制器需要同时监听Pod和Node两种资源,怎么保证它们的数据在时间上是一致的?
比如我需要根据Node的状态来决策Pod的处理逻辑——Node是Ready的才处理它的Pod。你可能会分别在Pod的Informer回调和Node的Informer回调里做业务逻辑。这就出现了一个问题:两个回调运行在不同的goroutine里,它们的执行时机是不确定的。你无法保证"我看到Node是Ready的时候,Pod的最新状态也一定已经更新到了Indexer里"。
社区对这种问题的经典解是:把事件统一入队,由同一个worker goroutine串行处理。也就是说,不管是Pod事件还是Node事件,都解析成key然后丢进同一个workqueue,worker从队列里取出key后,再去Indexer里查当前最新的Pod和Node状态,一起做决策。这样一来,决策的时刻读到的就是两个Indexer的当前状态,一致性就好很多。
当然,这仍然不是强一致(两个Indexer的更新时机仍然有极小的时间差),但对几乎所有控制器的业务场景来说已经足够了。Kubernetes本身就是最终一致的系统,你不可能也不需要做到绝对的强一致。
5. 读懂SharedInformerFactory:为什么大家都用工厂而不是裸写Informer
5.1 共享机制解决的是内存爆炸问题
我一开始写Informer的时候,是每个资源手动一个NewPodInformer,跑起来也正常。但随着监听的资源种类增多,我发现一个严重的问题:如果我的控制台同时需要Pod、Deployment、Service、Node、ConfigMap五种资源,每个Informer都会ListAndWatch一次。假设每种资源一万个对象,那就是五万次API请求打过去,内存里要归五份全量缓存。这还只是一个小项目,生产环境的大型controller可能要监听几十种资源。
SharedInformerFactory的核心机制就是同类型资源全局只创建一份Informer,所有消费者共享同一个Reflector、同一份Indexer缓存。
// 同一个factory,不管调用几次,返回的都是同一个PodInformer实例 factory := informers.NewSharedInformerFactory(client, 10*time.Minute) podInformer1 := factory.Core().V1().Pods() podInformer2 := factory.Core().V1().Pods() fmt.Println(podInformer1 == podInformer2) // true这意味着什么?假设你有12个控制器逻辑都要监听Pod,如果各自New一个Informer,就是12份Pod缓存,12条Watch连接;用SharedInformerFactory,它们全部归一,只有一份缓存、一条Watch连接,API Server的压力直接减少到原来的1/12。
多个控制器共享同一个Informer时,每个控制器可以注册自己独立的EventHandler,互不影响。这是"一对多"的广播模型——一个Informer的数据源,多个业务消费者。
5.2 Start:老生常谈但必须理解的启动顺序
SharedInformerFactory有一个很容易被忽视的细节:你必须在调用factory.Start(stopCh)之后,再使用Informer的Lister。如果顺序反了,你会拿到一份空缓存。
具体原因是,Start会为每个Informer启动一个goroutine执行Run,而Run里才真正开始执行ListAndWatch。在List完成之前,缓存是空的。如果你在Start之前就调用了Lister().List(),返回的就是空列表。
更稳妥的启动方式是先factory.Start(stopCh),再调用factory.WaitForCacheSync(stopCh)。WaitForCacheSync会阻塞等待所有Informer完成首次List,确保缓存可用后才放行。
factory := informers.NewSharedInformerFactory(client, 10*time.Minute) informer := factory.Core().V1().Pods() // 必须先启动 factory.Start(stopCh) // 必须等缓存同步完成 if !cache.WaitForCacheSync(stopCh, informer.Informer().HasSynced) { klog.Fatal("cache sync timeout") } // 到这里才能安全使用Lister pods, _ := informer.Lister().List(labels.Everything())这个WaitForCacheSync我当年第一次写的时候就跳过了,结果线上出现了一个非常尴尬的bug:控制器启动后的前几秒里,它认为集群里一个Pod都没有,把需要保留的Pod全部当成孤儿Pod清理了。这个问题后来我反思了很多,本质上是没有理解"缓存填充需要时间"这个基本事实。
5.3 多Informer之间的数据同步:等待关也是等
一个更进阶的场景:你的控制器在AddFunc里收到一个Pod事件,需要知道这个Pod属于哪个Deployment,所以去查Deployment的Indexer。但如果Deployment的Informer还没有完成首次List,查到的就是空的,你可能就会错误地认为"这个Pod不属于任何Deployment"。
这种跨资源依赖的问题,除了统一入队workqueue之外,还需要在启动阶段做额外的同步等待:
if !cache.WaitForCacheSync(stopCh, podInformer.Informer().HasSynced, deployInformer.Informer().HasSynced, ) { klog.Fatal("cache sync timeout") }把多个Informer的HasSynced都等一遍,确保启动阶段所有资源都缓存完整了,后续的事件处理就相对安全。
不过这依然有一种极端情况:运行过程中Deployment的Informer因为网络问题重连,导致一小段时间内Deployment缓存不完整。这正是Informer机制本身无法完全避免的,只能靠resync来兜底。
6. 从Informer到Workqueue:手写一个极简控制器的完整骨架
6.1 一个能跑的最小控制器
前面讲了那么多理论,我最后给你贴一个经量级的完整控制器骨架,你可以直接照着搭项目。这个骨架融合了我前面提到的所有实践要点:回调里只入队、worker里统一处理、跨资源同步等待、缓存读取。
package main import ( "context" "fmt" "time" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/clientcmd" "k8s.io/client-go/util/workqueue" "k8s.io/klog/v2" ) type PodController struct { informer cache.SharedIndexInformer queue workqueue.RateLimitingInterface client kubernetes.Interface } func NewPodController(client kubernetes.Interface) *PodController { // 只关注default命名空间的Pod,减少不必要的缓存 factory := informers.NewSharedInformerFactoryWithOptions( client, 10*time.Minute, informers.WithNamespace("default"), ) podInformer := factory.Core().V1().Pods() queue := workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter()) controller := &PodController{ informer: podInformer.Informer(), queue: queue, client: client, } podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { controller.enqueue(obj) }, UpdateFunc: func(oldObj, newObj interface{}) { controller.enqueue(newObj) }, DeleteFunc: func(obj interface{}) { controller.enqueue(obj) }, }) return controller } func (c *PodController) enqueue(obj interface{}) { key, err := cache.MetaNamespaceKeyFunc(obj) if err != nil { klog.ErrorS(err, "meta namespace key func failed") return } c.queue.Add(key) } func (c *PodController) processItem(ctx context.Context, key string) error { namespace, name, err := cache.SplitMetaNamespaceKey(key) if err != nil { return err } // 从Informer缓存中读取最新数据,而不是直接访问API Server pod, exists, err := c.informer.GetStore().GetByKey(key) if err != nil { return err } if !exists { // 对象已删除,只做清理逻辑 fmt.Printf("pod %s/%s has been deleted, cleanup\n", namespace, name) return nil } p := pod.(*corev1.Pod) fmt.Printf("processing pod %s/%s, phase=%s, node=%s\n", namespace, name, p.Status.Phase, p.Spec.NodeName) // 这里放你的业务逻辑 // 注意:不要阻塞过长时间,不要把外部API调用放在这里 return nil } func (c *PodController) Run(ctx context.Context, workers int) { defer c.queue.ShutDown() // 启动底层Informer,等待缓存同步 go c.informer.Run(ctx.Done()) if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) { klog.Error("cache sync timeout") return } klog.Info("cache synced, starting workers") // 启动N个worker并发消费队列 for i := 0; i < workers; i++ { go c.runWorker(ctx) } <-ctx.Done() } func (c *PodController) runWorker(ctx context.Context) { for c.processNextItem(ctx) { } } func (c *PodController) processNextItem(ctx context.Context) bool { key, shutdown := c.queue.Get() if shutdown { return false } defer c.queue.Done(key) err := c.processItem(ctx, key.(string)) if err != nil { // 处理失败,重新入队并限速重试 klog.ErrorS(err, "process item failed", "key", key) c.queue.AddRateLimited(key) return true } // 处理成功,遗忘这个key的失败历史 c.queue.Forget(key) return true } func main() { config, err := clientcmd.BuildConfigFromFlags("", "/root/.kube/config") if err != nil { panic(err) } client, err := kubernetes.NewForConfig(config) if err != nil { panic(err) } controller := NewPodController(client) ctx, cancel := context.WithCancel(context.Background()) defer cancel() controller.Run(ctx, 4) }这套骨架我用了很多次,每次新起operator项目都是在这个基础上改。它已经把事件回调、workqueue限速、缓存读取、worker并发这些核心问题都处理好了,你只需要专注在processItem里的业务逻辑上。
6.2 关于worker数量和处理性能的几个参考
worker数量不是越多越好。你要明白,worker并发数越高,对API Server的写请求并发也越高。如果你的业务逻辑里有大量的Update操作,建议worker数控制在2到4个,避免对API Server产生过大压力。
如果你的业务是纯计算型的,不需要写API Server,可以适当调高到8个或16个。
另外一个性能调优方向是informers.WithTweakListOptions。很多场景下你并不需要监听所有namespace的所有对象,可以通过LabelSelector、FieldSelector提前在源头上过滤。比如只监听带特定Label的Pod:
factory := informers.NewSharedInformerFactoryWithOptions( client, 10*time.Minute, informers.WithTweakListOptions(func(options *metav1.ListOptions) { options.LabelSelector = "app=my-app" }), )这样Reflector在List和Watch阶段就只关心符合条件的对象,Indexer缓存量级可能从几万降到几百,内存和CPU开销都会大幅下降。选型的时候这是第一个可以考虑的优化点。
7. 线下验证Informer机制的正确姿势
7.1 kube-apiserver的访问压力观察
写完代码总要验证一下Informer是不是真的在按预期工作。最直接的观察点就是API Server的访问日志或者请求指标。
在本地用kubectl反正随时可以观察,但更细致的做法是直接看client-go暴露的指标。Informer内部自带workqueue和reflector的metrics,可以通过/metrics端点暴露出来。核心指标有几个值得关注:
workqueue_depth:队列深度,如果长期不为0且持续增长,说明消费速度跟不上生产速度workqueue_adds_total:入队总数reflector_items_total:监听到的对象数量reflector_watch_events_total:Watch事件总数
如果你发现reflector_watch_events_total一直在增长,而业务逻辑没有对应的处理,大概率是回调里逻辑写得太重或者入队逻辑被遗漏了。
7.2 用代码模拟场景验证缓存一致性
我验证Informer缓存和API Server一致性的土办法是:起一个控制循环,每隔一段时间对比一次Indexer缓存里的对象列表和直接List API Server拿到的对象列表。
go func() { ticker := time.NewTicker(30 * time.Second) for range ticker.C { cachedPods, _ := informer.Lister().List(labels.Everything()) livePods, _ := client.CoreV1().Pods("").List(context.TODO(), metav1.ListOptions{}) if len(cachedPods) != len(livePods.Items) { klog.Warningf("cache size %d != live size %d", len(cachedPods), len(livePods.Items)) } } }()这个对比脚本看着简单,但能在早期发现很多问题,比如缓存恐慌、事件丢失、Indexer索引配置错误等。我在几个项目里靠这个办法抓到过两个隐蔽的缓存不一致问题——都是因为跨namespace的Informer配置参数写错了。
7.3 故障注入:按掉网络会发生什么
我强烈建议你在测试环境做一次"拔网线"实验。断开测试环境网络连接30秒,再恢复,看看你的Informer和控制器表现如何。
正常表现是这样的:
- Watch连接断开,Reflector开始退避重试(默认从1秒开始指数退避,最大到1000秒左右)
- 断网期间的业务事件不会被处理,队列里会堆积
- 恢复网络后,Reflector重新List全量数据,DeltaFIFO开始重新填充,队列里的key开始被消费
- 整个过程控制器不会崩溃,也不会死锁
如果这个实验里你的控制器表现异常,大概率是你在回调里做了太多依赖实时网络的操作——比如同步调用外部HTTP接口、访问数据库。这些都是Informer机制的死敌。正确做法是把这些操作全部丢到workqueue后面的worker里,由worker去执行,而不是在事件回调中直接执行。
8. 把Informer放进更复杂的系统里:多资源联动与权限边界
8.1 TPR监听带来的RBAC设计问题
当你用Informer监听自定义资源(CRD)时,有一个很重要但又很容易被忽略的环节:RBAC权限。Informer的ListAndWatch是直接访问API Server的,它需要拥有对应资源的list和watch权限。
我见过很多"能跑但一权限收紧就崩"的控制器。开发环境用的是kubeconfig,权限是admin,怎么跑怎么通。一上生产,集群管理员给了最小RBAC权限,结果控制器启动后Informer同步失败,所有缓存都是空的,业务逻辑全部失效。
如果你的控制器部署在集群内,正确做法是在ServiceAccount上绑定Role,并明确授予list和watch权限。这个在RBAC里经常被忽略,因为很多人只记得get就够了——但Informer偏偏用的不是get,而是list和watch。
apiVersion: rbac.authorization.k8s.io/v1 kind: Role metadata: namespace: my-system name: my-controller-role rules: - apiGroups: [""] resources: ["pods"] verbs: ["get", "list", "watch"] - apiGroups: [""] resources: ["pods/status"] verbs: ["get", "list", "watch", "update"]同时要注意,如果你同时监听了多个namespace的Pod(比如不用informers.WithNamespace限制),那你需要的就不是Role而是ClusterRole了,Grant范围要匹配Informer的监听范围。
8.2 多资源联动的编排能力从哪里来
Informer机制本身是没有"编排"能力的——它只是把每个资源的变化告诉你,至于"Pod变了之后要做哪几件事、这几件事的顺序是什么",完全由控制器自己来实现。
我当前的消息通知系统中就是这么做的:同时监听Deployment和Pod的Informer,一旦Deployment的副本数变化,就把事件入队到同一个workqueue;worker从队列取出key后,先查Deployment的缓存拿到期望副本数,再查Pod的缓存统计实际副本数,最后决定是否发出告警。这里的关键是"同一个workqueue串行处理不同资源的事件",避免并发处理引起的数据竞争。
这个模式在Kubernetes社区中就是标准的"多Informer + 单Workqueue + 串行处理"模式。几乎所有复杂的operator(比如etcd-operator、prometheus-operator)都是这么组织的。理解了Informer的机制,你看这些operator的源码会非常快,因为它们的高层逻辑就那么多,真正的复杂的是业务规则。
9. 把Scheduler的informer设计拿过来用
Kubernetes Scheduler是把Informer机制用得最极致的一个组件。它的调度逻辑里,Pod和Node的数据都不是从API Server实时读的,而是通过informer缓存维护的。调度器启动时会通过informer监听待调度的Pod和集群所有Node的实时状态(资源容量、亲和性、污点等),全部维护在本地缓存里。每来一个待调度Pod,调度器从本地缓存里快速找出符合约束的Node,绑定完成后把结果写回API Server。
这个模式的精髓在于:把"全量计算"变成"增量维护"。调度器不需要每次调度都去全量扫描集群,只在有事件发生时更新相应的缓存条目。这样即便集群有几千个Node,调度性能依然能维持在毫秒级。
所以当你在设计自己的系统时,如果你也需要一个"大量资源的状态查询 + 变化感知"的能力,Informer几乎是标配。它替你把最繁琐的数据同步、事件分发、缓存一致性都处理好了,让你能专注在自己的业务逻辑上——这才是Kubernetes生态里"约定优于配置"的一个绝佳例证。
从个人使用体验来说,Informer这套机制真正让我觉得舒服的地方在于它的"天然自愈"倾向。网络断了重连、事件丢失重拉、数据不一致靠resync兜底,处处体现着"在分布式环境下不要追求绝对精确,要追求最终一致"的设计哲学。理解了这套哲学,你会更容易看懂Kubernetes其他组件的设计意图,写出来的控制器质量也会明显上一个台阶。希望这篇文章对你有用。