Skip to content

Operator 模式与 controller-runtime

本篇将深入讲解 Operator 开发的两大基石:controller-runtime 库的架构,以及 Kubebuilder 脚手架的使用。上一篇我们理解了 CRD 和 Operator 模式的概念,这一篇我们要把这些概念落到具体的代码框架上。controller-runtime 是几乎所有 Go Operator 的运行时底座,理解它的组件划分和工作机制,是写出正确、高效 Operator 的前提。本篇末尾我们会用 Kubebuilder 生成一个完整的项目骨架,为第三章实现第一个 Operator 做好准备。

一、controller-runtime 架构概览

controller-runtime 是 sigs.k8s.io 团队维护的一个 Go 库,它在 client-go 之上做了封装,提供了一套更简洁、更声明式的 API 来构建 Controller 和 Webhook。Kubebuilder 和 Operator SDK 生成的项目都基于它。

它的核心组件可以用下面这张图概括:

┌──────────────────────────── Manager ────────────────────────────┐
│                                                                   │
│   ┌─────────────┐   ┌─────────────┐   ┌─────────────────────┐   │
│   │  Controller  │   │  Controller │   │  Webhook Server     │   │
│   │   (Redis)    │   │  (Backup)   │   │  (Mutating/Valid)   │   │
│   └──────┬───────┘   └──────┬──────┘   └─────────────────────┘   │
│          │                  │                                     │
│          ▼                  ▼                                     │
│   ┌──────────────────────────────────┐   ┌──────────────────┐   │
│   │            Cache                  │   │     Client        │   │
│   │  (本地缓存 + Informer + Watch)    │◄─►│  (读写 API Server)│   │
│   └──────────────────────────────────┘   └────────┬─────────┘   │
│                                                   │              │
│   ┌─────────────────┐  ┌────────────────┐         │              │
│   │  Leader Election│  │ Healthz/Metrics│         │              │
│   └─────────────────┘  └────────────────┘         │              │
└───────────────────────────────────────────────────┼──────────────┘

                                            ┌───────▼───────┐
                                            │  API Server   │
                                            └───────────────┘

下面逐个讲解这些组件。

二、Manager:管理一切的容器

Manager 是 controller-runtime 的顶层对象,由 ctrl.NewManager 创建。它的职责是:

  1. 持有共享的 Client 和 Cache:所有 Controller 共用同一个 Client 和 Cache,避免重复建立 Watch 连接。
  2. 管理 Controller 和 Webhook 的生命周期:通过 manager.Add 注册的组件会随 Manager 一起启动和退出。
  3. 提供 Leader Election:Manager 内置 Leader Election 机制,多副本部署时只有一个 Manager 实例真正运行 Controller。
  4. 提供 Metrics 和 Healthz 端点:内置 Prometheus metrics 和健康检查 HTTP server。

创建 Manager 的典型代码:

go
package main

import (
    "os"

    "k8s.io/apimachinery/pkg/runtime"
    utilruntime "k8s.io/apimachinery/pkg/util/runtime"
    clientgoscheme "k8s.io/client-go/kubernetes/scheme"
    ctrl "sigs.k8s.io/controller-runtime"
    "sigs.k8s.io/controller-runtime/pkg/healthz"
    "sigs.k8s.io/controller-runtime/pkg/log/zap"
    metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"

    cachev1 "example.com/redis-operator/api/v1"
)

var (
    scheme   = runtime.NewScheme()
    setupLog = ctrl.Log.WithName("setup")
)

func init() {
    utilruntime.Must(clientgoscheme.AddToScheme(scheme))
    utilruntime.Must(cachev1.AddToScheme(scheme))
}

func main() {
    opts := zap.Options{
        Development: true,
    }
    ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts)))

    mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
        Scheme:                  scheme,
        Metrics:                 metricsserver.Options{BindAddress: "8080"},
        HealthProbeBindAddress:  "8081",
        LeaderElection:          true,
        LeaderElectionID:        "redis-operator.example.com",
        LeaderElectionResourceLock: "leases",
    })
    if err != nil {
        setupLog.Error(err, "unable to start manager")
        os.Exit(1)
    }

    // 注册健康检查
    if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
        setupLog.Error(err, "unable to set up health check")
        os.Exit(1)
    }
    if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil {
        setupLog.Error(err, "unable to set up ready check")
        os.Exit(1)
    }

    setupLog.Info("starting manager")
    if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
        setupLog.Error(err, "problem running manager")
        os.Exit(1)
    }
}

几个关键点:

  • scheme 注册了所有需要识别的资源类型,包括内置资源和自定义资源。
  • LeaderElection: true 启用 Leader Election(第七章详述)。
  • mgr.Start(ctrl.SetupSignalHandler()) 启动 Manager,并监听 SIGINT/SIGTERM 优雅退出。

三、Controller 与 Reconcile 循环

Controller 是实际执行控制逻辑的组件。每个 Controller 负责一种自定义资源,核心是一个 Reconcile 函数。Controller 通过 Watch 机制订阅资源变化事件,把事件转换为「Reconcile 请求」交给 Reconcile 函数处理。

1. Reconcile 函数签名

Reconcile 函数实现 Reconciler 接口:

go
type Reconciler interface {
    Reconcile(context.Context, Request) (Result, error)
}

type Request struct {
    // 命名空间 + 名字,定位到具体的 CR 实例
    NamespacedName types.NamespacedName
}

type Result struct {
    // 是否立即重新入队
    Requeue bool
    // 多久之后重新入队
    RequeueAfter time.Duration
}

注意一个非常重要的设计:Reconcile 函数只接收资源名,不接收资源内容。这迫使你在函数内部重新去获取最新状态,避免使用过期缓存。这是 Reconcile 幂等性的关键。

2. 注册 Controller

go
type RedisClusterReconciler struct {
    client.Client
    Scheme *runtime.Scheme
}

func (r *RedisClusterReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
    // 控制逻辑写在这里
    return ctrl.Result{}, nil
}

// SetupWithManager 注册 Controller 并配置 Watch
func (r *RedisClusterReconciler) SetupWithManager(mgr ctrl.Manager) error {
    return ctrl.NewControllerManagedBy(mgr).
        For(&cachev1.RedisCluster{}).             // 主资源
        Owns(&appsv1.Deployment{}).               // 关联资源(CR 拥有 Deployment)
        Owns(&corev1.Service{}).                  // 关联资源(CR 拥有 Service)
        Complete(r)
}

For 表示主资源(Controller 监听它),Owns 表示关联资源(CR 是它的 Owner,变化时触发所属 CR 的 Reconcile)。

四、Client:与 API Server 通信

Client 是 controller-runtime 对 API Server 读写操作的封装。它有两个子客户端:

  • Cache Client:读操作走本地缓存,避免频繁打 API Server。
  • Live Client:写操作直接走 API Server。
go
// 读取:走 Cache
var redisCluster cachev1.RedisCluster
err := r.Get(ctx, req.NamespacedName, &redisCluster)

// 创建:走 API Server
err := r.Create(ctx, deployment)

// 更新:走 API Server
err := r.Update(ctx, deployment)

// 删除:走 API Server
err := r.Delete(ctx, deployment)

// 列表
var depList appsv1.DeploymentList
err := r.List(ctx, &depList, client.InNamespace(namespace))

Reconciler 通过嵌入 client.Client 拿到这些方法。需要注意:写操作不会立即反映到 Cache,Cache 是通过 Watch 异步更新的,通常几十毫秒后才会一致。所以「先 Update 再 Get」可能拿到旧数据,这种场景要用 Patch 或在 Reconcile 末尾返回重新入队。

五、Cache:本地缓存减少 API 调用

Cache 基于 client-go 的 Informer 机制实现。它会为每种被 Watch 的资源建立一个本地缓存,并通过 List + Watch 保持缓存与 API Server 的一致性。所有读操作都从 Cache 取,极大降低 API Server 压力。

Cache 的几个特点:

  1. 按 namespace 选择性缓存:可以通过 Namespaces 选项只缓存特定 namespace,减小内存占用。
  2. 按字段选择器过滤:用 ByObject 配置 FieldSelector 只缓存符合条件对象。
  3. 启动时 List 全量:Informer 启动会做一次全量 List 建立基线,之后只处理增量事件。
go
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
    Scheme: scheme,
    Cache: cache.Options{
        // 只缓存特定 namespace
        DefaultNamespaces: map[string]cache.Config{
            "operator-system": {},
        },
        ByObject: map[client.Object]cache.Object{
            &appsv1.Deployment{}: {
                Field: fields.Everything(),
            },
        },
    },
})

六、Source:事件源

Source 定义了 Controller 的事件来源。ForOwns 内部其实就是注册了 Kind 类型的 Source。

controller-runtime 提供两类 Source:

1. Kind Source

监听集群内某种资源的变化。最常用,ForOwns 已经封装好了。

go
import (
    "sigs.k8s.io/controller-runtime/pkg/source"
    "sigs.k8s.io/controller-runtime/pkg/handler"
)

// 手动注册一个 Kind source(等价于 Owns 的内部实现)
builder := ctrl.NewControllerManagedBy(mgr).
    For(&cachev1.RedisCluster{}).
    Watches(
        &corev1.ConfigMap{},
        handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, obj client.Object) []reconcile.Request {
            // 把 ConfigMap 变化映射为 CR 的 Reconcile 请求
            return []reconcile.Request{
                {NamespacedName: types.NamespacedName{
                    Name:      obj.Labels["redis-cluster"],
                    Namespace: obj.Namespace,
                }},
            }
        }),
    )

2. Channel Source

监听集群外部的事件源,比如外部消息队列、定时器、文件变更。事件通过 Go channel 投递。

go
import (
    "sigs.k8s.io/controller-runtime/pkg/source"
    "k8s.io/apimachinery/pkg/types"
)

ch := make(chan event.GenericEvent, 100)

// 在某个 goroutine 里投递事件
go func() {
    for range time.Tick(30 * time.Second) {
        ch <- event.GenericEvent{
            Object: &cachev1.RedisCluster{
                ObjectMeta: metav1.ObjectMeta{
                    Name:      "my-redis",
                    Namespace: "default",
                },
            },
        }
    }
}()

builder.WatchesRawSource(source.Channel(ch, &handler.EnqueueRequestForObject{}))

Channel source 适合做「定时同步」或「外部触发」场景,比如定时对账、接收外部告警后触发 Reconcile。

七、EventHandler:事件处理

EventHandler 决定一个事件如何被转换为 Reconcile 请求。controller-runtime 提供几种内置实现:

  • EnqueueRequestForObject:把变化对象本身作为 Reconcile 请求。For 默认用这个。
  • EnqueueRequestForOwner:把对象的 Owner 作为 Reconcile 请求。Owns 默认用这个。
  • EnqueueRequestsFromMapFunc:用自定义函数把任意对象映射为一组 Reconcile 请求。用于 Watch 不被 CR 拥有的资源(如 ConfigMap、Secret)。
go
// 示例:监听一个共享 Secret,变化时触发所有引用它的 RedisCluster 重新 Reconcile
handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, obj client.Object) []reconcile.Request {
    secret := obj.(*corev1.Secret)
    var requests []reconcile.Request

    // 列出所有 RedisCluster,找出引用了这个 Secret 的
    var list cachev1.RedisClusterList
    if err := r.List(ctx, &list, client.InNamespace(secret.Namespace)); err != nil {
        return nil
    }
    for _, rc := range list.Items {
        if rc.Spec.PasswordSecret == secret.Name {
            requests = append(requests, reconcile.Request{
                NamespacedName: types.NamespacedName{
                    Name:      rc.Name,
                    Namespace: rc.Namespace,
                },
            })
        }
    }
    return requests
})

八、Predicate:事件过滤

Predicate 是事件过滤器,决定一个事件是否要触发 Reconcile。合理使用 Predicate 可以避免无意义的 Reconcile,显著降低 CPU 和 API 调用。

controller-runtime 内置了常用 Predicate:

  • GenerationChangedPredicate:只有 metadata.generation 变化才触发。generationspec 变化时才递增,所以这能过滤掉纯 statuslabel 变化。这是最常用的 Predicate。
  • ResourceVersionChangedPredicate:任何变化都触发。
  • LabelChangedPredicate:只 label 变化触发。
  • AnnotationChangedPredicate:只 annotation 变化触发。
go
import (
    "sigs.k8s.io/controller-runtime/pkg/predicate"
)

builder := ctrl.NewControllerManagedBy(mgr).
    For(&cachev1.RedisCluster{}, builder.WithPredicates(predicate.GenerationChangedPredicate{})).
    Owns(&appsv1.Deployment{}).
    Owns(&corev1.Service{})

也可以自定义 Predicate:

go
import (
    "sigs.k8s.io/controller-runtime/pkg/event"
    "sigs.k8s.io/controller-runtime/pkg/predicate"
)

// 只在 spec.size 变化时触发
type SizeChangedPredicate struct {
    predicate.Funcs
}

func (SizeChangedPredicate) Update(e event.UpdateEvent) bool {
    oldRC, ok := e.ObjectOld.(*cachev1.RedisCluster)
    if !ok {
        return false
    }
    newRC, ok := e.ObjectNew.(*cachev1.RedisCluster)
    if !ok {
        return false
    }
    return oldRC.Spec.Size != newRC.Spec.Size
}

九、Reconcile 循环原理

1. 声明式 vs 命令式

Reconcile 循环是声明式的:它不关心「发生了什么事件」,只关心「当前期望状态是什么、实际状态是什么、差异在哪」。无论触发原因是用户改了 spec、还是 Pod 被人为删除、还是节点故障,Reconcile 的处理逻辑都一样——把实际状态向期望状态收敛。

这种设计带来一个重要好处:幂等。同一个 CR 实例无论 Reconcile 多少次,结果都一样。哪怕中间发生了网络抖动、重启,只要最终 Reconcile 跑完,状态就是对的。

2. Level-Triggered vs Edge-Triggered

这是两个来自电子学的术语,在 K8s 语境下含义如下:

  • Edge-Triggered(边沿触发):在「事件发生的瞬间」做出反应。比如「副本数从 3 变成 2」这个事件触发一次扩容。问题是如果错过了这个事件(比如 Controller 重启),就永远不会补副本了。
  • Level-Triggered(电平触发):在「状态存在差异期间」持续做出反应。Controller 不关心事件是否错过,它只看「现在是 2 个副本,期望是 3 个」就补一个。

K8s 的 Controller 都是 Level-Triggered 的,这也是为什么 Reconcile 函数只接收资源名——它每次都重新评估整体状态,而不依赖「上一次发生了什么」。这是 Operator 健壮性的核心来源。

3. 幂等性要求

写 Reconcile 必须保证幂等,这意味着:

  • 创建前先检查是否存在:用 Get 判断,存在就跳过 Create,避免 AlreadyExists 错误。
  • 更新前比较是否真的变了:避免无意义 Update 引起 generation 变化和无限 Reconcile。
  • 不要维护跨调用的内存状态:所有状态都从 API Server 获取,因为 Controller 随时可能重启。
  • 错误要可重试:返回 error 让 Controller 自动重试,不要 panic。

第四章会用完整代码演示这些原则。

十、controller-runtime vs client-go

很多初学者会困惑:既然有 client-go,为什么还要 controller-runtime?二者关系如下:

维度client-gocontroller-runtime
定位K8s 官方底层客户端库在 client-go 之上的 Operator 框架
抽象层级低,需手动管理 Informer/Workqueue高,Manager/Controller 封装好一切
代码量写一个 Controller 几百行同样功能几十行
Boilerplate多(Informer、Lister、Queue、Worker 都要自己写)少(只需写 Reconcile 函数)
学习曲线陡峭平缓
灵活性高,可完全自定义中,常见场景已封装

建议:除非有特殊定制需求,否则一律用 controller-runtime。它把 Informer、Workqueue、Leader Election、Metrics 这些样板代码都封装好了,让你专注于业务逻辑。本系列所有代码都基于 controller-runtime。

十一、Kubebuilder 脚手架

讲完理论,下面动手生成项目骨架。

1. kubebuilder init

先创建项目目录并初始化:

bash
mkdir -p redis-operator
cd redis-operator
go mod init example.com/redis-operator

# 初始化项目
kubebuilder init \
    --domain example.com \
    --repo example.com/redis-operator

执行后会生成如下结构:

redis-operator/
├── cmd/
│   └── main.go              # 入口
├── api/                     # CRD 类型定义(create api 后生成)
├── internal/
│   └── controller/          # Controller 实现(create api 后生成)
├── config/
│   ├── crd/                 # CRD 部署 YAML
│   ├── default/             # 默认 kustomize 配置
│   ├── manager/             # Operator Deployment YAML
│   ├── rbac/                # RBAC 权限 YAML
│   └── samples/             # CR 示例 YAML
├── Dockerfile
├── Makefile
├── go.mod
├── go.sum
└── PROJECT                  # Kubebuilder 项目元数据

2. kubebuilder create api

接着创建一个 API(CRD + Controller):

bash
kubebuilder create api \
    --group cache \
    --version v1 \
    --kind RedisCluster \
    --resource \
    --controller

执行后会增加:

api/v1/
├── rediscluster_types.go    # Spec/Status 类型定义
├── groupversion_info.go     # GroupVersion、SchemeBuilder
└── zz_generated.deepcopy.go # 自动生成的 DeepCopy 代码

internal/controller/
└── rediscluster_controller.go  # Reconcile 实现

同时 config/samples/ 下会生成 cache_v1_rediscluster.yaml 示例。

3. 生成 CRD 与代码

Kubebuilder 项目用 Makefile 封装了常用命令:

bash
# 生成 deepcopy 代码和 CRD manifests
make manifests

# 生成 deepcopy、clientset、informer 等(按需)
make generate

make manifests 会调用 controller-gen,从 api/v1/rediscluster_types.go 里的 +kubebuilder 注释生成 config/crd/bases/cache.example.com_redisclusters.yaml

4. 安装与运行

bash
# 把 CRD 安装到集群
make install

# 本地运行 Controller(连接当前 kubectl 上下文的集群)
make run

十二、项目结构详解

下面逐个解释关键文件的作用,这些是后续章节反复操作的文件。

1. cmd/main.go

程序入口。负责创建 Manager、注册 Scheme、注册 Controller、启动 Manager。后续添加 Webhook、Leader Election 配置都在这里改。

go
func main() {
    // ... flag 解析、日志初始化 ...

    mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
        Scheme:  scheme,
        Metrics: metricsserver.Options{BindAddress: metricsAddr},
        WebhookServer: webhook.NewServer(webhook.Options{
            Port: 9443,
        }),
        HealthProbeBindAddress: probeAddr,
        LeaderElection:         enableLeaderElection,
        LeaderElectionID:       "redis-operator.example.com",
    })
    // ...

    if err := (&controllers.RedisClusterReconciler{
        Client: mgr.GetClient(),
        Scheme: mgr.GetScheme(),
    }).SetupWithManager(mgr); err != nil {
        setupLog.Error(err, "unable to create controller", "controller", "RedisCluster")
        os.Exit(1)
    }

    if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
        os.Exit(1)
    }
}

2. api/v1/rediscluster_types.go

定义 CRD 的 Go 类型。这是 CRD 的「源代码」,CRD 的 YAML 由它生成。

go
package v1

import (
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)

// RedisClusterSpec 定义期望状态
type RedisClusterSpec struct {
    // +kubebuilder:validation:Minimum=1
    // +kubebuilder:validation:Maximum=10
    Size int32 `json:"size"`
    // +kubebuilder:validation:Required
    Image string `json:"image"`
    Port  int32 `json:"port,omitempty"`
}

// RedisClusterStatus 定义实际状态
type RedisClusterStatus struct {
    Phase string   `json:"phase,omitempty"`
    Nodes []string `json:"nodes,omitempty"`
}

// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
type RedisCluster struct {
    metav1.TypeMeta   `json:",inline"`
    metav1.ObjectMeta `json:"metadata,omitempty"`
    Spec   RedisClusterSpec   `json:"spec,omitempty"`
    Status RedisClusterStatus `json:"status,omitempty"`
}

// +kubebuilder:object:root=true
type RedisClusterList struct {
    metav1.TypeMeta `json:",inline"`
    metav1.ListMeta `json:"metadata,omitempty"`
    Items           []RedisCluster `json:"items"`
}

几个关键标记:

  • +kubebuilder:object:root=true:标记这是一个顶层资源类型。
  • +kubebuilder:subresource:status:启用 status 子资源(第五章详述)。
  • +kubebuilder:validation:Minimum=1:生成 OpenAPI 校验规则。

3. internal/controller/rediscluster_controller.go

Reconcile 逻辑所在文件。这是 Operator 的「大脑」,后续章节主要在这个文件里写代码。

4. config/ 目录

部署相关的 YAML 集合,通过 kustomize 组织:

  • crd/bases/:CRD 定义,make install 部署的就是它。
  • samples/:CR 示例,用于测试。
  • manager/:Operator 自己的 Deployment。
  • rbac/:Operator 运行所需的 ClusterRole、RoleBinding。
  • default/:kustomize 入口,把以上组合成完整部署清单。

5. Makefile

封装了所有常用操作,最重要几个 target:

Target作用
make manifests生成 CRD YAML
make generate生成 deepcopy 等代码
make install安装 CRD 到集群
make run本地运行 Controller
make docker-build构建 Operator 镜像
make deploy部署 Operator 到集群
make undeploy卸载 Operator
make test运行测试(含 envtest)

十三、小结

本篇系统讲解了 controller-runtime 的架构和 Kubebuilder 的使用。要点回顾:

  1. Manager 是顶层容器,持有共享的 Client/Cache,管理 Controller 和 Webhook 生命周期,提供 Leader Election 和健康检查。
  2. Controller 的核心是 Reconcile 循环,它通过 Watch 资源变化并把事件转换为 Reconcile 请求。Reconcile 只接收资源名,强制每次重新获取状态,保证幂等。
  3. Client 分读写两路:读走 Cache(Informer 本地缓存),写走 API Server。
  4. Source/EventHandler/Predicate 三件套控制「监听谁、如何映射、是否过滤」。For/Owns 是它们的高层封装。
  5. Level-Triggered 是 K8s 控制器的本质:不依赖事件是否被错过,每次都重新评估整体状态。
  6. controller-runtime 优于裸 client-go:屏蔽了 Informer/Workqueue 样板,专注业务逻辑。
  7. Kubebuilder 提供 init + create api 两步生成完整项目骨架,Makefile 封装常用命令。

下一篇我们就在这个骨架上,从零实现一个能真正管理 Redis 集群的 Operator,把理论变成可运行的代码。