前置:看完阶段 0 心智模型和client-go 与 CRD 教学。 本文对应学习路线的阶段 1,四步走: ① Clientset 列 Pod → ② Informer 打事件 → ③ 加 Workqueue → ④ 写出真正的 Reconcile。 阶段产出:一个约 180 行的
main.go,能 watch ConfigMap 并保证my-config永远存在——删掉自动重建、改坏自动纠正。 本文代码与 client-go 官方示例examples/workqueue对齐,使用当前(2026)的泛型 API;网上老教程的workqueue.NewRateLimitingQueue已弃用,别抄(文末有对照表)。
01〇、开工前准备(10 分钟)
你需要三样东西:
| 东西 | 要求 | 检查命令 |
|---|---|---|
| 一个 K8s 集群 | 能 kubectl 就行(kind / minikube / 测试集群均可) |
kubectl get nodes |
| Go | 1.22+ | go version |
| client-go 版本 | 和集群版本 ±1(教学文档第二节的规则) | kubectl version 看 Server Version |
建项目:
mkdir my-controller && cd my-controller
go mod init my-controller
# 规则:client-go 的 v0.X 对应 K8s 1.X,版本号对齐(允许 ±1,同版本最稳)
# 例如集群是 1.36 → 全用 v0.36.x;集群是 1.35 → 把下面的 36 改成 35
go get k8s.io/client-go@v0.36.0
go get k8s.io/api@v0.36.0
go get k8s.io/apimachinery@v0.36.0
go get k8s.io/klog/v2@latest # klog 只是日志库,不跟集群版本走全程只有一个
main.go文件,每一步都是完整可跑的程序,改一点、跑一点、看懂一点——别跳着抄最终版。
02第一步:Clientset 列出 Pod(15 行,热身)
目标:确认"代码能连上集群、能读东西"。新建 main.go:
package main
import (
"context"
"flag"
"fmt"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/clientcmd"
)
func main() {
kubeconfig := flag.String("kubeconfig", "", "kubeconfig 路径,留空自动找 ~/.kube/config")
flag.Parse()
// ① 怎么连(教学文档第二节:本地开发用 kubeconfig)
config, err := clientcmd.BuildConfigFromFlags("", *kubeconfig)
if err != nil {
panic(err)
}
// ② 拿到客户端
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
panic(err)
}
// ③ 链式调用:group → version → resource(教学文档第四节)
pods, err := clientset.CoreV1().Pods("default").List(context.TODO(), metav1.ListOptions{})
if err != nil {
panic(err)
}
for _, p := range pods.Items {
fmt.Println(p.Name)
}
}跑:
go run main.go
# 打出 default 命名空间下所有 Pod 的名字三个数字对应的零件,你在教学文档里都见过:BuildConfigFromFlags 是"两种连法"之一,NewForConfig 一行拿客户端,CoreV1().Pods().List() 是链式调用。唯一的认知提醒:这次 List 是你主动问 API Server 要的——下一秒它就过时了。轮询的办法阶段 0 已经否决了,所以下一步换 Informer。
03第二步:Informer 打印 Pod 事件(看到"订阅"长什么样)
目标:不主动问了,让事件推过来,亲眼看 Add/Update/Delete。
整体替换 main.go:
package main
import (
"context"
"flag"
"fmt"
"os"
"os/signal"
"syscall"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/tools/clientcmd"
)
// DeletionHandling 版本对普通对象和"删除墓碑"都能取出 key,三个回调通用
func keyOf(obj interface{}) string {
key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
if err != nil {
return "?"
}
return key
}
func main() {
kubeconfig := flag.String("kubeconfig", "", "kubeconfig 路径")
flag.Parse()
config, err := clientcmd.BuildConfigFromFlags("", *kubeconfig)
if err != nil {
panic(err)
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
panic(err)
}
// 按 Ctrl-C 时优雅退出(后面每一步都这么干)
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
// SharedInformerFactory:一种资源全进程共享一份 watch + 缓存(教学文档 6.1)
// 第二个参数 resync=0:关闭定时巡检,只要"真事件"
factory := informers.NewSharedInformerFactory(clientset, 0)
podInformer := factory.Core().V1().Pods().Informer()
// 固定动作①:注册回调(现在先打印,下一步改成放 key)
podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { fmt.Printf("[ADD] %s\n", keyOf(obj)) },
UpdateFunc: func(old, new interface{}) { fmt.Printf("[UPDATE] %s\n", keyOf(new)) },
DeleteFunc: func(obj interface{}) { fmt.Printf("[DELETE] %s\n", keyOf(obj)) },
})
factory.Start(ctx.Done()) // 启动所有 informer(非阻塞)
// 固定动作②:等缓存同步完
if !cache.WaitForCacheSync(ctx.Done(), podInformer.HasSynced) {
panic("缓存同步超时")
}
fmt.Println("=== 缓存同步完成,开始监听(去另一个终端操作 Pod)===")
<-ctx.Done()
}实验:开两个终端。
# 终端 1
go run main.go
# 终端 2(程序跑起来之后)
kubectl run nginx --image=nginx
kubectl label pod nginx env=test
kubectl delete pod nginx你会看到 [ADD] default/nginx、[UPDATE] default/nginx、[DELETE] default/nginx 依次出现。
启动瞬间的一个现象,务必看懂:程序刚跑起来时,会刷出一堆 [ADD]——集群里已有的每个 Pod 都触发了一次 Add。这不是有人批量建 Pod,是 informer 启动时 LIST 全量灌缓存,缓存里"从无到有"的每个对象都算一次 Add。这就是阶段 0"电平触发"的现场版:controller 一睁眼先看到全量现状,而不是从"什么都没有"开始。 你以后写的 reconcile,第一次执行就是在这波初始 Add 里跑完的。
04第三步:加 Workqueue(回调从"打印"退化成"放纸条")
目标:回调不再干活,只把 key 放进队列;另一端开 worker 慢慢处理。
在第二步基础上改三处:
① 创建队列(注意:这是当前的泛型写法):
import "k8s.io/client-go/util/workqueue"
queue := workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[string]())② 三个回调全部退化成放 key:
podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { queue.Add(keyOf(obj)) },
UpdateFunc: func(old, new interface{}) { queue.Add(keyOf(new)) },
DeleteFunc: func(obj interface{}) { queue.Add(keyOf(obj)) },
})③ 缓存同步完成后,启动 worker 消费:
fmt.Println("=== 缓存同步完成,worker 开工 ===")
for i := 0; i < 2; i++ { // 2 个工人
go func() {
for processNextItem(queue) {
}
}()
}worker 骨架(这个结构背下来,所有裸写 controller 都长这样):
func processNextItem(queue workqueue.TypedRateLimitingInterface[string]) bool {
key, quit := queue.Get() // 阻塞等一张纸条
if quit {
return false
}
defer queue.Done(key) // 必须 Done,否则队列认为你永远在处理它
err := business(key) // ← 你的业务逻辑(现在还是打印)
switch {
case err == nil:
queue.Forget(key) // 成功:清掉这个 key 的退避计数(千万别漏!)
case queue.NumRequeues(key) < 5:
queue.AddRateLimited(key) // 失败:指数退避后重试
default:
queue.Forget(key) // 连挂 5 次:放弃,记日志告警
fmt.Printf("放弃 %s: %v\n", key, err)
}
return true
}
func business(key string) error {
fmt.Printf("处理: %s\n", key)
time.Sleep(2 * time.Second) // 故意放慢,为了做下面的去重实验
return nil
}实验(亲眼看到"去重"和"拿最新状态"):worker 故意 sleep 2 秒,然后你飞快连打 5 次 label:
for i in 1 2 3 4 5; do kubectl label pod nginx a=$i --overwrite; done事件来了 5 次,但"处理"只打印了 2 次左右,不是 5 次——worker 慢吞吞处理第一张纸条时,后面 4 个事件在队列里合并成了一张(教学文档第七节:队列里放 key 的第一个原因)。等 worker 再取到它时,去缓存里读到的已经是 a=5 的最新状态。电平触发:中间过程全被压缩掉了,只看最终状态。
05第四步:真正的 Reconcile——保证 my-config 永远存在
前三步是零件练习,这一步把它们拧成真正的 controller。业务目标(就是学习路线阶段 1 的产出要求):
期望状态:
default下必须有一个叫my-config的 ConfigMap,内容是固定的。 它不存在 → 创建;内容被人改了 → 改回来。
整体替换为最终版(约 180 行,注释即讲解):
package main
import (
"context"
"flag"
"os"
"os/signal"
"reflect"
"syscall"
"time"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
corelisters "k8s.io/client-go/listers/core/v1"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/tools/clientcmd"
"k8s.io/client-go/util/workqueue"
"k8s.io/klog/v2"
)
// ========== 期望状态:这个 controller 唯一的"业务知识" ==========
const (
targetNamespace = "default"
targetName = "my-config"
)
var desiredData = map[string]string{
"app.conf": "env=prod\nreplicas=3",
}
type Controller struct {
clientset kubernetes.Interface
cmLister corelisters.ConfigMapLister // 读:走缓存(教学文档 6.3)
cmSynced cache.InformerSynced
queue workqueue.TypedRateLimitingInterface[string]
}
func main() {
klog.InitFlags(nil)
kubeconfig := flag.String("kubeconfig", "", "kubeconfig 路径")
flag.Parse()
config, err := clientcmd.BuildConfigFromFlags("", *kubeconfig)
if err != nil {
klog.Fatalf("加载 kubeconfig 失败: %v", err)
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
klog.Fatalf("创建 clientset 失败: %v", err)
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
// 只看 default;resync=30s:方便你亲眼看到"周期性对账"(生产里通常是小时级)
factory := informers.NewSharedInformerFactoryWithOptions(clientset, 30*time.Second,
informers.WithNamespace(targetNamespace))
cmInformer := factory.Core().V1().ConfigMaps()
informer := cmInformer.Informer()
c := &Controller{
clientset: clientset,
cmLister: cmInformer.Lister(),
cmSynced: informer.HasSynced,
queue: workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[string]()),
}
// 固定动作①:回调只放 key
informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { c.enqueue(obj) },
UpdateFunc: func(old, new interface{}) { c.enqueue(new) },
DeleteFunc: func(obj interface{}) { c.enqueue(obj) },
})
factory.Start(ctx.Done())
// 固定动作②:缓存不同步完,worker 绝不开工
if !cache.WaitForCacheSync(ctx.Done(), c.cmSynced) {
klog.Fatal("缓存同步超时")
}
klog.Info("缓存同步完成")
// ★ 主动踢一脚:my-config 可能压根不存在,不存在就永远不会有它的事件,
// reconcile 就永远跑不到它。启动时手动入队,保证对账至少跑一次。
c.queue.Add(targetNamespace + "/" + targetName)
for i := 0; i < 2; i++ {
go wait.UntilWithContext(ctx, c.runWorker, time.Second)
}
<-ctx.Done()
klog.Info("收到退出信号,停机")
c.queue.ShutDown()
}
func (c *Controller) enqueue(obj interface{}) {
key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
if err == nil {
c.queue.Add(key)
}
}
// ---------- worker 骨架:和第三步一模一样 ----------
func (c *Controller) runWorker(ctx context.Context) {
for c.processNextItem(ctx) {
}
}
func (c *Controller) processNextItem(ctx context.Context) bool {
key, quit := c.queue.Get()
if quit {
return false
}
defer c.queue.Done(key)
err := c.reconcile(ctx, key) // ← 唯一的业务逻辑在这
switch {
case err == nil:
c.queue.Forget(key)
case c.queue.NumRequeues(key) < 5:
klog.Infof("对账失败,稍后重试 (%s): %v", key, err)
c.queue.AddRateLimited(key)
default:
klog.Errorf("对账连挂 5 次,放弃 (%s): %v", key, err)
c.queue.Forget(key)
}
return true
}
// ---------- reconcile:你 90% 的工作量(以后)都在这种函数里 ----------
func (c *Controller) reconcile(ctx context.Context, key string) error {
namespace, name, err := cache.SplitMetaNamespaceKey(key)
if err != nil {
return err
}
// 不是我们的目标对象,看一眼就放行(别的事件也会进来,这很正常)
if namespace != targetNamespace || name != targetName {
return nil
}
// 读:走缓存,不问 API Server
cm, err := c.cmLister.ConfigMaps(namespace).Get(name)
switch {
case apierrors.IsNotFound(err):
// 期望存在,实际不存在 → 创建(这就是"自愈")
klog.Infof("%s 不存在,创建中...", key)
desired := &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: namespace},
Data: desiredData,
}
_, err = c.clientset.CoreV1().ConfigMaps(namespace).Create(ctx, desired, metav1.CreateOptions{})
return err
case err != nil:
return err
}
// 存在 → 对比内容,不一致才纠正(这就是"补差")
if reflect.DeepEqual(cm.Data, desiredData) {
klog.V(1).Infof("%s 与期望一致,无需动作", key)
return nil // ★ 一致就什么都不写 —— 防 hot loop 的关键
}
klog.Infof("%s 内容偏离期望,纠正中...", key)
cmCopy := cm.DeepCopy() // ★ 缓存里的对象是共享的,改前必须 DeepCopy!
cmCopy.Data = desiredData
_, err = c.clientset.CoreV1().ConfigMaps(namespace).Update(ctx, cmCopy, metav1.UpdateOptions{})
return err
}跑起来,做四个实验:
# 实验 1:启动即创建
go run main.go
# 日志: default/my-config 不存在,创建中...
kubectl get cm my-config -o yaml # 已经有了
# 实验 2:自愈 —— 删掉它
kubectl delete cm my-config
# 日志立刻: default/my-config 不存在,创建中...
kubectl get cm my-config # 又回来了 ★ 这就是学习路线的阶段验收
# 实验 3:纠偏 —— 手贱改内容
kubectl patch cm my-config --type merge -p '{"data":{"app.conf":"env=dev"}}'
# 日志: default/my-config 内容偏离期望,纠正中...
kubectl get cm my-config -o jsonpath='{.data.app\.conf}' # 被改回 env=prod
# 实验 4:幂等 —— Ctrl-C 后再跑一次
go run main.go
# 什么"创建"日志都没有:已存在且一致,直接收工。跑 10 遍 = 跑 1 遍。再加 go run main.go -v=1 等 30 秒,会看到周期性的"与期望一致,无需动作"——那就是 resync 定时巡检在对账,教学文档 6.2 第③条的现场版。
这段代码里藏着的四个"为什么"(逐个钉死)
① 为什么启动要"踢一脚"(queue.Add)? 事件只描述"变化",不描述"应该"。my-config 如果从来没存在过,就不会有它的任何事件,reconcile 永远轮不到它。启动后手动入队一次,保证对账至少跑一遍。——顺便预告:阶段 2 有了 CRD 之后,"期望状态"会变成集群里的一个 CR 对象,它一被创建就有 Add 事件,就不用踢这一脚了。这就是 CRD 存在的意义之一。
② 为什么不会 hot loop? 我们的 Create/Update 会触发新事件 → 又入队 → 又 reconcile——但下一圈读到缓存发现"已一致",直接 return nil 什么都不写,循环就安静了。写之前先对比,一致绝不写——阶段 0 的闭环铁律,落实为 reflect.DeepEqual 这一个判断。
③ 为什么改之前要 DeepCopy()? lister 返回的指针指向共享缓存里的那一份,直接改它的字段,所有其他订阅者读到的都是你的脏数据,而且这种污染不触发任何事件、不报任何错,是最阴的隐形 bug。规矩:从缓存读出来的对象,想改先 DeepCopy。
④ 为什么"读"和"写"用两个对象? cmLister.Get 走缓存(每 30 秒 100 个 key 对账也不惊动 API Server);clientset.Create/Update 直连 API Server。教学文档 6.3 的"读走内存,写才走 API Server",在这一个 reconcile 里齐了。
06五、把代码对回阶段 0 那张图
| 阶段 0 图里的零件 | 本文代码里的谁 |
|---|---|
| API Server 的变化通知(watch) | factory.Start → informer 内部的 Reflector |
| 本地缓存(内存拷贝) | informer 内部的 indexer,cmLister 是它的只读窗口 |
| 小纸条(只有名字) | key:字符串 "default/my-config" |
| Workqueue(去重/重试/节奏) | c.queue |
| 工人 worker | 2 个 runWorker goroutine |
| Reconcile(你写的对账) | c.reconcile() |
| 你的运维知识 | desiredData + reconcile 里的三段判断(没有→建 / 偏了→改 / 一致→收工) |
阶段 0 说的"你 90% 的工作量是写对账函数",现在你有体感了:刨掉连接、信号、worker 骨架(全是模板),真正动脑子的只有 desiredData 和 reconcile 那 30 行。
07六、踩坑清单(本阶段你一定会遇到的)
| 症状 | 原因 | 解法 |
|---|---|---|
编译错 cannot use queue.Add(key) 之类 |
抄了老教程,新旧 workqueue API 混用 | 见下方对照表,全套用泛型版 |
| controller 行为诡异,数据互相串 | 直接改了 lister 返回的对象 | 改前 DeepCopy() |
| 失败后重试越来越慢 | 成功路径忘了 Forget(key) |
重看 worker 骨架 switch |
| 启动时把没同步当不存在,误建/误删 | 没等 WaitForCacheSync 就开工 |
固定动作②不能省 |
configmaps is forbidden |
集群内跑,ServiceAccount 没权限 | 见下方折叠 RBAC |
connection refused / 找不到集群 |
kubeconfig 路径错,或走了 InCluster | 本地跑传 --kubeconfig 或留空 |
Update 报 409 Conflict |
乐观锁:别人先改了(resourceVersion 过期) | 不是 bug,返回 err 让队列重试,下圈读到新的就好 |
老教程 API 对照表(搜到的中文教程大多还停在左边):
| 老教程(已弃用,别抄) | 现在(v0.31+,本文用法) |
|---|---|
workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter()) |
workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[string]()) |
queue.Get() 返回 interface{},要 key.(string) 强转 |
泛型,Get() 直接返回 string |
workqueue.RateLimitingInterface |
workqueue.TypedRateLimitingInterface[string] |
折叠:以后要部署进集群时,RBAC 最小权限长这样(本地 kubeconfig 是 admin,现在不需要)
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: my-controller
namespace: default
rules:
- apiGroups: [""]
resources: ["configmaps"]
verbs: ["get", "list", "watch", "create", "update"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: my-controller
namespace: default
subjects:
- kind: ServiceAccount
name: my-controller
roleRef:
{ kind: Role, name: my-controller, apiGroup: rbac.authorization.k8s.io }get/list/watch 是 informer 要的,create/update 是 reconcile 要的——RBAC 需求可以从代码里推出来,这是阶段 3 的伏笔。
08七、我们的玩具和"真家伙"差在哪(下一阶段预告)
官方 sample-controller 就是这套骨架的完整版,对照看差距:
| 本文的玩具 | 生产 controller(sample-controller / kubebuilder) |
|---|---|
| 期望状态写死在代码里,启动要"踢一脚" | 期望状态是集群里的 CR 对象(CRD),它的创建本身就是事件 |
| 只盯一种内置资源 | 自己的 CRD + Owns() 盯子资源,子资源变化反查 owner 入队 |
| 删除只靠默认行为 | Owner Reference 级联删除 + Finalizer 清理外部资源 |
| 出事只打日志 | Events: kubectl describe 能看到 |
| 单副本裸奔 | leader election,多副本选主 |
| 没有 status | /status 子资源回写观测状态 |
右边这列,就是学习路线阶段 2-3 要补的课。但你注意:骨架一字未变——informer、workqueue、reconcile 三段式到 kubebuilder 里还是它,只是被框架藏起来了。这就是"先裸写再上框架"的原因。
09八、学完自检(6 题)
Q1. 启动时为什么会刷出一堆 [ADD] 事件?
informer 启动先 LIST 全量灌缓存,已有对象"从无到有"各触发一次 Add。这是电平触发的现场版:对账从全量现状开始。
Q2. worker 处理很慢时,一个对象连变 5 次,会处理几次?为什么?
远少于 5 次。队列里放的是 key,key 唯一,未处理期间重复入队被合并(去重);worker 再取到它时从缓存读的是最新状态。中间过程被压缩,只看最终态。
Q3. 为什么启动时要 queue.Add("default/my-config") 踢一脚?
事件只描述"变化"不描述"应该"。目标对象从没存在过就没有它的事件,reconcile 永远轮不到它;手动入队保证对账至少跑一次。阶段 2 用 CRD 后,CR 的创建本身就是事件,不用踢了。
Q4. 从 lister 读出的对象,改之前必须做什么?为什么?
DeepCopy()。它指向共享缓存,直接改会污染所有订阅者读到的数据,且不触发事件、不报错,是隐形 bug。
Q5. 这个 controller 为什么不会 hot loop?
每次写(Create/Update)虽会触发新事件再对账,但下圈发现"缓存已一致"就直接返回、绝不再写。写之前先对比(DeepEqual),一致不写,循环收敛。
Q6. Update 返回 409 Conflict 怎么办?
不用处理成错误——这是乐观锁(resourceVersion)生效:你读完后有人先改了。返回 err 让队列退避重试,下一轮从缓存读到新版本再写即可。
10九、参考资料
官方(代码与本文对齐,优先):
- client-go/examples/workqueue — 本文骨架的出处
- kubernetes/sample-controller — 完整形态,阶段 1 收尾时通读
- community/controllers.md — controller 写法的官方设计说明
- client-go tags — 按集群版本选 client-go 版本
中文(思路可看,API 偏旧,注意对照第六节弃用表):
- client-go实战之九:手写一个 kubernetes 的 controller(程序员欣宸) — 含 Reflector/DeltaFIFO 内部结构图
- client-go Informer、Reflector、Indexer、Workqueue 介绍(Jimyag) — 数据流讲得细
下一步:把第四步的最终版亲手敲一遍、四个实验都做出来(尤其 kubectl delete 后看到自动重建),阶段 1 就毕业了——然后进阶段 2:kubebuilder,看框架怎么把你今天手写的这些全藏起来。