Kubernetes 基础体系 · 第 77/83 篇。示例基于 Kubernetes 当前稳定 API;弃用、版本偏差、云厂商差异和生产风险会明确说明。
Kubernetes Controller 开发:client-go、Cache、Queue、Reconcile 和幂等
Kubernetes Controller 是一种持续运行的控制程序:它观察 Kubernetes API 中的对象,将观察到的状态与期望状态进行比较,然后执行有限的 API 操作,使实际状态逐步接近期望状态。
一个典型的 client-go Controller 可以抽象为:
API Server
│
├── List/Watch
▼
SharedInformer ──> 本地 Cache ──> Event Handler ──> WorkQueue
│
▼
Reconcile
│
┌────────────────┴───────────────┐
▼ ▼
读取 Cache 写入 API Server
其中:
- client-go 提供访问 Kubernetes API、Informer、Cache、WorkQueue 等基础组件;
- Cache 保存 Informer 从 API Server 同步到的本地对象快照;
- Queue 保存“需要重新处理的对象键”,而不是完整对象;
- Reconcile 根据当前观察状态计算并执行期望状态;
- 幂等 保证同一个对象被重复调谐不会不断产生无意义变化或破坏状态。
这几个概念不是相互独立的工具。Informer 决定了事件如何进入系统,Queue 决定了事件如何合并和重试,Reconcile 决定了控制逻辑,幂等性决定了整个系统能否在重复事件、延迟、重试和并发下稳定收敛。
一、先建立 Controller 的状态模型
1. 观察状态、期望状态和实际状态
对某个 Kubernetes 资源,可以区分三种状态:
- 期望状态(Desired State):用户或上层控制器声明希望系统达到的状态;
- 实际状态(Actual State):集群中相关资源当前真实存在的状态;
- 观察状态(Observed State):Controller 通过 Cache 或 API 读取到的状态。
在理想情况下:
Observed State = Actual State
但在分布式系统中,这个等式并不总是立即成立:
- API Server 写入成功后,Informer 可能尚未收到事件;
- Cache 可能暂时落后;
- 外部系统可能在 Controller 读取后再次修改对象;
- 多个 Controller 可能同时更新同一资源。
因此,Reconcile 不应假设一次读取就获得永久正确的全局状态。它应当在每次运行时重新计算,并允许后续事件再次触发。
设:
- 表示期望状态;
- 表示第 次调谐时观察到的状态;
- 表示调谐动作;
- 表示动作执行后系统最终呈现的状态。
可以写成:
Controller 的目标不是“处理某一个事件”,而是让系统逐渐满足:
这里的“约等于”不是 Go 中的指针相等,而是业务语义上的一致。例如 Deployment 的副本数、Pod 模板和选择器达到期望关系。
2. 为什么事件不是调谐输入本身
一个常见误解是:
收到一个 Pod 更新事件,就直接根据事件对象修改另一个资源。
这会把事件误认为完整事实。实际上,事件更像一个提示:
“某个对象附近的状态可能发生了变化,请重新检查。”
事件可能:
- 被合并;
- 重复到达;
- 到达时对象已经发生下一次变化;
- 到达时 Cache 仍未更新;
- 因删除事件而只剩下墓碑对象;
- 因 Controller 重启而完全丢失。
因此,事件处理器通常只把对象转换成队列键:
namespace/name
真正的对象读取和状态判断发生在 Reconcile 中。
二、client-go 提供了哪些基础能力
client-go 是 Kubernetes 官方 Go 客户端库。它包含多层能力:
Typed Client
├── CoreV1().Pods(...)
├── AppsV1().Deployments(...)
└── ...
Informer / SharedInformerFactory
├── List/Watch
├── Local Store
├── Event Handlers
└── Resync
WorkQueue
├── 去重
├── 并发安全
├── 延迟入队
└── 限速重试
1. Typed Client
Typed Client 为 Kubernetes 内置资源提供类型安全的访问接口:
pods, err := clientset.CoreV1().
Pods("default").
Get(ctx, "demo", metav1.GetOptions{})
写入对象时:
_, err := clientset.CoreV1().
ConfigMaps("default").
Create(ctx, cm, metav1.CreateOptions{})
Typed Client 的优点是:
- 编译期类型检查;
- 使用 Kubernetes API 类型;
- 对
Get、Create、Update、Patch、Delete等操作有清晰接口; - 能正确携带
ResourceVersion、UID等元数据。
它不负责:
- 自动缓存;
- 自动重试所有业务冲突;
- 自动实现期望状态;
- 自动保证幂等。
这些属于 Controller 逻辑或其他 client-go 组件的职责。
2. REST Client 与动态客户端
除了 Typed Client,client-go 还提供:
- REST Client:更底层,直接构造 API 请求;
- Dynamic Client:使用
unstructured.Unstructured访问未知或动态资源。
对于已知的内置资源,Typed Client 通常更适合。对于需要同时支持未知 CRD、多个版本或通用资源处理器的程序,Dynamic Client 可能更合适,但代价是类型检查减少,字段访问容易出现运行时错误。
本文重点使用 Typed Client 和 Informer。
三、Informer、Cache 与 List/Watch
1. Informer 的工作方式
一个 SharedInformer 通常经历以下流程:
启动
│
├── List:获取当前资源集合
│
├── 将对象写入本地 Store
│
├── 触发 Add 事件
│
├── Watch:持续接收新增、修改、删除事件
│
└── Watch 断开时重新建立 List/Watch
其中:
- Store 是本地对象存储;
- Indexer 是带索引能力的 Store;
- SharedInformer 将 List/Watch、Store 和事件分发组合起来;
- SharedInformerFactory 可以让同一进程中的多个 Controller 共享 Informer 和 Cache。
共享 Informer 的核心价值是避免每个 Controller 都单独执行一套 List/Watch。多个消费者可以共享一次 API Server 观察流和本地缓存。
2. Cache 不是 API Server 的实时镜像
Cache 是异步同步的本地快照,而不是强一致读取接口。
时间线可能如下:
t0: Controller 读取 Cache,看到 replicas=1
t1: 另一个客户端把 replicas 改为 3
t2: Controller 仍可能从 Cache 读到 replicas=1
t3: Informer 收到 Watch 事件并更新 Cache
t4: Controller 再次读取 Cache,看到 replicas=3
因此:
- 使用 Cache 读取通常更高效;
- Cache 读取可能暂时落后;
- 写入必须通过 Client 发往 API Server;
- 不应直接修改 Cache 中的对象;
- 需要强制读取最新 API 状态时,可以调用 Typed Client,但这增加 API Server 压力,也不能消除并发冲突。
3. Cache 中的对象不能直接修改
Informer Cache 中的对象通常是共享对象。错误示例:
obj, exists, err := informer.GetStore().GetByKey(key)
if err != nil || !exists {
return err
}
cm := obj.(*corev1.ConfigMap)
cm.Labels["managed-by"] = "demo-controller" // 错误
这里修改的是 Cache 中的对象,可能导致:
- 其他处理逻辑看到未经 API Server 确认的状态;
- 数据竞争;
- 后续 DeepEqual 比较失真;
- Cache 与 API Server 状态脱节。
如果必须构造修改对象,应先复制:
desired := cm.DeepCopy()
desired.Labels = maps.Clone(cm.Labels)
更常见的做法是:从 Cache 读取对象,只提取必要字段;创建或更新另一个独立对象时,使用 API Server 返回的最新对象。
4. Cache 同步完成后才能启动 Worker
Informer 启动后,Cache 不是立即完整的。Controller 通常必须等待:
if !cache.WaitForCacheSync(stopCh, informer.HasSynced) {
return
}
如果不等待就开始 Reconcile,可能把“Cache 中尚未出现”误判成“资源不存在”,从而:
- 错误创建重复资源;
- 错误删除资源;
- 错误更新状态;
- 在启动期间产生大量无效重试。
HasSynced 表示 Informer 完成了初始同步,但不表示之后永远没有延迟,也不表示所有相关资源都已经处于业务稳定状态。
四、事件处理器:只负责入队
Informer 的事件处理器通常不执行复杂业务逻辑,而是完成:
事件对象 → namespace/name → Queue
1. Add 和 Update
informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
enqueue(obj)
},
UpdateFunc: func(oldObj, newObj interface{}) {
enqueue(newObj)
},
})
如果 Update 事件每次都入队,可能增加无效调谐,但通常比在事件处理器中进行复杂判断更安全。更精细的 Controller 可以比较 ResourceVersion、Spec 或其他字段,只在相关字段变化时入队。
必须注意:ResourceVersion 变化不等于业务字段变化。Status 更新、其他控制器更新 Annotation,都可能改变 ResourceVersion。
2. Delete 与 Tombstone
删除事件不一定直接携带资源对象。Informer 可能传递 DeletedFinalStateUnknown,这被称为 Tombstone:
func enqueue(obj interface{}) {
key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
if err != nil {
return
}
queue.Add(key)
}
DeletionHandlingMetaNamespaceKeyFunc 同时处理:
- 正常对象;
- 删除墓碑对象。
删除事件进入 Queue 后,Reconcile 再从 Cache 读取时通常会得到 NotFound。这个 NotFound 不是异常,而是正常控制路径:
对象存在 → 读取并确保子资源
对象删除 → Cache 读不到 → 清理或退出
3. 事件与队列的关系
如果同一个键连续收到多次事件:
default/demo
default/demo
default/demo
Queue 通常不会让 Worker 必须连续处理三次完全相同的任务。它具有去重语义:同一个键已经在队列中时,再次入队不会无限复制相同任务。
但这不意味着事件数量完全没有影响:
- 如果第一次任务已经被取出,后续事件仍可能再次入队;
- 业务写入可能产生新的事件;
- 不同键仍然分别排队;
- Resync 可能周期性重新入队。
因此,Reconcile 必须容忍重复执行。
五、WorkQueue:保存“需要处理什么”,不保存“完整对象”
1. Queue 中应该放什么
推荐放:
namespace/name
不推荐把事件对象直接放入队列:
queue.Add(obj) // 不推荐
原因是队列中的对象可能:
- 已经过期;
- 与 Cache 中的对象不一致;
- 被后续事件替代;
- 被误修改;
- 占用更多内存。
键的意义是“请重新检查这个对象”,而不是“请处理我入队时看到的那个快照”。
2. Queue 的基本生命周期
一个 Worker 的循环通常是:
Get
↓
处理 item
↓
成功:Forget
失败:AddRateLimited
↓
Done
Done 必须调用,否则 Queue 会认为该项目仍在处理中。
典型结构:
func (c *Controller) runWorker(ctx context.Context) {
for c.processNextItem(ctx) {
}
}
func (c *Controller) processNextItem(ctx context.Context) bool {
item, shutdown := c.queue.Get()
if shutdown {
return false
}
defer c.queue.Done(item)
key, ok := item.(string)
if !ok {
c.queue.Forget(item)
return true
}
if err := c.reconcile(ctx, key); err != nil {
c.queue.AddRateLimited(key)
return true
}
c.queue.Forget(item)
return true
}
Forget 的含义不是删除资源,而是清除该键在 RateLimiter 中的失败记录。成功后不调用 Forget,后续失败可能继续沿用旧的退避次数。
3. 限速重试
如果 Reconcile 失败,不能无延迟地死循环重试:
for {
err := reconcile()
if err != nil {
continue // 错误:可能打满 CPU 和 API Server
}
}
AddRateLimited 通常使用指数退避,具体实现和默认参数属于 client-go 实现细节,不应将某个默认延迟数字当成 Kubernetes API 保证。
退避的直觉是:
第 1 次失败:较快重试
第 2 次失败:等待更久
第 3 次失败:继续拉长间隔
这能降低以下故障的放大效应:
- API Server 暂时不可用;
- RBAC 配置错误;
- 资源冲突;
- 外部依赖不可用;
- Controller 自身逻辑错误。
永久错误不应无限重试。例如对象格式始终不合法时,应记录清晰错误,并根据业务需要:
- 继续限速重试,等待用户修复;
- 更新 Status 说明错误;
- 使用最大重试次数;
- 将错误分为可恢复和不可恢复两类。
六、Reconcile:从当前状态计算动作
1. Reconcile 的标准步骤
一个可靠的 Reconcile 通常包含:
1. 从 Queue 获取 namespace/name
2. 从 Cache 读取主资源
3. 如果资源不存在:
- 执行清理,或直接返回
4. 计算期望的子资源状态
5. 读取当前子资源
6. 不存在则创建
7. 存在但不符合期望则更新或 Patch
8. 已符合期望则不写入
9. 处理冲突和暂时性错误
形式化地说,Reconcile 可以视为:
其中:
- 是对象键;
- 是当前观察状态;
- 是创建、更新、删除、状态更新或重新入队等动作。
事件本身通常不是这个函数的核心输入,事件只决定何时调用它。
2. NotFound 的两种含义
从 Cache 读取主资源时得到 NotFound,可能表示:
- 资源真的已经删除;
- Informer 尚未观察到创建事件;
- Cache 暂时落后;
- Reconcile 处理的是一个过期队列键。
在 Cache 已完成初始同步后,NotFound 通常可以按“资源已不存在”处理。但是否要清理子资源,取决于生命周期设计:
- 使用 OwnerReference,让 Kubernetes Garbage Collector 清理;
- 使用 Finalizer,在主资源删除前执行外部清理;
- 子资源独立保留;
- 删除时只记录状态。
3. 更新必须基于最新 ResourceVersion
Kubernetes 对对象更新使用乐观并发控制。对象的 metadata.resourceVersion 表示某个观察版本。
典型冲突:
t0: Controller A 读取 rv=10
t1: Controller B 更新对象,rv=11
t2: Controller A 带 rv=10 更新
t3: API Server 返回 Conflict
Controller 不应盲目覆盖冲突:
if apierrors.IsConflict(err) {
// 重新读取最新对象,重新计算期望状态
}
在使用 Informer Cache 的情况下,重试时通常返回错误并让 Queue 重新调谐;下一次 Reconcile 会从 Cache 获取新对象。若冲突很紧急,也可以直接通过 Client 重新读取 API Server,但仍然需要重新计算,而不是复用旧对象。
七、幂等:重复调谐必须安全且可收敛
1. 幂等的定义
对同一状态重复应用同一调谐动作,如果第一次已经达到期望状态,后续执行不应继续产生新的业务变化:
这里的 表示一次完整的“读取当前状态并执行必要动作”。
需要区分两个概念:
- 函数幂等:重复调用返回相同结果;
- 操作幂等:重复执行不会造成额外副作用。
Controller 需要的是后者的业务版本。例如:
如果目标 ConfigMap 已经存在且内容正确,则不再 Update。
2. 一个完整算例
假设主资源 default/demo 的期望是拥有一个子 ConfigMap:
default/demo-observed
其数据为:
parentUID: <demo 的 UID>
第一次 Reconcile:
Cache:
demo 存在
demo-observed 不存在
动作:
Create demo-observed
结果:
API Server 创建成功
第二次 Reconcile 可能由以下任意原因触发:
- Create 产生事件;
- 初始同步产生重复入队;
- Controller 重试;
- 周期性 Resync;
- 其他无关事件导致重新入队。
再次读取:
Cache:
demo 存在
demo-observed 存在
data.parentUID 已正确
动作:
无写入
此时:
系统已经达到固定点。
3. 反例:每次 Reconcile 都 Update
错误逻辑:
child.Data["parentUID"] = string(parent.UID)
_, err := client.ConfigMaps(ns).Update(ctx, child, metav1.UpdateOptions{})
即使内容没有变化,也可能每次都发起 Update。后果包括:
Update ConfigMap
↓
ResourceVersion 变化
↓
Informer 收到 Update
↓
Queue 再次入队
↓
Reconcile 再次 Update
这可能形成高频自触发循环。即使事件去重和调度延迟阻止了严格的无限同步循环,也会产生:
- 不必要的 API Server 写压力;
- 更高的 Watch 流量;
- 更多冲突;
- 审计日志噪声;
- 其他 Watcher 被无意义唤醒。
正确方式是只在语义不一致时更新:
if child.Data["parentUID"] == desiredUID &&
child.Labels["managed-by"] == "demo-controller" {
return nil
}
4. 非幂等名称是常见错误
错误示例:
name := parent.Name + "-" + strconv.FormatInt(time.Now().UnixNano(), 10)
每次 Reconcile 都生成新名称:
demo-1710000000000
demo-1710000001000
demo-1710000002000
Controller 永远无法判断“已经存在的目标对象”,因为目标标识本身每次都变化。
常见解决方案:
- 使用确定性名称;
- 使用稳定的 OwnerReference;
- 使用标签和索引查找;
- 对外部系统使用稳定的业务键;
- 将生成的标识保存到 Status 或 Annotation。
5. 不要用时间戳伪造变化
以下代码也破坏幂等性:
desired.Annotations["last-reconciled-at"] = time.Now().Format(time.RFC3339)
如果该 Annotation 属于业务可见状态,它会让每次 Reconcile 都产生变化。时间信息应当只在确有必要时写入,例如:
- 记录最后一次成功动作;
- 记录外部操作时间;
- 只在状态发生语义变化时更新;
- 写入日志或指标,而不是资源对象。
八、一个基于 client-go 的完整示例
下面的示例实现一个简单 Controller:
- 监听带有
app.example.io/managed=true标签的 ConfigMap; - 为每个主 ConfigMap 确保存在一个确定性名称的子 ConfigMap;
- 子 ConfigMap 的数据记录主对象 UID;
- 设置 OwnerReference,使主 ConfigMap 删除后由 Garbage Collector 回收子对象;
- 只在子对象确实不符合期望时更新。
该示例使用 Kubernetes 内置 ConfigMap,避免引入 CRD 生成代码,便于说明 client-go 的核心流程。
1. go.mod
版本号应与目标集群和项目依赖策略匹配。以下只展示依赖形态,实际项目应使用项目锁定的兼容版本:
module example.com/demo-controller
go 1.22
require (
k8s.io/api v0.XX.X
k8s.io/apimachinery v0.XX.X
k8s.io/client-go v0.XX.X
)
v0.XX.X 不是可直接编译的版本号。项目中必须替换为真实、相互匹配的 Kubernetes 依赖版本;client-go、api 和 apimachinery 通常应保持同一发布线。
2. Controller 实现
package main
import (
"context"
"fmt"
"os"
"os/signal"
"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/labels"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go"
"k8s.io/client-go/informers"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/util/workqueue"
"k8s.io/klog/v2"
)
const (
managedLabel = "app.example.io/managed"
managedValue = "true"
childLabel = "app.example.io/child"
childValue = "true"
)
type Controller struct {
clientset clientset.Interface
informer cache.SharedIndexInformer
queue workqueue.RateLimitingInterface
}
func NewController(
clientset clientset.Interface,
informer cache.SharedIndexInformer,
) *Controller {
c := &Controller{
clientset: clientset,
informer: informer,
queue: workqueue.NewNamedRateLimitingQueue(
workqueue.DefaultControllerRateLimiter(),
"demo-configmaps",
),
}
informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: c.enqueue,
UpdateFunc: func(_, newObj interface{}) {
c.enqueue(newObj)
},
DeleteFunc: c.enqueue,
})
return c
}
func (c *Controller) enqueue(obj interface{}) {
key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
if err != nil {
utilruntime.HandleError(fmt.Errorf("build queue key: %w", err))
return
}
c.queue.Add(key)
}
func (c *Controller) Run(ctx context.Context, workers int) error {
defer runtime.HandleCrash()
defer c.queue.ShutDown()
klog.Info("waiting for informer cache sync")
if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) {
return fmt.Errorf("informer cache sync failed")
}
for i := 0; i < workers; i++ {
go wait.UntilWithContext(ctx, c.runWorker, time.Second)
}
<-ctx.Done()
return nil
}
func (c *Controller) runWorker(ctx context.Context) {
for c.processNextItem(ctx) {
}
}
func (c *Controller) processNextItem(ctx context.Context) bool {
item, shutdown := c.queue.Get()
if shutdown {
return false
}
defer c.queue.Done(item)
key, ok := item.(string)
if !ok {
c.queue.Forget(item)
utilruntime.HandleError(fmt.Errorf("unexpected queue item type %T", item))
return true
}
if err := c.reconcile(ctx, key); err != nil {
c.queue.AddRateLimited(key)
utilruntime.HandleError(fmt.Errorf("reconcile %q failed: %w", key, err))
return true
}
c.queue.Forget(item)
return true
}
func (c *Controller) reconcile(ctx context.Context, key string) error {
obj, exists, err := c.informer.GetIndexer().GetByKey(key)
if err != nil {
return err
}
if !exists {
// 主对象已经从 Cache 中消失。
// 子对象由 OwnerReference 负责回收,因此这里无需主动删除。
klog.Infof("ConfigMap %q no longer exists", key)
return nil
}
parent, ok := obj.(*corev1.ConfigMap)
if !ok {
return fmt.Errorf("expected ConfigMap, got %T", obj)
}
// 计算确定性子资源名称,不能使用时间戳或随机数。
childName := parent.Name + "-observed"
owner := true
blockOwnerDeletion := true
desired := &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Name: childName,
Namespace: parent.Namespace,
Labels: map[string]string{
childLabel: childValue,
},
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: corev1.SchemeGroupVersion.String(),
Kind: "ConfigMap",
Name: parent.Name,
UID: parent.UID,
Controller: &owner,
BlockOwnerDeletion: &blockOwnerDeletion,
},
},
},
Data: map[string]string{
"parentUID": string(parent.UID),
},
}
configmaps := c.clientset.CoreV1().
ConfigMaps(parent.Namespace)
current, err := configmaps.Get(ctx, childName, metav1.GetOptions{})
if apierrors.IsNotFound(err) {
_, err = configmaps.Create(ctx, desired, metav1.CreateOptions{})
if apierrors.IsAlreadyExists(err) {
// 另一个 Worker 或另一个 Controller 可能已经创建。
// 返回 nil 会丢失本次校验机会;更稳妥的做法是返回错误,
// 让 Queue 稍后重新读取并判断。
return fmt.Errorf("child %s appeared during create: %w", childName, err)
}
return err
}
if err != nil {
return err
}
if equalDesired(current, desired) {
return nil
}
// Update 必须基于 API Server 返回的 current,保留最新 ResourceVersion。
updated := current.DeepCopy()
updated.Labels = desired.Labels
updated.Data = desired.Data
updated.OwnerReferences = desired.OwnerReferences
_, err = configmaps.Update(ctx, updated, metav1.UpdateOptions{})
return err
}
func equalDesired(current, desired *corev1.ConfigMap) bool {
return current.Labels[childLabel] == childValue &&
current.Data["parentUID"] == desired.Data["parentUID"] &&
sameOwner(current.OwnerReferences, desired.OwnerReferences)
}
func sameOwner(a, b []metav1.OwnerReference) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i].UID != b[i].UID ||
a[i].Name != b[i].Name ||
a[i].Kind != b[i].Kind ||
a[i].APIVersion != b[i].APIVersion {
return false
}
}
return true
}
上面代码表达了几个关键因果关系:
- 事件处理器只入队,不直接执行创建和更新;
- Reconcile 从 Informer Cache 读取主对象;
- 写入子对象使用 Typed Client;
- 创建名称是确定性的;
- 更新使用 API Server 返回的对象,因而包含最新
ResourceVersion; - 已符合期望时不写入,从而满足幂等;
- 主对象不存在时正常返回,而不是把 NotFound 当成故障;
- OwnerReference 负责生命周期关联。
示例中的 clientset.Interface 和 corev1.SchemeGroupVersion 需要在完整工程中正确导入和初始化。不同 Kubernetes 依赖版本对包名和辅助函数可能存在小的 API 差异;实际编译时应以锁定版本的 go doc 为准。
3. 初始化 Client、Informer 和信号处理
func main() {
config, err := rest.InClusterConfig()
if err != nil {
klog.ErrorS(err, "build in-cluster config")
os.Exit(1)
}
clientset, err := clientset.NewForConfig(config)
if err != nil {
klog.ErrorS(err, "build clientset")
os.Exit(1)
}
factory := informers.NewSharedInformerFactoryWithOptions(
clientset,
10*time.Minute,
informers.WithTweakListOptions(func(options *metav1.ListOptions) {
options.LabelSelector = labels.Set{
managedLabel: managedValue,
}.AsSelector().String()
}),
)
informer := factory.Core().V1().ConfigMaps().Informer()
controller := NewController(clientset, informer)
ctx, cancel := signal.NotifyContext(
context.Background(),
syscall.SIGINT,
syscall.SIGTERM,
)
defer cancel()
factory.Start(ctx.Done())
if err := controller.Run(ctx, 2); err != nil {
klog.ErrorS(err, "controller stopped")
os.Exit(1)
}
}
为了完整编译,还需要导入:
import (
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
)
并将前文的 clientset.Interface、clientset.NewForConfig 替换为:
kubernetes.Interface
kubernetes.NewForConfig
也就是完整字段形式:
type Controller struct {
clientset kubernetes.Interface
informer cache.SharedIndexInformer
queue workqueue.RateLimitingInterface
}
这里有两个重要前置条件:
- 程序必须运行在 Kubernetes 集群内,或者改用
clientcmd.BuildConfigFromFlags加载本地 kubeconfig; - ServiceAccount 必须拥有读取和写入 ConfigMap 的 RBAC 权限。
4. 示例资源和预期结果
创建一个主 ConfigMap:
apiVersion: v1
kind: ConfigMap
metadata:
name: demo
namespace: default
labels:
app.example.io/managed: "true"
data:
input: hello
应用:
kubectl apply -f demo.yaml
预期子资源:
kubectl get configmaps -n default
可能看到:
NAME DATA AGE
demo 1 ...
demo-observed 1 ...
检查内容:
kubectl get configmap demo-observed \
-n default \
-o jsonpath='{.data.parentUID}{"\n"}'
输出应为主 ConfigMap 的 UID:
7f1b...
再次执行:
kubectl apply -f demo.yaml
Controller 可能收到新的事件,但 demo-observed 已经符合期望,因此不会产生无意义的 Update。
删除主对象:
kubectl delete configmap demo -n default
由于子 ConfigMap 设置了主对象的 OwnerReference,Kubernetes Garbage Collector 通常会回收 demo-observed。这依赖:
- OwnerReference 指向真实存在的对象;
- UID 正确;
- API Group、Kind、Name 正确;
- Kubernetes 的垃圾回收机制正常工作;
- RBAC 和对象作用域关系合法。
OwnerReference 不是所有资源组合都可以任意设置。跨命名空间 OwnerReference、无效 UID 或不匹配的 GVK 可能不会产生预期的垃圾回收效果。
九、并发模型:一个键通常串行,不同键可以并行
1. 多 Worker 的语义
如果启动两个 Worker:
controller.Run(ctx, 2)
通常意味着:
default/a ──> Worker 1
default/b ──> Worker 2
default/c ──> 等待
同一时刻,不同键可以并行处理,从而提高吞吐;但 Controller 不能依赖全局串行执行。
一个健壮的 Reconcile 必须假设:
- 不同对象会并发处理;
- 同一对象可能因事件和重试再次入队;
- 多个 Controller 可能修改同一资源;
- Worker 可能在写入后立即再次收到事件;
- Controller 可能在任意中间步骤被终止。
2. 不应把内存字段当作真实状态
错误设计:
c.lastProcessed[name] = desiredVersion
然后根据这个内存 Map 判断是否需要更新。
Controller 重启后,这个 Map 会丢失;多个副本之间也不会共享;故障恢复时会产生错误判断。真实状态应存放在:
- Kubernetes 对象的 Spec;
- Kubernetes 对象的 Status;
- Annotation 或 Label;
- 外部系统的持久化状态。
内存缓存适合优化,不适合作为一致性依据。
3. 防止同一资源的业务竞争
Queue 通常保证同一键不会被多个 Worker 同时处理,但这不是跨进程全局锁。以下场景仍然可能并发:
Controller 副本 A ──> 更新对象
Controller 副本 B ──> 更新同一对象
另一个 Kubernetes Controller ──> 更新同一对象
因此必须依赖 API Server 的资源版本冲突检测,而不是依赖本地锁。
十、更新、Patch 与字段所有权
1. Update 的风险
Update 会提交整个对象。若 Controller 使用从旧 Cache 读取的对象执行 Update,可能覆盖其他客户端已经修改的字段。
错误模式:
读取旧对象
修改自己的字段
Update 整个对象
覆盖其他客户端的新字段
更安全的方式包括:
- 基于 API Server 最新对象更新;
- 只管理自己拥有的字段;
- 使用 Patch;
- 使用 Server-Side Apply 并声明 FieldManager。
2. Patch 的语义
Patch 可以只修改部分字段,降低覆盖无关字段的风险。例如 JSON Merge Patch:
patch := []byte(`{
"metadata": {
"labels": {
"app.example.io/managed": "true"
}
}
}`)
_, err := client.ConfigMaps(ns).Patch(
ctx,
name,
types.MergePatchType,
patch,
metav1.PatchOptions{},
)
但 Patch 仍然必须幂等。重复执行同一个 Patch,如果目标字段已经是目标值,不应造成业务上的无限变化。
3. Server-Side Apply
Server-Side Apply 使用声明式字段管理。它适合多个参与者分别拥有对象的不同字段,但需要明确:
- FieldManager 名称;
- 哪些字段由 Controller 负责;
- 冲突时是报错还是强制接管;
- 是否会改变现有字段所有权。
Apply 不是“自动解决所有冲突”。如果两个 FieldManager 都声明拥有同一字段,仍然可能发生冲突。
十一、Status、Spec 和 Metadata 的调谐边界
对于 CRD 或其他声明式资源,通常应区分:
spec = 用户希望系统做什么
status = Controller 观察到系统做到了什么
metadata = 标识、标签、注解、版本和生命周期信息
Controller 不应把用户的 Spec 当作自己的工作区,也不应随意覆盖用户控制的字段。
例如:
spec:
replicas: 3
status:
observedReplicas: 3
conditions:
- type: Ready
status: "True"
如果 Controller 负责更新 Status,应注意:
- Status 更新可能再次触发 Update 事件;
- 必须只在 Status 语义发生变化时更新;
- Conditions 的
LastTransitionTime不应在每次 Reconcile 都刷新; - Status 更新与 Spec 更新可能产生资源版本冲突;
- 使用 Status 子资源时,应调用对应的 Status 更新接口。
一个错误的 Condition 实现:
condition.LastTransitionTime = metav1.Now()
status.Conditions = []metav1.Condition{condition}
UpdateStatus(...)
每次调谐都会改变时间戳,导致无意义写入和自触发。正确做法是只有当 Condition 的状态、原因、消息等语义字段发生变化时,才更新时间。
十二、最终一致性与重新入队
1. 为什么一次 Reconcile 不一定完成任务
Controller 可能依赖多个异步系统:
创建 Deployment
↓
Deployment Controller 创建 ReplicaSet
↓
ReplicaSet Controller 创建 Pod
↓
Scheduler 调度 Pod
↓
Kubelet 启动容器
如果 Controller 的目标是“应用已经可用”,它不能在创建 Deployment 后立即认为目标完成,而应观察后续状态,并在必要时重新入队。
设目标状态为:
Deployment.spec.replicas = 3
Deployment.status.availableReplicas = 3
中间状态可能是:
第 1 次:desired=3,available=0
第 2 次:desired=3,available=1
第 3 次:desired=3,available=3
Reconcile 的结果可能不是错误,而是“当前未完成,稍后再次检查”。
2. 延迟重新入队
某些 client-go Controller 会使用:
queue.AddAfter(key, 30*time.Second)
这适合:
- 外部系统没有 Watch;
- 等待异步操作完成;
- 定期检查 TTL 或过期状态;
- 防止长期没有事件时状态停滞。
但固定周期轮询会增加 API 和 CPU 负担。若 Kubernetes 资源本身有 Watch 事件,应优先通过相关资源事件重新触发,而不是无条件轮询。
3. Requeue 与错误重试的区别
两者语义不同:
- 业务未完成但没有错误:可以
AddAfter; - 本次读取或写入失败:通常返回错误并
AddRateLimited; - 永久性业务错误:更新 Status 或记录明确错误,避免无意义高速重试。
应避免把所有情况都当成错误重试,否则监控中无法区分“系统暂时故障”和“业务尚未达到期望”。
十三、Cache、Queue 和 Reconcile 的故障路径
1. Watch 断开
Watch 可能因为以下原因断开:
- 网络连接中断;
- API Server 重启;
- ResourceVersion 太旧;
- 代理或负载均衡关闭连接;
- 服务端超时。
Informer 会重新建立同步过程。Controller 不应自行假设 Watch 永远存在,也不应只依赖第一次收到的事件。
2. ResourceVersion 过旧
如果 Watch 使用的 ResourceVersion 已经无法从 API Server 的历史事件存储中恢复,服务端可能返回类似“Gone”的错误。Informer 通常通过重新 List 获取当前状态,再建立新的 Watch。
这也是为什么 Reconcile 必须是“读取当前状态并计算”,而不是依赖每一条历史事件都被精确处理。
3. Controller 重启
重启会丢失:
- 内存 Queue;
- 未持久化的重试次数;
- 内存索引;
- 未写入 Kubernetes 的中间状态。
但重启后 Informer 会重新 List,Controller 可以再次把当前对象加入处理流程。只要 Reconcile 幂等,重启通常只是延迟,而不是破坏状态。
4. 创建成功但响应丢失
这是典型故障:
Controller 发送 Create
API Server 已创建对象
网络在响应返回前断开
Controller 认为 Create 失败
Queue 重试
再次 Create
第二次 Create 返回 AlreadyExists。这个结果不能简单当成系统故障。Controller 应重新读取对象并检查其是否符合期望:
Create 返回 AlreadyExists
↓
Get 当前对象
↓
符合期望:成功
不符合期望:Update/Patch
这正是幂等设计要处理的故障路径。
十四、常见误区与诊断方法
误区一:把事件对象当成最新对象
表现:
- 使用旧对象覆盖新字段;
- 处理删除事件时类型断言失败;
- 事件对象已经不存在却继续执行更新。
诊断:
- 打印 Queue key,而不是只打印事件对象;
- 在 Reconcile 中记录 Cache 读取结果;
- 检查
ResourceVersion; - 对 Delete 使用
DeletionHandlingMetaNamespaceKeyFunc。
误区二:直接修改 Cache 对象
表现:
- 逻辑偶尔看到 API Server 中不存在的字段;
- 数据竞争;
- DeepEqual 判断异常;
- 重启后状态恢复。
诊断:
- 检查是否对 Informer 返回对象调用了字段赋值;
- 使用
DeepCopy(); - 只把 Cache 对象当作只读输入。
误区三:每次调谐都写入资源
表现:
- Controller 日志持续出现 Update;
kubectl get -o yaml中 ResourceVersion 持续变化;- API Server 审计日志出现大量重复写请求;
- Controller 自己不断触发 Update 事件。
诊断:
- 记录“期望对象”和“当前对象”的语义差异;
- 分别统计 Get、Create、Update、Patch 次数;
- 检查时间戳、随机值、无条件排序和 Map 重建;
- 检查 Status Condition 是否每次都刷新时间。
误区四:把所有 NotFound 都当错误
表现:
- 删除对象后 Queue 不断重试;
- 日志充斥 NotFound;
- Controller 无法稳定清理。
诊断:
区分:
主资源 NotFound:
可能表示正常删除
子资源 NotFound:
可能表示需要创建
依赖资源 NotFound:
可能表示等待依赖出现
不同资源的 NotFound 具有不同业务含义。
误区五:忽略 RBAC 与作用域
表现:
forbidden
cannot list resource
cannot update resource
诊断:
检查 ServiceAccount、Role/ClusterRole 和 RoleBinding/ClusterRoleBinding:
kubectl auth can-i \
--as=system:serviceaccount:default:demo-controller \
get configmaps -n default
kubectl auth can-i \
--as=system:serviceaccount:default:demo-controller \
create configmaps -n default
kubectl auth can-i \
--as=system:serviceaccount:default:demo-controller \
update configmaps -n default
Informer 的 List 和 Watch 也需要权限。只允许 Get 而没有 List/Watch,Informer 无法正常工作。
十五、规范保证、实现细节和工程建议的边界
需要明确区分三类结论。
Kubernetes API 的规范性行为
通常可以依赖:
- API 对象具有
metadata.name、namespace、UID、resourceVersion等元数据; - 更新会进行资源版本并发检查;
- Watch 是对资源变化的观察机制;
- OwnerReference 参与垃圾回收;
- Status、Finalizer、Label、Annotation 有明确 API 语义;
- ResourceVersion 过旧时可能需要重新 List。
client-go 的常见实现
通常可以观察到:
- SharedInformer 使用本地 Store/Indexer;
- 同一个 SharedInformer 的多个 Handler 共享缓存;
- WorkQueue 对相同键具有去重和并发安全特征;
- RateLimitingQueue 维护失败次数并执行退避;
- Informer 在 Watch 断开后重新同步。
但具体内部类型、默认退避参数、线程调度和实现字段不应作为跨版本稳定契约。
Controller 的工程建议
属于应用设计选择:
- Worker 数量;
- 是否使用 Patch 或 Apply;
- 是否使用最终一致性轮询;
- 是否设置最大重试次数;
- 是否使用 Finalizer;
- 是否更新 Status;
- 如何划分字段所有权;
- 如何处理永久性错误。
这些选择必须结合资源规模、API Server 负载、业务副作用和恢复要求验证,不能仅因为某个示例采用了某种方案就认为它是 Kubernetes 的强制要求。
十六、一个可验证的 Controller 正确性检查
可以用以下状态序列验证 Controller 是否具备基本幂等性:
初始:
主对象存在
子对象不存在
R1:
Create 子对象成功
R2:
重复 Reconcile
子对象内容正确
不发送 Update
外部修改:
子对象被改成错误内容
R3:
读取错误内容
Update 为期望内容
网络故障:
R3 的响应丢失
Queue 重试
R4:
读取子对象已经正确
不重复创建,不重复更新
主对象删除:
Cache 返回 NotFound
Controller 正常退出该键的处理
OwnerReference 最终回收子对象
如果一个 Controller 在这组序列中能够稳定工作,它通常已经具备了几个重要性质:
结语
client-go Controller 的核心不是某个回调函数,而是一条完整的状态处理链:
Informer 从 API Server 观察变化
↓
Cache 提供本地只读快照
↓
Event Handler 把对象转换成稳定键
↓
Queue 合并、调度和重试键
↓
Reconcile 重新读取当前状态并计算动作
↓
Typed Client 修改 API Server
↓
新的状态再次通过 Watch 进入系统
Queue 中保存键而不是旧对象,使 Reconcile 能够重新观察;Cache 提供高效读取,但接受短暂滞后;API Server 提供资源版本控制,防止旧对象无条件覆盖新对象;幂等逻辑则让重复事件、重试、重启和响应丢失都能最终回到同一个稳定状态。
可以把 Controller 的基本正确性概括为:
只有将这些机制组合起来,Controller 才不是“收到事件后执行一次脚本”,而是一个能够在分布式故障和状态变化中持续恢复、反复校正并最终一致的 Kubernetes 控制回路。
系列导航与关联阅读
- 系列入口:Kubernetes 完整学习路线:从 Pod 与控制面到安全、运维和 Operator
- 上一篇:Kubernetes CRD:Schema、版本、Defaulting、Validation、Conversion 和存储
- 下一篇:Kubernetes Operator 模式:领域状态机、升级、备份、恢复和测试
- 延伸:Kubernetes 调谐循环:Desired State、Watch、Queue、幂等和最终一致
官方资料
本文依据 Kubernetes、CNCF 与相关项目官方文档重新梳理;正文和生产清单由 WR BLOG 编写。

评论
0 条讨论