1. 为什么搞 Kubernetes 自动化必须先啃透 client-go
如果你写过 Kubernetes 控制器、Operator 或者任何跟集群自动化沾边的工具,大概率已经跟 client-go 打过照面。这个库是 Kubernetes 官方维护的 Go 客户端,几乎所有周边生态——从 kubectl 这样的命令行工具,到 kube-controller-manager 里的内置控制器,再到你手写的自定义调度器、Device Plugin、巡检脚本——底层都是通过它跟 apiserver 通信。
但很多人对 client-go 的理解停留在"哦,就是个 SDK,封装了 API 调用"这个层面。真到自己写代码的时候,会遇到一堆问题:informers 和 workqueue 是什么关系?为什么不直接用 List 接口轮询就够了?Update 和 Patch 到底该用哪个?为什么集群里跑着跑着权限就 Forbidden 了?这些问题不搞清楚,写出来的程序要么性能稀烂,要么三天两头出诡异故障。
这篇文章我把 client-go 的核心机制、日常用法和踩过的坑一次性讲透。适合两类人:一是刚开始写 Kubernetes 控制器、想搞清楚底层原理的开发者;二是已经写过一些自动化工具但总觉得哪里不对劲、想系统性补课的人。文章里会有原理拆解、可直接跑的代码、配置参数的经验值,以及我在生产环境里真实踩过的问题记录。
2. client-go 的整体架构:一个核心,四条主线
2.1 一个核心:RESTClient
client-go 的底层,是一个叫 RESTClient 的东西。可以把它理解成"翻译官":负责把 Go 结构体变成 HTTP 请求发给 apiserver,再把响应解析回 Go 结构体。路径拼接、认证头注入、请求体编码、错误处理,这些脏活累活全在它里面。
你几乎不会直接调用 RESTClient,但它是所有上层能力的基石。上层那些 typed client、dynamic client、discovery client,本质都是基于 RESTClient 加了一层类型约束或通用处理逻辑。所以调优、排错的时候,最终都要回到这一层去看请求是怎么发出去的。
2.2 四条主线拆解
我们可以把 client-go 按功能切成四块,这样理解起来特别清晰:
ClientSet:日常打交道最多的入口。把内置资源按 API 分组封装成类型安全的方法集合。比如操作 Deployment 用
client.AppsV1().Deployments(ns).Get(...),操作 Pod 用client.CoreV1().Pods(ns).Get(...)。类型安全的意思是参数、返回值都是具体结构体,编译期就能抓住大部分类型错误。DynamicClient:专门处理没有固定结构的资源,比如 CRD。因为 CRD 的结构事先不知道,它用
unstructured.Unstructured这种通用 map 结构承接任意 JSON 数据。适合写通用工具、CRD 控制器这类场景。DiscoveryClient:负责探测集群支持哪些 API 组、版本、资源类型。写工具时经常要先确认资源是否存在、支持哪些操作,再决定走哪条路。
Informer / Lister / Indexer:这是 client-go 最精华的部分。一句话概括:它让你在本地有一份集群数据的实时缓存,读数据不用每次打 apiserver,变更靠推送而不是轮询。能不能写出高性能控制器,就看这块用得怎么样。
另外还有几个辅助模块也经常用到:WorkQueue带去重和延迟能力的工作队列,处理事件时攒着慢慢干活;EventRecorder往集群写事件,方便用kubectl describe排查问题;leaderelection选主机制,多副本控制器避免重复干活;RESTMapper把资源类型名映射到 REST 路径。
2.3 为什么要设计成这么多层
第一次看 client-go 源码的人都会觉得层数太多。但用久了会发现每一层都有存在的道理。
分层最大的好处是复用。RESTClient 处理"怎么发请求"这种通用逻辑,不管上层是内置资源还是 CRD,都共用同一套传输、认证、重试代码。接口隔离也做得很好:缓存、监听、索引这块重逻辑单独抽成 informer,你写业务时不用把"怎么跟 apiserver 保持同步"和"拿到变更后干什么"搅在一起。控制器模式能成为 Kubernetes 生态的事实标准,跟这种拆分有直接关系。
3. 核心机制拆解:Informer 是理解 client-go 的分水岭
3.1 Informer 的四个核心行为
一个 informer 干四件事:
- List:启动时全量拉取某类资源一次,拿到集群当前快照,构建本地缓存初始副本。
- Watch:跟 apiserver 建立长连接,持续接收增删改事件流,实时更新缓存。
- Store / Indexer:本地缓存,一个带索引的内存数据库,可以按 namespace、label、自定义字段查对象。
- EventHandler:注册回调函数。有对象被添加、更新、删除时,对应回调被触发。
这四个行为拼起来是一套"推拉结合"模型:启动拉一次全量,之后全靠推。Informer 性能好的关键,是它几乎不在热点路径上打 apiserver。读数据优先走本地缓存,计算也基于本地缓存,集群规模越大,这个优势越明显。
3.2 Reflector、DeltaFIFO、Indexer 的分工
Informer 内部有三个角色,名字唬人但职责清楚:
- Reflector:跟 apiserver 打交道的"采购员"。负责 List 和 Watch,收到事件后塞进 DeltaFIFO。
- DeltaFIFO:既能排队又能去重的中间层。每个变更封装成 Delta(变更类型 + 对象),比如 Added、Updated、Deleted、Sync。同一对象的变更会合并,处理顺序先进先出。
- Indexer:处理完的 Delta 被交给 Indexer 更新本地缓存,同时触发回调。Indexer 本质是带索引的内存 map,支持按字段查询。
整个流程可以这样理解:Reflector 负责"菜市场进货",DeltaFIFO 是"备菜区"(先来后到、相同食材合并),Indexer 是"冰箱"(存好随时取),EventHandler 就是"厨师"(食材到了通知你做菜)。
3.3 事件回调与 Resync 机制
有个概念必须搞清楚:Resync(重新同步)。这是新手最容易懵的地方。
Reflector 不只是 Watch,还会周期性触发一次"假更新":把本地缓存里的所有对象重新推一遍 EventHandler,但不会重新 List apiserver。默认周期可以通过 informer factory 的第二个参数配置。
Resync 有两个作用:一是处理事件失败丢掉了变更时,给你一次"对账"机会;二是处理依赖关系的场景,比如你只 watch Deployment,但 Deployment 依赖的 ConfigMap 变了,resync 能让你重新评估。
注意,resync 推的是缓存里的对象,不是 apiserver 最新数据。如果回调里直接处理,拿到的不一定是最新状态,需要自己再 Get 一次,或者用 workqueue 的延迟机制兜底。这也是为什么很多事件回调里还要再查一遍 informer cache——确认对象确实存在,且拿到的是最新版本。
4. 动手写第一个 client-go 程序:从配置到 Informer
4.1 准备环境与依赖
动手之前,先把环境备好。你需要:Go 1.21+(client-go 新版对 Go 版本有要求)、一个能访问的 Kubernetes 集群(本地的 kind、minikube 都行)、集群的 kubeconfig(默认在~/.kube/config)。
新建项目并初始化:
mkdir client-go-demo && cd client-go-demo go mod init client-go-demo拉取依赖时注意,k8s.io/client-go、k8s.io/apimachinery、k8s.io/api这三个模块版本必须严格一致:
go get k8s.io/client-go@v0.29.0 go get k8s.io/apimachinery@v0.29.0 go get k8s.io/api@v0.29.0提示:版本不一致是编译期最常见的坑。我见过太多人只升了 client-go 不升 apimachinery,然后报一堆类型不匹配的错误。更稳妥的做法是直接用
go mod tidy让工具帮你解析。
4.2 加载 kubeconfig 并创建 ClientSet
写main.go,先做基础配置:
package main import ( "fmt" "os" "path/filepath" "k8s.io/client-go/kubernetes" "k8s.io/client-go/tools/clientcmd" "k8s.io/client-go/util/homedir" ) func main() { // 1. 加载 kubeconfig,支持显式指定或使用默认位置 kubeconfig := filepath.Join(homedir.HomeDir(), ".kube", "config") if env := os.Getenv("KUBECONFIG"); env != "" { kubeconfig = env } // 2. 构建 REST 配置 config, err := clientcmd.BuildConfigFromFlags("", kubeconfig) if err != nil { panic(err.Error()) } // 3. 创建 ClientSet clientset, err := kubernetes.NewForConfig(config) if err != nil { panic(err.Error()) } // 4. 验证连通性:先拿一下集群版本 version, err := clientset.Discovery().ServerVersion() if err != nil { panic(err.Error()) } fmt.Printf("Connected to Kubernetes %s\n", version.String()) }这里的几个细节值得说。
BuildConfigFromFlags("", kubeconfig)第一个参数传空,意思是不手动指定 apiserver 地址,完全从 kubeconfig 里读。如果你写的是跑在集群内部的程序(比如 Deployment 里的容器),要换成rest.InClusterConfig(),它会自动读取服务账号挂载的 Token 和 CA 证书。两者差别很大:本地开发用 kubeconfig,集群内运行用 InClusterConfig。我见过不少人在集群里跑本地模式,结果每次都说权限不足——因为 kubeconfig 用的是本机身份,根本不是 pod 里的服务账号。
4.3 用 ClientSet 操作资源(增删改查实战)
配置通了,先写几个最基础的 CRUD,感受一下 typed client。
// 获取 default namespace 下所有 Pod pods, err := clientset.CoreV1().Pods("default").List(context.TODO(), metav1.ListOptions{}) if err != nil { panic(err.Error()) } fmt.Printf("There are %d pods in namespace default\n", len(pods.Items)) // 获取集群所有 Deployment(所有 namespace) deployments, err := clientset.AppsV1().Deployments("").List(context.TODO(), metav1.ListOptions{}) if err != nil { panic(err.Error()) } fmt.Printf("There are %d deployments in cluster\n", len(deployments.Items))注意三个坑:
List()的 namespace 传空字符串表示所有 namespace,传具体名字只列那个 namespace。ListOptions{}可以加字段选择器和标签选择器,比如metav1.ListOptions{LabelSelector: "app=nginx"}只返回打了该标签的对象。这个筛选是 apiserver 执行的,不是本地过滤,大集群里一定要用选择器,别全量拉回来再自己过滤。- List 返回的对象列表是某个时间点的快照,不是持续的。动态监听要靠 informer,别用轮询。
接下来创建、更新、删除一个 Deployment:
// 创建 Deployment deploy := &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ Name: "demo-deploy", Namespace: "default", }, Spec: appsv1.DeploymentSpec{ Replicas: ptr.To[int32](3), // 指针包一下,不然 0 会被当成未指定 Selector: &metav1.LabelSelector{ MatchLabels: map[string]string{"app": "demo"}, }, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: map[string]string{"app": "demo"}, }, Spec: corev1.PodSpec{ Containers: []corev1.Container{ { Name: "demo", Image: "nginx:1.25", }, }, }, }, }, } created, err := clientset.AppsV1().Deployments("default").Create(context.TODO(), deploy, metav1.CreateOptions{}) if err != nil { panic(err.Error()) } fmt.Printf("Created deployment %s\n", created.Name)创建这里有一个特别容易踩的坑:Replicas是指针类型,不能直接填数字。这是 Kubernetes API 的约定,为了区分"没设置"和"设置为 0"。Go 里可以用ptr.To[int32](3)(client-go 提供的工具函数),或者先取变量地址再赋值。
更新和删除:
// 更新:先 Get 出来,改完再整体 Update got, err := clientset.AppsV1().Deployments("default").Get(context.TODO(), "demo-deploy", metav1.GetOptions{}) if err != nil { panic(err.Error()) } got.Spec.Replicas = ptr.To[int32](5) updated, err := clientset.AppsV1().Deployments("default").Update(context.TODO(), got, metav1.UpdateOptions{}) if err != nil { panic(err.Error()) } fmt.Printf("Updated deployment replicas to %d\n", *updated.Spec.Replicas) // 删除 err = clientset.AppsV1().Deployments("default").Delete(context.TODO(), "demo-deploy", metav1.DeleteOptions{}) if err != nil { panic(err.Error()) } fmt.Println("Deleted deployment")这里要特别提醒:Update是整体替换,不是部分更新。如果你拿一个只有名字的新对象去 Update,其他字段(selector、template)会被清空成默认值,apiserver 校验失败或者直接把 Deployment 改坏。所以必须先从集群 Get 出来,改完再整个 Update 回去。想局部更新就用 Patch,后面专门讲。
4.4 接入 Informer:监听资源变化
CRUD 只是热身。接下来写真正有价值的东西:用 Informer 监听 Deployment 的变化。
package main import ( "context" "fmt" "time" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/clientcmd" ) func main() { // ... 同上,构建 clientset ... // 1. 创建 informer factory,第二个参数是 resync 周期 factory := informers.NewSharedInformerFactory(clientset, 30*time.Second) // 2. 获取 Deployment informer deployInformer := factory.Apps().V1().Deployments() // 3. 注册事件回调 _, err := deployInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { deploy := obj.(*appsv1.Deployment) fmt.Printf("[ADD] %s/%s replicas=%d\n", deploy.Namespace, deploy.Name, *deploy.Spec.Replicas) }, UpdateFunc: func(oldObj, newObj interface{}) { oldDeploy := oldObj.(*appsv1.Deployment) newDeploy := newObj.(*appsv1.Deployment) if oldDeploy.ResourceVersion != newDeploy.ResourceVersion { fmt.Printf("[UPDATE] %s/%s replicas: %d -> %d\n", newDeploy.Namespace, newDeploy.Name, *oldDeploy.Spec.Replicas, *newDeploy.Spec.Replicas) } }, DeleteFunc: func(obj interface{}) { deploy, ok := obj.(*appsv1.Deployment) if !ok { // 删除时有时会收到 tombstone(墓碑)对象,需要特殊处理 tombstone, ok := obj.(cache.DeletedFinalStateUnknown) if !ok { return } deploy, ok = tombstone.Obj.(*appsv1.Deployment) if !ok { return } } fmt.Printf("[DELETE] %s/%s\n", deploy.Namespace, deploy.Name) }, }) if err != nil { panic(err.Error()) } // 4. 启动 informer factory.Start(wait.NeverStop) // 5. 等待缓存同步 if !cache.WaitForCacheSync(wait.NeverStop, deployInformer.Informer().HasSynced) { fmt.Println("Failed to sync cache") return } fmt.Println("Cache synced, watching deployments...") // 6. 阻塞主协程 select {} }这段代码有几个关键点:
informers.NewSharedInformerFactory是工厂模式,管理着集群里所有资源的 informer。同一个资源的 informer 全局只有一份,你监听 Deployment 和 ReplicaSet 两个资源时,它们共享同一个 factory,不重复消耗连接和缓存。- 第二个参数
30*time.Second是 resync 周期。注意这只是周期推一遍缓存,不是重新 List。 AddEventHandler里UpdateFunc的 oldObj 和 newObj 是缓存里的新旧版本。resync 时 ResourceVersion 相同的情况会出现,要做过滤避免日志刷屏。- 删除回调里处理 tombstone 是必要的。为什么?因为 informer 缓存可能滞后于 apiserver,如果 apiserver 已删除对象而你的缓存里还有,收到删除事件时对象可能已经被 GC 清理,直接类型断言会 panic。tombstone 机制就是兜底这种情况。
factory.Start(wait.NeverStop)的入参是 stop channel,wait.NeverStop表示不停止。生产环境应该用信号通道做优雅退出。cache.WaitForCacheSync一定要等,否则 informer 本地缓存还没建好,事件回调可能收不全初始事件。
跑起来之后,你随便kubectl scale deployment xxx --replicas=2或者kubectl delete deployment xxx,程序会实时打印对应事件。这就是控制器的心脏部分。
4.5 用 Workqueue 串起事件处理
事件回调里直接干活有个问题:如果某个事件处理特别慢,会阻塞 informer 的事件分发线程,拖垮所有资源的监听。正确做法是用 WorkQueue 把事件先缓存起来,由独立 worker 消费。
来看一个典型的 controller 模式:
import ( "k8s.io/client-go/util/workqueue" "k8s.io/apimachinery/pkg/util/runtime" ) type Controller struct { informer cache.SharedIndexInformer queue workqueue.TypedRateLimitingInterface[string] // 新版用泛型 } func NewController(informer cache.SharedIndexInformer) *Controller { c := &Controller{ informer: informer, queue: workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[string]()), } informer.AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { key, _ := cache.MetaNamespaceKeyFunc(obj) // 生成 "namespace/name" 格式的 key c.queue.Add(key) }, UpdateFunc: func(oldObj, newObj interface{}) { key, _ := cache.MetaNamespaceKeyFunc(newObj) c.queue.Add(key) }, DeleteFunc: func(obj interface{}) { key, _ := cache.DeletionHandlingMetaNamespaceKeyFunc(obj) // 处理 tombstone c.queue.Add(key) }, }) return c } func (c *Controller) Run(ctx context.Context) { defer runtime.HandleCrash() defer c.queue.ShutDown() go c.informer.Run(ctx.Done()) if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) { return } // 启动两个 worker 并发消费 for i := 0; i < 2; i++ { go wait.UntilWithContext(ctx, c.processNextItem, time.Second) } <-ctx.Done() } func (c *Controller) processNextItem(ctx context.Context) bool { key, shutdown := c.queue.Get() if shutdown { return false } defer c.queue.Done(key) obj, exists, err := c.informer.GetIndexer().GetByKey(key.(string)) if err != nil { return false } if exists { fmt.Printf("Handling key: %s\n", key) // 这里做真正的业务逻辑... } c.queue.Forget(key) return true }为什么要用 workqueue 而不是直接在回调里处理事件?三个原因:
- 削峰填谷:apiserver 可能短时间内推送几百个事件,同步处理根本扛不住。队列把积压任务排好,worker 按自己节奏消费。
- 去重:同一对象短时间内收到多次更新事件,workqueue 里只保留一个 key,避免重复处理。
- 重试与限速:处理失败可以
AddRateLimited(key)重新入队,内置限速器指数退避,不会因为出错把 apiserver 打爆。
这四块拼起来,就是标准 controller-runtime 里 controller 的简化版。理解了这套逻辑,再去看 kubebuilder 生成的 Operator 代码,会发现都是这套东西换了个壳。
5. 实用进阶:DynamicClient、Patch 与字段选择器
5.1 DynamicClient 处理 CRD
集群里大概率有大量自定义资源。CRD 没有预定义的 Go 结构(除非你用代码生成器生成),这时候就得靠 DynamicClient 加 Unstructured 对象。
核心用法:
import ( "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/client-go/dynamic" ) func manageCRD(dynamicClient dynamic.Interface) error { // 定义资源类型,注意 schema 的写法 gvr := schema.GroupVersionResource{ Group: "example.com", Version: "v1", Resource: "myresources", } // 用 map 构造一个 Unstructured 对象 obj := &unstructured.Unstructured{ Object: map[string]interface{}{ "apiVersion": "example.com/v1", "kind": "MyResource", "metadata": map[string]interface{}{ "name": "demo", "namespace": "default", }, "spec": map[string]interface{}{ "size": 3, }, }, } // 创建 created, err := dynamicClient.Resource(gvr).Namespace("default"). Create(context.TODO(), obj, metav1.CreateOptions{}) if err != nil { return err } // 读取,拿到的还是 Unstructured got, err := dynamicClient.Resource(gvr).Namespace("default"). Get(context.TODO(), "demo", metav1.GetOptions{}) if err != nil { return err } // 用 NestedInt64 安全读取嵌套字段 size, found, err := unstructured.NestedInt64(got.Object, "spec", "size") if err != nil || !found { return fmt.Errorf("spec.size not found") } fmt.Printf("CRD size: %d\n", size) return nil }这里面最容易出错的是 GVR 写错。CRD 定义文件里有group、version、plural三个字段,GVR 必须跟它们对应。注意 Resource 用的是复数形式。
访问嵌套字段建议用unstructured.Nested*这一组工具函数。别直接断言obj.Object["spec"].(map[string]interface{}),一旦结构不对 panic 直接把程序打崩。工具函数会返回found bool,更安全。
5.2 Patch 和 Update 到底怎么选
前面说 Update 是整体替换,Patch 是局部更新。什么时候用谁?
简单场景(你的控制器只修改一个字段)用 Patch 更合适:逻辑清晰、避免 commit 冲突。三种 Patch 类型:
- Strategic Merge Patch(默认,最常用):Kubernetes 特有的一种合并策略。对 list 字段会按
patchStrategy合并(比如容器数组按 name 合并),而不是直接覆盖。比如给 Deployment 加一个 env 变量,用这个最省心。 - JSON Patch:一个操作数组
[{op: "replace", path: "/spec/replicas", value: 5}],精准指定路径修改、删除字段。 - Merge Patch:标准的 JSON Merge Patch,对 map 直接合并,但 list 是整体替换,一般不推荐用于 Kubernetes 资源。
示例代码:
// Strategic Merge Patch:修改副本数为 5 patchData := []byte(`{"spec":{"replicas":5}}`) _, err := clientset.AppsV1().Deployments("default").Patch( context.TODO(), "demo-deploy", types.StrategicMergePatchType, patchData, metav1.PatchOptions{})再叠加一个标签:
patchData := []byte(`{ "metadata": {"labels": {"tier": "backend"}}, "spec": {"replicas": 5} }`)Patch 大多数时候比 Update 快且不容易踩冲突。但也不是万能:如果改动很多字段,多个 patch 反而比 Update 麻烦。我的经验是:控制器内"以 informer 缓存为读源 + 单字段 Patch 写回"是最稳的组合;批量修改或者对 CRD 做复杂结构更新时再用 Update 或者 DynamicClient。
5.3 字段选择器与标签选择器
ListOptions 里的 FieldSelector 和 LabelSelector 容易被忽略,但在生产环境非常重要。apiserver 对标签选择器的支持很完善,kubectl get pods -l app=nginx就是这么过滤的。client-go 里同样能用:
// 只列出需要处理的资源 opts := metav1.ListOptions{ LabelSelector: "app=nginx", FieldSelector: "metadata.namespace=default", }这比你全量拉数据再本地过滤高效得多。apiserver 端做过滤,网络传输的数据量小一个量级,控制器处理也轻松。特别在几百上千节点的集群里,这个习惯必须养成。
6. 高可用与性能调优:限流、缓存同步、调参指南
6.1 限流:QPS 和 Burst 如何配
apiserver 是集群的中枢神经,你写控制器最怕把自己或者别人打挂。client-go 默认限制 QPS 是 5,Burst 是 10。对小规模集群够用,但对大规模控制器来说太小。可以根据场景调整:
config.QPS = 100 config.Burst = 200QPS 和 Burst 可以类比成地铁闸机:QPS 是常速每分钟能过多少人,Burst 是高峰期一次涌进来多少人允许短暂放行。client-go 的限流是令牌桶算法:平时每秒补 QPS 个令牌,桶最多存 Burst 个令牌。请求来了先取令牌,取不到就排队等待。
我的经验值:
- 简单巡检、偶尔 List 的小工具:默认值就够了。
- 管理几百个 Deployment 的控制器:QPS 50、Burst 100 差不多。
- 大型 Operator(管理几千个 CRD 实例):QPS 200+、Burst 400+,同时要配合多副本选主。
注意 QPS 不要调得太大。apiserver 有自己一套优先级和公平性控制,但你过度压榨它会拖垮整个集群的 kubelet 和其他组件。调参前先看一下 apiserver 监控里的 request latency。
6.2 大规模缓存:Indexer 与字段索引
Informer 本地缓存的容量等于集群对象数量。如果集群有 10 万个 Pod,缓存 map 至少有 10 万个 Pod 对象,内存按百 MB 起步。更可怕的是没用对索引,查一次就全量遍历一次。
Indexer 支持自定义索引。比如按某个 label 的 value 查 Pod:
indexer.AddIndexers(cache.Indexers{ "byLabelApp": func(obj interface{}) ([]string, error) { pod, ok := obj.(*corev1.Pod) if !ok { return []string{}, nil } if app, ok := pod.Labels["app"]; ok { return []string{app}, nil } return []string{}, nil }, })之后用indexer.ByIndex("byLabelApp", "nginx")快速拿到所有带app=nginx的 Pod。查询是纯内存的,对 apiserver 零压力。当你需要在事件处理里频繁查关联资源时,自定义索引能省掉无数个 List 请求。
6.3 多副本部署与 Leader Election
生产环境里控制器通常两个副本以上。不做选主,两个副本同时干活重复处理还互相踩。client-go 的选主机制核心代码:
import ( "k8s.io/client-go/tools/leaderelection" "k8s.io/client-go/tools/leaderelection/resourcelock" ) func runLeaderElection(config *rest.Config, runFunc func(ctx context.Context)) { lock := &resourcelock.LeaseLock{ LeaseMeta: metav1.ObjectMeta{ Name: "my-controller", Namespace: "kube-system", }, Client: clientset.CoordinationV1(), LockConfig: resourcelock.ResourceLockConfig{ Identity: os.Getenv("POD_NAME"), // 每个副本必须唯一 }, } leaderelection.RunOrDie(context.TODO(), leaderelection.LeaderElectionConfig{ Lock: lock, ReleaseOnCancel: true, LeaseDuration: 15 * time.Second, RenewDeadline: 10 * time.Second, RetryPeriod: 2 * time.Second, Callbacks: leaderelection.LeaderCallbacks{ OnStartedLeading: func(ctx context.Context) { runFunc(ctx) }, OnStoppedLeading: func() { // 选举失败或丢主,退出让 Kubernetes 重启 os.Exit(0) }, }, }) }几个经验值:LeaseDuration 默认 15 秒,RenewDeadline 10 秒,RetryPeriod 2 秒,这套参数被大量项目验证过,别随便改。Identity必须每副本唯一,否则选主会乱。
6.4 监控与可观测性
生产环境控制器没监控等于裸奔。至少做到:
- Prometheus metrics:记录 queue 长度、处理耗时、错误率。client-go 自带 workqueue 的 metrics(
workqueue_depth、workqueue_adds_total等),用 controller-runtime 时会自动暴露。手写 client-go 可以引入k8s.io/component-base/metrics,或者简单起一个promhttp.Handler。 - 结构化日志:使用
k8s.io/klog/v2,设置-v=4能看到详细的 HTTP 请求日志,排查认证、权限问题很有帮助。 - Events:用 EventRecorder 往集群写事件,用户
kubectl describe就能看到控制器做了什么。
7. 常见的坑和排错实战记录
7.1 权限不足(Forbidden)问题
症状:运行程序直接报一堆deployments.apps is forbidden: User "system:serviceaccount:default:xxx" cannot list resource "deployments" in API group "apps" at the cluster scope。
根因:服务账号(ServiceAccount)没有对应 RBAC 权限。本地 kubeconfig 用的是当前用户权限,集群内跑的就要给 ServiceAccount 授权。
排查步骤:
- 先确认程序用的什么身份。看报错里 User 字段,如果是
system:serviceaccount,那肯定是集群内运行,走的是 InClusterConfig。 - 检查 ServiceAccount 是否存在:
kubectl get sa -n default。 - 创建对应的 Role/ClusterRole 和 RoleBinding/ClusterRoleBinding。比如给 default 的 sa 加部署资源的读权限:
apiVersion: rbac.authorization.k8s.io/v1 kind: ClusterRole metadata: name: deploy-reader rules: - apiGroups: ["apps"] resources: ["deployments"] verbs: ["get", "list", "watch"] --- apiVersion: rbac.authorization.k8s.io/v1 kind: ClusterRoleBinding metadata: name: deploy-reader-binding subjects: - kind: ServiceAccount name: default namespace: default roleRef: kind: ClusterRole name: deploy-reader apiGroup: rbac.authorization.k8s.io写完kubectl apply -f之后重新部署 pods 才生效。RBAC 变更不会热更新到已运行进程。
经验:给控制器的权限尽量遵循最小权限原则。只给需要读写的资源、动词配权限。我看到很多项目直接给控制器绑cluster-admin,跑是能跑,但一旦容器被入侵,攻击者直接拿到整个集群控制权。这个坏习惯真的别养成。
7.2 ObservedGeneration 与状态回写
控制器还有一个专业细节:更新 Status 子资源。注意 Update 接口不能更新 Status,必须用 UpdateStatus:
_, err := clientset.AppsV1().Deployments("default").UpdateStatus( context.TODO(), deploy, metav1.UpdateOptions{})这里引出ObservedGeneration字段。它用来让用户知道:控制器看到的资源版本跟当前最新版本差多少。如果控制器处理慢,Status 里记录的 ObservedGeneration 小于 Generation,说明状态可能滞后。规范做法是每次处理完资源后,把deploy.Status.ObservedGeneration = deploy.Generation写回去。
还有一个细节:对内置资源的 status 回写,尽量不要改 spec;对 CRD 反而要区分 spec 和 status 两个子资源,有些 CRD 框架(比如 kubebuilder)会自动处理,手写 client-go 时就要自己注意。
7.3 事件丢失与处理失败重试
Informer 事件处理是"以内存为准"的。如果程序崩溃,内存缓存和 event handler 状态都会丢。重启后会重新 List 一次全量,所以事件理论上不会"永久丢失"。但处理过程中失败怎么办?workqueue 里AddRateLimited可以重新排队并带限速,但若一直失败会无限重试,直到队列满——注意,限速队列严格来说不是无限重试,它会指数退避到最大值后一直以固定间隔重试。所以你要自己加一个最大重试次数或者超时逻辑,超过阈值就上报错误并丢弃事件。
我见过一个真实案例:一个控制器处理某 CRD 时因为状态字段结构不对,一处理就 panic,workqueue 不断重试,最后把内存和日志全吃满。加了一个重试 5 次就放弃并写 Event 的逻辑后,问题立刻缓解。
7.4 版本兼容问题
client-go 版本和集群版本不一致,最常见的表现是某些字段不认识、某些 API 版本不存在。比如本地用 v0.29 连一个 v1.23 的集群,不一定会立即报错,但某些新 API 字段会被丢弃。
经验做法:
- 开发、测试、生产环境的集群版本尽量统一。
- 升级时先看 client-go 仓库的 compatibility matrix(每个 release 都标注了支持的 Kubernetes 版本区间)。
- 如果你的控制器发布成二进制给不同集群用,注意不要把 client-go 版本对应的 API 行为差异引入到同一个二进制里。
7.5 Stop channel 与优雅退出
写生产代码时别用wait.NeverStop当永久不退出。应该用signal.NotifyContext捕获 SIGTERM、SIGINT:
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() // 传给 informer 和 worker factory.Start(ctx.Done()) <-ctx.Done()这样 Pod 被删除时,Kubernetes 先发 SIGTERM,程序能优雅关停:停止接收新事件、把队列里的任务尽力处理完、刷新缓存再退出。用NeverStop的话,Kubernetes 等不到进程退出,最后只能 SIGKILL,控制器运行时状态可能不一致。
- 什么场景必须用 DynamicClient?如果你操作的资源是 CRD 且没有用代码生成器生成 client,就必须用它。判断标准很简单:你的代码里有没有一个类型能对应到那个 CRD 的具体结构。没有,就只能用
unstructured。 - WaitForCacheSync 卡住怎么办?先确认你的 RBAC 权限有没有 list/watch。Reflector 启动时会先 List,权限不足会一直报错。
- Update 和 Patch 的选择:如果只想改一个字段,无脑用 Patch。如果想做复杂验证或者批量改,再用 Update。CRD 的 status 更新别忘了用 UpdateStatus。
- 为什么 informer 缓存的对象不能直接修改?informer 缓存里的对象是共享的,你在事件回调里拿到的是指针,直接改会污染缓存。要用
obj.DeepCopy()复制一份再改。
8. 从 client-go 到 controller-runtime:下一站
写到这里,对 client-go 应该有个比较全面的认识了。最后聊聊它和 controller-runtime(也就是 kubebuilder 背后的库)的关系。
controller-runtime 本质上是 client-go 的"亲儿子",在 client-go 基础上又封装了一层:
- 管理多个 controller 的生命周期,每个 controller 对应一个 reconciler 函数。
- 自动处理 informer 工厂管理、缓存同步、事件分发。
- 内置更完善的 RBAC 注解(kubebuilder 里的
+kubebuilder:rbac:groups=...),代码生成器帮你生成权限配置。 - 集成了 webhook、metrics、leader election 等生产级功能。
如果你的目标是写大型 Operator,直接用 controller-runtime 的脚手架更高效。但我不建议跳过 client-go 直接学 controller-runtime。因为 controller-runtime 的抽象太优雅了,优雅到你不知道它在底下干了什么。一旦遇到性能瓶颈、奇怪的事件重复、缓存不一致,扒开源码看到的还是 informer、workqueue、cache.Indexer 这套东西。地基扎实,上层才不会塌。
我在实际运维和开发里最大的体会是:client-go 是一个值得花几个晚上把核心机制啃透的库。那些限流、缓存、队列、选主、重试的设计,放到任何需要跟外部系统高效交互的场景里都是通用的。把它读明白了,写任何大规模分布式系统的客户端,脑子里都会有非常清晰的架构感。
建议下一步做两件事:一是把今天的 informer + workqueue 代码跑起来,在测试集群里亲手演练增删改和崩溃恢复;二是找一份 kubebuilder 生成的 Operator 代码,跟本文讲的这套结构对照着看,你会发现之前看不懂的地方全都清晰了。如果这篇文章帮你少走了一段弯路,那这些时间就花得值了。