

















读完本篇,你应该能回答:sample-controller 的 Controller 结构体有哪些字段,各自承担什么职责?NewController 中如何注册 Informer 和事件处理器?Reconciler 的 syncHandler 完整流程是什么?主调谐循环 Run 如何启动 Workers 并等待缓存同步?Deployment 调谐和 Service 调谐有什么不同?RBAC 如何配置才能让控制器访问所需资源?如何用 FilteredResourceEventHandler 过滤无关事件?如何通过 OwnerReference 实现资源级联删除?
Kubernetes Go Operator Controller sample-controller client-go Informer k8s v1.36.1
学习重点提示 — 建议先通读全文,再重点回顾标注内容
重点掌握(必须)
- Controller 结构体设计:理解每个字段(clientset、Lister、Synced、workqueue、recorder)的职责和来源
- Reconciler 完整流程:syncHandler 中 Get → Create → Check OwnerReference → Update → UpdateStatus 五步法
- 事件处理链路:enqueueFoo / handleObject 的注册时机和触发条件
- Worker 主循环:processNextWorkItem 中 Get → syncHandler → Done/Forget/AddRateLimited 三分支
次重点(了解即可)
- EventRecorder 的事件发布机制和事件限频
- FilteredResourceEventHandler 的过滤函数设计
- RBAC ClusterRole/Role 的最小权限原则
文章目录
思考记忆提示 — 本节是全篇的"入口"——理解为什么需要自定义控制器,才能理解 sample-controller 解决什么问题
Kubernetes 的核心哲学是声明式 API + 控制器模式。你声明期望状态(Spec),Kubernetes 通过控制器持续调谐,使实际状态向期望状态收敛。
但 Kubernetes 原生只提供了 Deployment、StatefulSet、DaemonSet、Job 等几种控制器的实现。如果你需要管理一个自定义资源(Custom Resource),比如 "Foo"(一个需要关联 Deployment 的抽象概念),Kubernetes 本身并不知道如何处理它——你需要自己写一个控制器。
sample-controller(位于 staging/src/k8s.io/sample-controller/)是 Kubernetes 官方提供的最小可用控制器模板。它演示了一个完整的控制器应该包含哪些组件:自定义资源定义(Foo CRD)、Informers、Workqueue、Reconciler、RBAC 配置。所有 Operator(包括 Kube-builder、Operator-SDK 生成的代码)都脱胎于这个模式。
我的理解的意思是说
可以把 sample-controller 想象成一个自动备课机器人:
sample-controller 的核心逻辑就是:Foo 资源声明了"我要管理一个 Deployment",Controller 发现没有就创建,发现规格变了就更新。这是一个最简单的"声明一个,管理一个"的模式。
思考记忆提示 — 本节是全篇的"地图"——在深入每个细节之前,先对整个架构有一个全局视图
sample-controller 的核心文件结构如下:
staging/src/k8s.io/sample-controller/
main.go # 程序入口:创建 clientset、InformerFactory、启动 Controller
controller.go # 核心逻辑:Controller 结构体 + 所有方法实现
controller_test.go # 单元测试:fixture 模式 + 行为验证
pkg/apis/samplecontroller/v1alpha1/
types.go # Foo 自定义资源定义(Spec + Status)
zz_generated.deepcopy.go # DeepCopy 代码(自动生成)
pkg/generated/
clientset/ # Foo 资源的 typed client(自动生成)
informers/ # Foo 资源的 SharedInformer(自动生成)
listers/ # Foo 资源的缓存读取接口(自动生成)
applyconfiguration/ # Foo 资源的 apply 模式配置(自动生成)
artifacts/rbac/
leader-election-role.yaml # Leader 选举 RBAC 配置
namespace.yaml # 命名空间
role.yaml # Controller 的 RBAC 权限
role-binding.yaml # ServiceAccount 与 Role 的绑定
整体架构可以用一张数据流图来描述:
┌─────────────────────────────────────────────────────────────────────────────────────┐
│ sample-controller 完整数据流 │
│ │
│ ┌──────────────────┐ main.go │
│ │ kubeconfig/ │ 创建 Config │
│ │ API Server │◄──────► kubeClient ────► kubeInformerFactory (Deployment) │
│ └──────────────────┘ (kubernetes.Interface) │ │
│ ▲ │ Start() │
│ │ ▼ │
│ exampleClient ┌─────────────────┐ │
│ (Foo CRD Client) │ Deployment │ │
│ ▲ │ Informer │ │
│ │ │ + Indexer │ │
│ ┌───────────────────────┴──────────────┐ └────────┬────────┘ │
│ │ exampleInformerFactory │ │ │
│ │ (Foo Informer + Indexer + Lister) │ │ AddEventHandler │
│ └────────────────┬─────────────────────┘ ▼ │
│ │ ┌─────────────────┐ │
│ │ Start() │ Foo Informer │ │
│ │ │ + Indexer │ │
│ │ └────────┬────────┘ │
│ │ │ │
│ │ │ enqueueFoo / │
│ │ │ handleObject │
│ │ ▼ │
│ │ ┌─────────────────┐ │
│ │ │ Workqueue │ │
│ │ │ (rate-limited) │ │
│ │ └────────┬────────┘ │
│ │ │ │
│ │ │ Get() │
│ │ ▼ │
│ │ ┌─────────────────┐ │
│ │ │ syncHandler │ │
│ │ │ (Reconciler) │ │
│ │ └────────┬────────┘ │
│ │ │ │
│ │ ┌────────────────────┼────────────────────┐│
│ │ │ │ ││
│ │ ▼ ▼ ▼│
│ │ Create Deployment Update Deployment Update Status│
│ │ │ │ ││
│ │ └────────────────────┼────────────────────┘│
│ │ kubeClient │
│ └──────────────────────────────────────► AppsV1() │
└─────────────────────────────────────────────────────────────────┘
设计精髓
sample-controller 的最大特点是同时监听两个资源类型(Foo 和 Deployment),并通过 OwnerReference 建立它们之间的拥有关系。这种设计实现了两个关键能力:
思考记忆提示 — Controller 结构体是全篇的核心骨架,每个字段都有明确的职责
Controller 结构体定义在 staging/src/k8s.io/sample-controller/controller.go:69~89:
// staging/src/k8s.io/sample-controller/controller.go(行 69~89)
// Controller is the controller implementation for Foo resources
type Controller struct {
// kubeclientset is a standard kubernetes clientset
// 用来操作 Kubernetes 内置资源(Deployment、Pod、Service 等)
kubeclientset kubernetes.Interface
// sampleclientset is a clientset for our own API group
// 用来操作 Foo CRD(创建/更新/删除 Foo、更新 FooStatus)
sampleclientset clientset.Interface
// deploymentsLister + deploymentsSynced
// 监听 Deployment 变化,从本地缓存(Indexer)快速读取
// 对应的 Informer 在 main.go 中由 kubeInformerFactory.Apps().V1().Deployments() 创建
deploymentsLister appslisters.DeploymentLister
deploymentsSynced cache.InformerSynced // 返回函数:() bool,判断缓存是否已同步
// foosLister + foosSynced
// 监听 Foo 自定义资源变化,从本地缓存读取
// 对应的 Informer 由 exampleInformerFactory.Samplecontroller().V1alpha1().Foos() 创建
foosLister listers.FooLister
foosSynced cache.InformerSynced
// workqueue is a rate limited work queue.
// 这是整个控制器的"任务调度中心":去重、延时、限速、重试
// 泛型类型为 cache.ObjectName(代表 namespace/name)
workqueue workqueue.TypedRateLimitingInterface[cache.ObjectName]
// recorder is an event recorder for recording Event resources to the
// Kubernetes API. 用来向 Foo 资源写入事件(Normal/Warning),便于 kubectl describe 查
recorder record.EventRecorder
}
小贴士
cache.InformerSynced 的类型是 func() bool,不是 bool。这是因为 Informer 的缓存同步是一个过程,不是瞬间完成的——HasSynced 是一个查询函数,每次调用都返回当前的同步状态。在 cache.WaitForCacheSync(ctx.Done(), c.deploymentsSynced, c.foosSynced) 中,这个函数会被反复调用,直到返回 true。
思考记忆提示 — NewController 是整个控制器的"组装工厂"——理解它就知道所有组件是如何拼接在一起的
// staging/src/k8s.io/sample-controller/controller.go(行 92~156)
// NewController returns a new sample controller
func NewController(
ctx context.Context,
kubeclientset kubernetes.Interface, // 操作内置资源
sampleclientset clientset.Interface, // 操作 Foo CRD
deploymentInformer appsinformers.DeploymentInformer, // Deployment 监听器
fooInformer informers.FooInformer, // Foo 监听器
) *Controller { ... }
NewController 遵循依赖注入(Dependency Injection)模式:它不自己创建 Kubernetes 客户端,而是接收外部传入的 clientset 和 Informer 实例。这种设计使得控制器可以轻松通过 fake clientset 和 fake informer 进行单元测试。
// NewController 内部(controller.go:100~109)
// 1. 将 sample-controller 的类型注册到 Scheme
// 这样 Event 对象才能正确序列化 Foo 相关的 OwnerReference
utilruntime.Must(samplescheme.AddToScheme(scheme.Scheme))
// 2. 创建事件广播器(Broadcaster)
eventBroadcaster := record.NewBroadcaster(record.WithContext(ctx))
// 3. 启动结构化日志输出(将事件写到 klog,不是 API Server)
eventBroadcaster.StartStructuredLogging(0)
// 4. 启动将事件写入 API Server 的 Sink(每条事件都创建一个 Event 资源)
eventBroadcaster.StartRecordingToSink(&typedcorev1.EventSinkImpl{
Interface: kubeclientset.CoreV1().Events(""),
})
// 5. 创建 Recorder:写入事件时使用 controllerAgentName 作为 Component
recorder := eventBroadcaster.NewRecorder(scheme.Scheme,
corev1.EventSource{Component: controllerAgentName})
注意
StartRecordingToSink 会为每条事件调用 API Server 创建 Event 资源。在生产环境中,如果控制器的调谐频率很高(比如每秒调谐数百次),会产生大量 Event 对象,占用 etcd 空间。建议在高吞吐量场景下使用 record.EventAggregator 或直接关闭事件记录(不调用 StartRecordingToSink)。
// NewController 内部(controller.go:110~113)
// 创建限速器:两个限速策略取最大值
ratelimiter := workqueue.NewTypedMaxOfRateLimiter(
// 策略1:指数退避(per-item)
// 第1次失败:5ms 第2次:10ms 第3次:20ms ... 第18次:~655s(接近上限)
workqueue.NewTypedItemExponentialFailureRateLimiter[cache.ObjectName](
5*time.Millisecond, // baseDelay
1000*time.Second, // maxDelay
),
// 策略2:令牌桶(全局)
// 50 QPS,允许突发最多 300 个请求
&workqueue.TypedBucketRateLimiter[cache.ObjectName]{
Limiter: rate.NewLimiter(rate.Limit(50), 300),
},
)
NewController 中最核心的部分是两个事件处理器的注册:
// NewController 内部(controller.go:126~153)
logger.Info("Setting up event handlers")
// ========== 注册 Foo 的事件处理器(正向调谐)==========
fooInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: controller.enqueueFoo, // Foo 创建 → 入队
UpdateFunc: func(old, new interface{}) {
controller.enqueueFoo(new) // Foo 更新 → 入队(直接用新的)
},
// DeleteFunc 未注册:因为 Foo 被删除时,其拥有的 Deployment
// 会被 Kubernetes GC 级联删除,Deployment 的 DeleteFunc 会触发 handleObject
})
// ========== 注册 Deployment 的事件处理器(级联调谐)==========
deploymentInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: controller.handleObject, // Deployment 被外部创建 → 触发调谐
UpdateFunc: func(old, new interface{}) {
newDepl := new.(*appsv1.Deployment)
oldDepl := old.(*appsv1.Deployment)
// 如果只是 ResourceVersion 变化(Periodic Resync),忽略
if newDepl.ResourceVersion == oldDepl.ResourceVersion {
return
}
controller.handleObject(new)
},
DeleteFunc: controller.handleObject, // Deployment 被删除 → 触发调谐
})
我的理解的意思是说
为什么 Foo 的 DeleteFunc 没有注册?因为当 Foo 被删除时,Kubernetes 的垃圾收集器(Garbage Collector)会根据 OwnerReference 自动删除其拥有的 Deployment。Deployment 被删除后,Deployment Informer 的 DeleteFunc 会被触发,从而通过 handleObject 重新触发 Foo 的调谐。此时 Foo 已经不存在了,foosLister.Get() 会返回 NotFound,Reconciler 就会安静地返回 nil——什么都不做。
这种设计叫GC 驱动的级联删除,是 Kubernetes 推荐的所有权管理模式。
思考记忆提示 — Run 方法是控制器的"启动入口"——理解它就知道控制器启动时的完整顺序
// staging/src/k8s.io/sample-controller/controller.go(行 162~188)
// Run will set up the event handlers for types we are interested in, as well
// as syncing informer caches and starting workers.
func (c *Controller) Run(ctx context.Context, workers int) error {
defer utilruntime.HandleCrash() // 崩溃时记录日志,不让进程退出
defer c.workqueue.ShutDown() // 退出前关闭队列
logger := klog.FromContext(ctx)
logger.Info("Starting Foo controller")
// ========== 阶段1:等待所有 Informer 缓存同步完成 ==========
// 这是关键!如果缓存还没同步完,Reconciler 从 Lister 拿到的可能是空数据
logger.Info("Waiting for informer caches to sync")
if ok := cache.WaitForCacheSync(ctx.Done(), c.deploymentsSynced, c.foosSynced); !ok {
return fmt.Errorf("failed to wait for caches to sync")
}
// ========== 阶段2:启动 Workers ==========
logger.Info("Starting workers", "count", workers)
for i := 0; i < workers; i++ {
// 每个 Worker 是一个独立 goroutine,持续从 Workqueue 取任务
// wait.UntilWithContext:每 1 秒执行一次 runWorker(带上下文取消)
go wait.UntilWithContext(ctx, c.runWorker, time.Second)
}
logger.Info("Started workers")
<-ctx.Done() // 阻塞,直到收到 ctx 的取消信号(SIGTERM/SIGINT)
logger.Info("Shutting down workers")
return nil
}
小贴士
WaitForCacheSync 内部会阻塞,直到所有传入的 Synced 函数返回 true。如果 Informer 启动失败(比如 API Server 不可达),这个函数会一直阻塞——控制器永远无法启动。在生产环境中,需要配合 ctx 的超时机制,避免无限等待。
思考记忆提示 — Worker 是整个控制器的"发动机"——processNextWorkItem 是每个 Worker 反复执行的核心循环
// staging/src/k8s.io/sample-controller/controller.go(行 190~236)
// runWorker 驱动 processNextWorkItem 循环
func (c *Controller) runWorker(ctx context.Context) {
for c.processNextWorkItem(ctx) {
// 无限循环:每次处理完一个任务后,继续处理下一个
// 当 processNextWorkItem 返回 false 时(队列关闭),goroutine 退出
}
}
// processNextWorkItem:从队列取一个任务,尝试处理
func (c *Controller) processNextWorkItem(ctx context.Context) bool {
// Get() 阻塞直到拿到任务或队列关闭
objRef, shutdown := c.workqueue.Get()
logger := klog.FromContext(ctx)
if shutdown {
return false // 队列已关闭,Worker 退出
}
// defer 确保无论成功失败,Done() 都会被调用
// Done() 的作用是从 processing 集合中移除 item,允许同一个 item 被重新入队
defer c.workqueue.Done(objRef)
// 调用 syncHandler 执行实际的调谐逻辑
err := c.syncHandler(ctx, objRef)
if err == nil {
// 成功:Forget 从队列中彻底清除追踪,AddRateLimited 次数清零
c.workqueue.Forget(objRef)
logger.Info("Successfully synced", "objectName", objRef)
return true
}
// 失败:记录错误,并将任务重新入队(带限速)
utilruntime.HandleErrorWithContext(ctx, err,
"Error syncing; requeuing for later retry", "objectReference", objRef)
c.workqueue.AddRateLimited(objRef) // 根据 RateLimiter 计算等待时间后重新入队
return true
}
我的理解的意思是说
Workqueue 的去重机制可以通过三个集合来理解:
当一个 key 被 Get 时,它从 queue 移到 processing;当 Done 被调用时,它从 processing 移除。如果在 processing 期间同一个 key 被再次 Add,它会被加到 dirty。Done 执行后,如果 dirty 中有这个 key,就重新入队。这意味着:同一个 key 在 processing 期间无论被入队多少次,都只会被处理一次。
思考记忆提示 — syncHandler 是全篇最重要的方法——控制器的所有业务逻辑都在这里
syncHandler 是控制器的"调谐引擎",实现了完整的五步调谐逻辑(controller.go:241~312):
// staging/src/k8s.io/sample-controller/controller.go(行 241~312)
// syncHandler compares the actual state with the desired, and attempts to
// converge the two. It then updates the Status block of the Foo resource.
func (c *Controller) syncHandler(ctx context.Context, objectRef cache.ObjectName) error {
logger := klog.LoggerWithValues(klog.FromContext(ctx), "objectRef", objectRef)
// ========== 步骤1:从 Foo Informer 缓存获取 Foo 资源 ==========
// objectRef 是 namespace/name,通过 foosLister.Foos(ns).Get(name) 查询本地缓存
foo, err := c.foosLister.Foos(objectRef.Namespace).Get(objectRef.Name)
if err != nil {
// Foo 可能已被删除(GC 触发 handleObject 后,Foo 已不存在)
if errors.IsNotFound(err) {
utilruntime.HandleErrorWithContext(ctx, err,
"Foo referenced by item in work queue no longer exists",
"objectReference", objectRef)
return nil // 安静地返回 nil:任务完成,无需重试
}
return err // 其他错误(网络问题等):重新入队,稍后重试
}
// ========== 步骤2:获取期望的 Deployment 名称 ==========
deploymentName := foo.Spec.DeploymentName
if deploymentName == "" {
// Spec 中没有声明 DeploymentName:无法调谐,记录错误
utilruntime.HandleErrorWithContext(ctx, nil,
"Deployment name missing from object reference",
"objectReference", objectRef)
return nil // 返回 nil 而不是 error:这不是临时问题,重试也无用
}
// ========== 步骤3:从 Deployment Informer 缓存获取 Deployment ==========
// 注意:先查缓存,缓存没有再通过 clientset.Create
deployment, err := c.deploymentsLister.Deployments(foo.Namespace).Get(deploymentName)
// ========== 步骤3(扩展):Deployment 不存在则创建 ==========
if errors.IsNotFound(err) {
// 调用 kubeclientset 创建真实的 Deployment 资源
deployment, err = c.kubeclientset.AppsV1().Deployments(foo.Namespace).Create(
ctx, newDeployment(foo), // newDeployment() 负责构造 Deployment 对象
metav1.CreateOptions{FieldManager: FieldManager},
)
// 创建后立即返回,等下一次调谐(Deployment 的 AddEventHandler 会再次触发)
}
// 如果 Get/Create 出错(网络错误等),重新入队重试
if err != nil {
return err
}
// ========== 步骤4:检查 OwnerReference(防止误管他人 Deployment)==========
// 如果该 Deployment 不是由这个 Foo 管理的,记录警告事件并报错
// 这是防止"名称冲突但非同一人管理"的边界情况
if !metav1.IsControlledBy(deployment, foo) {
msg := fmt.Sprintf(MessageResourceExists, deployment.Name)
c.recorder.Event(foo, corev1.EventTypeWarning, ErrResourceExists, msg)
return fmt.Errorf("%s", msg) // 返回 error:这是配置错误,重试无意义(会触发指数退避)
}
// ========== 步骤5:检查副本数是否匹配 ==========
// 如果 Foo.Spec.Replicas 与 Deployment.Spec.Replicas 不一致,更新 Deployment
if foo.Spec.Replicas != nil && *foo.Spec.Replicas != *deployment.Spec.Replicas {
logger.V(4).Info("Update deployment resource",
"currentReplicas", *deployment.Spec.Replicas,
"desiredReplicas", *foo.Spec.Replicas)
deployment, err = c.kubeclientset.AppsV1().Deployments(foo.Namespace).Update(
ctx, newDeployment(foo), // 用最新的 Foo Spec 重新构造 Deployment
metav1.UpdateOptions{FieldManager: FieldManager},
)
}
if err != nil {
return err // 更新失败,重新入队重试
}
// ========== 步骤6:更新 Foo 的 Status(反映 Deployment 的实际状态)==========
err = c.updateFooStatus(ctx, foo, deployment)
if err != nil {
return err
}
// 全部成功:记录成功事件,返回 nil(触发 Forget)
c.recorder.Event(foo, corev1.EventTypeNormal, SuccessSynced, MessageResourceSynced)
return nil
}
设计精髓
syncHandler 遵循一个重要的原则:幂等性(Idempotency)。无论是 Create 还是 Update,操作都应该可以安全地重复执行而不会产生副作用。每次调谐都从最新的 Foo Spec 出发,计算期望的 Deployment 状态,然后与实际状态对比。如果实际状态与期望状态不一致,就调用 clientset Update。这个模式叫"声明式调谐"(Declarative Reconciliation)。
// staging/src/k8s.io/sample-controller/controller.go(行 314~326)
func (c *Controller) updateFooStatus(
ctx context.Context,
foo *samplev1alpha1.Foo,
deployment *appsv1.Deployment,
) error {
// NEVER modify objects from the store.
// 使用 DeepCopy 创建副本,避免修改原始缓存对象
fooCopy := foo.DeepCopy()
// 将 Deployment 的实际可用副本数同步到 Foo.Status
fooCopy.Status.AvailableReplicas = deployment.Status.AvailableReplicas
// 使用 UpdateStatus(而不是 Update)只更新 Status 字段
// FieldManager 确保该字段的更新来源可追溯
_, err := c.sampleclientset.SamplecontrollerV1alpha1().Foos(foo.Namespace).UpdateStatus(
ctx, fooCopy, metav1.UpdateOptions{FieldManager: FieldManager},
)
return err
}
注意
UpdateStatus 是 Kubernetes 提供的子资源更新操作,它只允许更新 .status 字段,.spec 和其他字段会被 API Server 忽略或拒绝。这保证了 Foo 的规格(Spec)不会被意外覆盖。如果 Foo CRD 没有启用 subresources: status(自定义资源默认没有),则需要用普通的 Update 方法。
思考记忆提示 — newDeployment 是构建 Deployment 对象的核心函数——理解它就知道如何正确设置 OwnerReference 和标签
// staging/src/k8s.io/sample-controller/controller.go(行 388~421)
// newDeployment creates a new Deployment for a Foo resource. It also sets
// the appropriate OwnerReferences on the resource so handleObject can discover
// the Foo resource that 'owns' it.
func newDeployment(foo *samplev1alpha1.Foo) *appsv1.Deployment {
// 标签必须同时在 Deployment.Spec.Selector、Deployment.Spec.Template.Labels 中一致设置
// Selector 是 Deployment 的核心标识,不能修改(修改会导致 Pod 被重建)
// Template.Labels 是 Pod 的标签,ReplicaSet 根据 Selector 创建 Pod
labels := map[string]string{
"app": "nginx", // 应用名标签(可选,由用户指定)
"controller": foo.Name, // 关键标签:以 Foo 名称为值,便于追溯
}
return &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: foo.Spec.DeploymentName,
Namespace: foo.Namespace,
// OwnerReference 是 Kubernetes GC 的核心:
// 设置后,当 Foo 被删除时,Kubernetes 会自动删除这个 Deployment
OwnerReferences: []metav1.OwnerReference{
*metav1.NewControllerRef(foo, samplev1alpha1.SchemeGroupVersion.WithKind("Foo")),
},
},
Spec: appsv1.DeploymentSpec{
// 副本数:直接使用 Foo.Spec.Replicas(可能为 nil,表示默认 1)
Replicas: foo.Spec.Replicas,
Selector: &metav1.LabelSelector{
MatchLabels: labels,
},
Template: corev1.PodTemplateSpec{
// Pod 标签必须与 Selector.MatchLabels 完全一致
// 否则 ReplicaSet 创建的 Pod 不会被 Deployment 匹配到
ObjectMeta: metav1.ObjectMeta{
Labels: labels,
},
Spec: corev1.PodSpec{
Containers: []corev1.Container{
{
Name: "nginx",
Image: "nginx:latest",
},
},
},
},
},
}
}
我的理解的意思是说
newDeployment 的关键在于三个标签的一致性:
如果 Template 标签和 Selector 标签不一致,ReplicaSet 创建的 Pod 不会被 Deployment 匹配,导致"已创建但未被管理"的状态。
思考记忆提示 — handleObject 是实现"级联调谐"的核心——Deployment 变化时如何反过来触发 Foo 的调谐
// staging/src/k8s.io/sample-controller/controller.go(行 345~383)
// handleObject will take any resource implementing metav1.Object and attempt
// to find the Foo resource that 'owns' it.
// It does this by looking at the object's metadata.ownerReferences field.
func (c *Controller) handleObject(obj interface{}) {
var object metav1.Object
var ok bool
// ========== 步骤1:从 obj 提取 metav1.Object ==========
// obj 可能是两种类型:
// 1. 直接是 metav1.Object(正常 Informer 事件)
// 2. cache.DeletedFinalStateUnknown(Indexer 中对象已被删除,但 DeleteFunc 仍被触发)
if object, ok = obj.(metav1.Object); !ok {
// Tombstone 处理:Indexer 缓存中对象已被删除,
// 但 Workqueue 中还有对它的引用(常见的 GC 场景)
tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
if !ok {
utilruntime.HandleErrorWithContext(context.Background(), nil,
"Error decoding object, invalid type", "type", fmt.Sprintf("%T", obj))
return
}
object, ok = tombstone.Obj.(metav1.Object)
if !ok {
utilruntime.HandleErrorWithContext(context.Background(), nil,
"Error decoding object tombstone, invalid type",
"type", fmt.Sprintf("%T", tombstone.Obj))
return
}
klog.V(4).Info("Recovered deleted object", "resourceName", object.GetName())
}
// ========== 步骤2:从对象中获取 OwnerReference ==========
// GetControllerOf 返回指向直接控制者的 OwnerReference
ownerRef := metav1.GetControllerOf(object)
if ownerRef == nil {
return // 该对象没有控制器(比如系统默认创建的 Deployment):忽略
}
// ========== 步骤3:检查控制器类型是否为 Foo ==========
if ownerRef.Kind != "Foo" {
return // 该对象被其他控制器管理(比如 HPA 控制 Deployment):忽略
}
// ========== 步骤4:获取拥有该 Deployment 的 Foo ==========
foo, err := c.foosLister.Foos(object.GetNamespace()).Get(ownerRef.Name)
if err != nil {
// Foo 也可能已被删除(级联删除进行中)
klog.V(4).Info("Ignore orphaned object",
"object", klog.KObj(object), "foo", ownerRef.Name)
return
}
// ========== 步骤5:将 Foo 入队,触发调谐 ==========
c.enqueueFoo(foo)
}
小贴士
DeletedFinalStateUnknown 是一种"墓碑"机制。当 Kubernetes GC 删除一个对象时,Indexer 会立即从缓存中移除该对象。但此时 Informer 的 DeleteFunc 可能还没有被处理(或者 Workqueue 中还有对它的引用)。Tombstone 就是来解决这个问题的:Indexer 在删除对象时,用一个包含原始对象副本的 DeletedFinalStateUnknown 包装替换它。Reconciler 仍然可以通过 Tombstone 拿到对象的 key 和名称。
思考记忆提示 — 这一节扩展到 Service 调谐——展示一个完整的"多资源调谐"模式
在 sample-controller 的基础模式上扩展 Service 调谐,我们需要同时管理三个资源:Foo CRD、Deployment、Service。与 Deployment 调谐类似,Service 调谐的核心是确保 Service 存在且配置正确。
// Service 调谐 Reconciler 示例(扩展模式)
func (c *Controller) syncService(ctx context.Context, foo *samplev1alpha1.Foo) error {
svcClient := c.kubeclientset.CoreV1().Services(foo.Namespace)
// 获取已有的 Service(如果存在)
existingSvc, err := c.servicesLister.Services(foo.Namespace).Get(foo.Name)
if errors.IsNotFound(err) {
// 不存在则创建
_, err = svcClient.Create(ctx, c.newService(foo), metav1.CreateOptions{
FieldManager: FieldManager,
})
return err
}
if err != nil {
return err
}
// 检查 OwnerReference(非自管理则忽略)
if !metav1.IsControlledBy(existingSvc, foo) {
return fmt.Errorf("Service %s is not controlled by Foo %s",
existingSvc.Name, foo.Name)
}
// Service 通常不需要更新(Labels/Ports 一般不变)
// 如果需要更新(如端口变了),执行 Update
return nil
}
// newService 构建一个指向 Deployment Pod 的 Service
func (c *Controller) newService(foo *samplev1alpha1.Foo) *corev1.Service {
return &corev1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: foo.Name, // Service 名称
Namespace: foo.Namespace,
OwnerReferences: []metav1.OwnerReference{
*metav1.NewControllerRef(foo, samplev1alpha1.SchemeGroupVersion.WithKind("Foo")),
},
},
Spec: corev1.ServiceSpec{
// ClusterIP: "" 表示由 API Server 自动分配
// Selector 必须与 Deployment.Spec.Template.Labels 完全一致
// 这样 Service 会自动将流量路由到 Deployment 创建的 Pod
Selector: map[string]string{
"app": "nginx",
"controller": foo.Name,
},
Ports: []corev1.ServicePort{
{
Port: 80, // Service 暴露的端口
TargetPort: intstr.IntOrString{Type: intstr.Int, IntVal: 80},
Protocol: corev1.ProtocolTCP,
},
},
},
}
}
设计精髓
Service 和 Deployment 的调谐逻辑是独立的但协调的:
这种"声明式 Selector"设计让 Kubernetes 能够自动维护 Pod 到 Service 的映射关系,无需手动管理。
// 扩展后的完整 syncHandler:同时调谐 Deployment 和 Service
func (c *Controller) syncHandlerExtended(ctx context.Context,
objectRef cache.ObjectName) error {
// 获取 Foo
foo, err := c.foosLister.Foos(objectRef.Namespace).Get(objectRef.Name)
if err != nil {
if errors.IsNotFound(err) { return nil }
return err
}
// 1. 调谐 Deployment
if err := c.syncDeployment(ctx, foo); err != nil {
return err
}
// 2. 调谐 Service
if err := c.syncService(ctx, foo); err != nil {
return err
}
// 3. 更新 Foo Status(汇总 Deployment 状态)
deployment, _ := c.deploymentsLister.Deployments(foo.Namespace).Get(foo.Spec.DeploymentName)
if deployment != nil {
c.updateFooStatus(ctx, foo, deployment)
}
c.recorder.Event(foo, corev1.EventTypeNormal, SuccessSynced, MessageResourceSynced)
return nil
}
思考记忆提示 — RBAC 是控制器的"通行证"——理解 Role/ClusterRole 的设计才能正确部署控制器
sample-controller 的 RBAC 配置位于 staging/src/k8s.io/sample-controller/artifacts/rbac/:
# staging/src/k8s.io/sample-controller/artifacts/rbac/role.yaml
# 命名空间级别的 Role:只允许访问 sample-controller 所在命名空间的资源
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: sample-controller-role
namespace: sample-controller # 指定命名空间
rules:
# ========== 操作 Foo CRD(自定义资源)==========
- apiGroups: ["samplecontroller.k8s.io"] # CRD 的 API Group
resources: ["foos"] # 资源名称(复数形式)
verbs: ["get", "list", "watch"] # 只读操作
- apiGroups: ["samplecontroller.k8s.io"]
resources: ["foos/status"] # 子资源
verbs: ["get", "update", "patch"]
# ========== 操作 Deployment(内置资源)==========
- apiGroups: ["apps"]
resources: ["deployments"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
# ========== 操作 Events(核心 API)==========
# Events 没有 apiGroup,属于核心 API(apiGroups 为 "")
- apiGroups: [""]
resources: ["events"]
verbs: ["get", "list", "watch", "create", "update", "patch"]
小贴士
sample-controller 使用 Role(命名空间级别)而不是 ClusterRole(集群级别),这是最小权限原则的体现:控制器只需要管理自己命名空间内的资源,不需要跨命名空间访问,也不需要访问集群级别的资源。如果需要跨命名空间管理,可以使用 RoleBinding 引用 ClusterRole(授予特定 ClusterRole 到本命名空间)。
# staging/src/k8s.io/sample-controller/artifacts/rbac/role-binding.yaml
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: sample-controller-role-binding
namespace: sample-controller
subjects:
# 绑定到控制器的 ServiceAccount
- kind: ServiceAccount
name: sample-controller
apiGroup: ""
roleRef:
kind: Role
name: sample-controller-role
apiGroup: rbac.authorization.k8s.io
思考记忆提示 — 事件过滤是生产级控制器的必备优化——减少不必要的调谐,提升系统效率
在实际生产环境中,同一个资源类型可能有大量对象(比如集群中有几万个 Pod),但控制器只需要管理其中一小部分(比如带有特定标签的 Pod)。通过 FilteredResourceEventHandler 在事件源头过滤,可以避免大量无用的 Workqueue 入队和 Reconciler 调用。
// FilteredResourceEventHandler 示例:只处理带有 "managed-by: my-controller" 标签的 Foo
fooInformer.Informer().AddEventHandler(cache.FilteredResourceEventHandler{
// FilterFunc:在事件处理器执行前判断是否需要处理
// 返回 true = 继续执行 Handler;返回 false = 忽略该事件
FilterFunc: func(obj interface{}) bool {
foo, ok := obj.(*samplev1alpha1.Foo)
if !ok {
return false
}
// 只处理标记了 "managed-by: my-controller" 的 Foo
managedBy, exists := foo.Labels["managed-by"]
if !exists || managedBy != "my-controller" {
return false
}
// 可选:过滤暂停中的 Foo(通过 annotation 判断)
if foo.Annotations["paused"] == "true" {
return false
}
return true
},
Handler: cache.ResourceEventHandlerFuncs{
AddFunc: controller.enqueueFoo,
UpdateFunc: func(old, new interface{}) { controller.enqueueFoo(new) },
},
})
注意
FilteredResourceEventHandler 在事件入队前过滤,好处是:Workqueue 中不会有无用的 key,减少内存占用和调谐次数。而在 Reconciler 中过滤的问题是:即使过滤掉,Workqueue 中仍然有一个 key,Worker 仍然需要 Get 和处理它(在 Reconciler 中发现不符合条件后再返回 nil)。对于大规模集群,前者效率显著更高。
思考记忆提示 — 这一节讲解一些"边缘情况"——它们不常见但很重要
// staging/src/k8s.io/sample-controller/controller.go(行 331~338)
// enqueueFoo takes a Foo resource and converts it into a namespace/name
// string which is then put onto the work queue.
func (c *Controller) enqueueFoo(obj interface{}) {
if objectRef, err := cache.ObjectToName(obj); err != nil {
// obj 可能已经被 Indexer 删除(变成 DeletedFinalStateUnknown)
// 此时无法提取 key,直接忽略
utilruntime.HandleError(err)
return
} else {
// 使用 TypedRateLimitingQueue 的 Add 方法入队
// 会经过 RateLimiter 限速处理
c.workqueue.Add(objectRef)
}
}
小贴士
cache.ObjectToName 是 k8s v1.36.1 新增的 API,它返回 cache.ObjectName(namespace/name 结构化类型),而不是字符串。在老版本中使用的是 cache.DeletionHandlingMetaNamespaceKeyFunc,返回字符串 key。新的 ObjectName 类型更安全,避免了字符串解析错误,且支持更好的类型检查。
时间线:Kubernetes GC 删除 Foo 及其拥有资源的过程
T1: API Server 收到删除 Foo 的请求
│
▼
T2: Kubernetes GC 发现 Foo 拥有 Deployment(通过 OwnerReference)
└─► 删除 Deployment
│
▼
T3: Deployment Informer 收到 Deleted 事件
│
├─► DeploymentIndexer 删除 Deployment 对象
│
▼
T4: Deployment Informer 调用 DeleteFunc(传入 DeletedFinalStateUnknown 包装)
│
▼
T5: handleObject(new) 被调用:
├─► obj 是 DeletedFinalStateUnknown { Obj: 原 Deployment 对象副本 }
├─► 提取 OwnerReference(Kind="Foo", Name="foo-xxx")
├─► foosLister.Get("foo-xxx") → NotFound(Foo 也被 GC 删除了)
└─► 安静地忽略(Orphan)
│
▼
T6: Foo 也被 GC 删除,Foo Informer 收到 Deleted 事件
│
▼
T7: FooIndexer 删除 Foo 对象(没有 handleObject 注册 DeleteFunc)
│
✓ 调谐结束
必记闭环逻辑(核心考点)
sample-controller 的核心闭环是:Foo 声明期望状态,Controller 持续调谐 Deployment 向期望状态收敛,OwnerReference 驱动 GC 级联删除,handleObject 实现级联调谐。这四个环节缺一不可,共同构成了 Kubernetes 控制器的标准模式。
思考记忆提示 — 本节预告后续文章,帮助读者建立完整的学习路线图
本篇以 sample-controller 为蓝本,深入分析了 Kubernetes 自定义控制器的完整实现。后续将分两期展开:
全篇必记总纲
sample-controller 的完整链路是:NewController(组装)→ Informer(监听)→ Workqueue(调度)→ Worker(消费)→ syncHandler(调谐)→ ClientSet(操作 API Server)→ updateFooStatus(反馈)。理解这条链路,就掌握了 Kubernetes 所有控制器的通用骨架。
思考记忆提示 — FAQ 是全篇的"临考前速背"模块,20 组覆盖全部关键知识点
Controller 有 6 个核心字段:两个 clientset(内置/CRD)、两个 Lister+Synced(缓存读取)、workqueue(任务调度)、recorder(事件记录)。kubeclientset 用于操作 Deployment 等内置资源;sampleclientset 用于操作 Foo CRD;deploymentsLister 和 foosLister 从本地缓存读取(快);deploymentsSynced 和 foosSynced 是函数类型,用于判断缓存是否已同步;workqueue 是限速队列,实现去重和重试;recorder 负责向 API Server 写入事件(kubectl describe 时可见)。
Lister 从本地 Indexer 缓存读取(快,零网络开销),ClientSet 从 API Server 读取(准确,有网络延迟)。Reconciler 中优先用 Lister 获取对象(因为 Informer 已经预热了缓存);只有当 Lister 查不到(Indexer 未同步)或需要最新数据时(写入前的乐观锁冲突检测)才用 ClientSet。
InformerSynced 是 func() bool 函数类型,不是 bool。因为缓存同步是一个过程,不是瞬间完成。HasSynced 方法每次调用都查询当前的同步状态。在 WaitForCacheSync 中,该函数会被反复轮询,直到返回 true 才解除阻塞。
只监听 Foo 不足以实现完整的控制闭环。如果 Deployment 被外部工具(如 kubectl)手动修改,Foo Informer 收不到事件,Controller 就不知道需要重新调谐。同时监听两个资源,通过 handleObject 实现"级联调谐":Deployment 的任何变化都会触发 Foo 的调谐,确保 Foo 的期望状态最终被应用。
enqueueFoo 将 Foo 对象转换为 namespace/name key 并入 Workqueue。AddFunc 在 Foo 创建时触发;UpdateFunc 在 Foo 更新时触发(直接使用 new 对象,忽略 old,因为 Workqueue 本身有去重机制)。两者都入队的原因是:Reconciler 需要重新检查 Foo 的 Spec 是否与实际 Deployment 状态一致。
handleObject 从 Deployment 的 OwnerReference 找到其拥有者 Foo 并将 Foo 入队。DeletedFinalStateUnknown(Tombstone)是 Indexer 在删除对象时留下的包装。当 Kubernetes GC 删除了 Foo,进而级联删除了 Deployment 时,Deployment Informer 的 DeleteFunc 会被触发,此时原始对象已经从 Indexer 中移除,只能从 Tombstone 中恢复。handleObject 专门处理这种场景。
OwnerReference 有三个作用:建立 Foo→Deployment 的所有权关系、驱动 GC 级联删除、供 handleObject 追溯拥有者。通过 metav1.NewControllerRef(foo, GroupVersion.WithKind("Foo")) 创建,其中包含 UID 确保唯一性。当 Foo 被删除时,Kubernetes GC 根据 OwnerReference 自动删除其拥有的 Deployment,无需控制器手动处理。
因为 Foo 删除触发 GC,GC 删除了 Deployment,Deployment 的 AddFunc(handleObject)被触发。此时 handleObject 通过 OwnerReference 查找 Foo,发现 Foo 已不存在(NotFound),安静地忽略。真正的级联删除由 Kubernetes GC 完成,不需要控制器介入。
因为 NotFound 表示 Foo 已被删除(通常是 GC 级联删除触发的),这是正常状态,无需重试。返回 nil 意味着 Workqueue.Forget 会被调用,key 被彻底清除。如果返回 error,该 key 会被重新入队(AddRateLimited),造成无意义的重试。
五步法:Get Foo → Get Deployment(不存在则Create)→ Check OwnerReference → Update Deployment(副本不匹配时)→ Update Foo Status。每一步都可能出错并返回 error(触发重试)。Create 成功后直接返回 nil,等下一次 Deployment AddEvent 触发完整的五步调谐检查。
因为 newDeployment(foo) 基于最新的 Foo.Spec 生成期望状态。Reconciler 遵循幂等性原则:每次调谐都从期望状态出发,而不是基于上一次的状态做增量修改。这样即使中间有任何变化(比如有人改了 Foo.Spec),下一次调谐都能正确收敛。
UpdateStatus 只更新 .status 字段,防止 .spec 被意外覆盖。Foo 的 .spec 定义了期望状态,是声明式契约;.status 是实际状态,反映 Deployment 的可用副本数。用 UpdateStatus 确保只有 status 发生变化,spec 保持原样。
FieldManager 记录资源的"最后修改者",用于冲突检测和审计。当多个控制器同时修改同一个资源时,API Server 根据 FieldManager 判断哪个修改是最新的。sample-controller 使用 controllerAgentName("sample-controller")作为 FieldManager,便于 kubectl apply 了解谁是修改来源。
因为 Deployment Informer 注册了 UpdateFunc 处理器。当 Deployment 被修改时,Deployment Informer 收到 Update 事件,调用 handleObject。handleObject 通过 OwnerReference 找到 Foo,将 Foo 入队,触发 syncHandler 重新检查并修正 Deployment 的实际状态向期望状态收敛。
Done() 将 key 从 processing 集合移除,允许同一个 key 被重新入队。如果一个 key 在 processing 期间被多次 Add,它只在 dirty 集合中记录,不在 queue 中重复。Done() 执行后,Workqueue 检查 dirty:如果有该 key,就重新入队。这样保证了一个 key 不会被多个 Worker 并发处理。
因为 sample-controller 只需要访问自己命名空间内的资源,不需要集群级别的权限。Role 的权限范围仅限于特定命名空间,遵循最小权限原则。如果控制器需要跨命名空间管理资源,可以使用 RoleBinding 引用 ClusterRole(授予集群级别权限到特定命名空间)。
FilteredResourceEventHandler 在入队前过滤(避免无用的 Workqueue key),Reconciler 中过滤在入队后过滤(key 仍存在但被跳过)。前者更好:减少了 Workqueue 的内存占用、减少了 Worker 的无效处理次数、减少了日志噪音。在大规模集群中,这种差异非常显著。
Service 的 Selector 必须与 Deployment.Spec.Template.Labels 完全一致,才能将流量正确路由到 Pod。当两者一致时,EndpointsController(内置)会自动创建和更新 Endpoints(记录 Pod 的 IP 和端口)。Endpoints 更新后,kube-proxy 感知变化,更新 iptables/ipvs 规则,Service 流量自动路由到新的 Pod。
每条事件都会创建 Event 资源,高频调谐场景下会占用大量 etcd 存储。建议在高吞吐量场景使用事件聚合(EventAggregator)将同类型事件合并,或直接关闭事件记录(不调用 StartRecordingToSink)。可以在事件广播器中使用 rate limiter 限制事件写入速率。
fixture 模式通过 fake clientset 和预填充 Informer Indexer 实现无网络依赖的控制器测试。测试流程:1. 用 fake clientset 创建 Controller;2. 手动将测试对象添加到 Informer Indexer(模拟缓存预热);3. 直接调用 syncHandler;4. 验证 fake clientset 中产生的 Actions(Create/Update/Delete)。这种方式不依赖真实 API Server,测试运行快且稳定。
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。