Skip to content

第一个 Operator:从零开始

本篇我们动手实现第一个真正的 Operator——一个 Redis Operator。它管理 RedisCluster 自定义资源,根据用户声明的期望状态自动创建并维护对应的 Deployment 和 Service。这是把前两章理论变成代码的关键一步。学完本篇,你将掌握 Operator 开发的完整闭环:定义 API、生成 CRD、实现 Reconcile、本地运行、部署验证。后续章节的状态管理、Finalizer、Webhook 都是在这个基础上叠加的。

一、需求分析

1. 我们要做什么

我们要构建一个 Redis Operator,让用户能用一行声明创建一个 Redis 集群:

yaml
apiVersion: cache.example.com/v1
kind: RedisCluster
metadata:
  name: my-redis
spec:
  size: 3              # 副本数
  image: redis:7.0     # 镜像
  port: 6379           # 端口
  password: "secret"   # 密码(演示用,生产应从 Secret 引用)

提交这个 CR 后,Operator 应当:

  1. 创建一个 Deployment,副本数等于 spec.size,使用 spec.image 镜像。
  2. 创建一个 Service,暴露 spec.port
  3. status.phase 更新为 Runningstatus.nodes 填上实际 Pod 名。
  4. 当用户修改 spec.size 时,自动扩缩容。
  5. 当用户删除 RedisCluster 时,关联资源被级联删除。

2. 设计要点

在写代码前,先理清几个设计决策:

  • CRD 名称RedisCluster,API Group cache.example.com,Version v1
  • 关联方式:Deployment 和 Service 的 ownerReferences 指向 RedisCluster,实现级联删除。
  • 命名规则:关联资源名字与 CR 名字相同,便于通过 req.NamespacedName 定位。
  • 标签:所有关联资源打上 app.kubernetes.io/managed-by: redis-operatorapp.kubernetes.io/instance: <cr-name>,便于查询和过滤。
  • 状态:本篇先用简单的 phase + nodes,第五章会升级为完整的 Conditions 模式。

二、初始化项目

按第二章的方法初始化项目并创建 API:

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

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

执行完后项目结构如下,我们主要改三个文件:

  • api/v1/rediscluster_types.go:定义 Spec 和 Status。
  • internal/controller/rediscluster_controller.go:实现 Reconcile。
  • cmd/main.go:确认入口正确(一般 Kubebuilder 已生成好)。

三、定义 API:types.go

打开 api/v1/rediscluster_types.go,把 Spec 和 Status 改成下面这样。+kubebuilder 注释会被 controller-gen 解析成 CRD 的 OpenAPI 校验规则。

go
package v1

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

// RedisClusterSpec 定义 Redis 集群的期望状态
type RedisClusterSpec struct {
    // Size 是 Redis 副本数
    // +kubebuilder:validation:Minimum=1
    // +kubebuilder:validation:Maximum=10
    // +kubebuilder:default=1
    Size int32 `json:"size"`

    // Image 是 Redis 镜像地址
    // +kubebuilder:validation:Required
    Image string `json:"image"`

    // Port 是 Redis 服务端口
    // +kubebuilder:validation:Minimum=1
    // +kubebuilder:validation:Maximum=65535
    // +kubebuilder:default=6379
    Port int32 `json:"port,omitempty"`

    // Password 是 Redis 访问密码(演示用,生产环境请改用 passwordSecret)
    // +kubebuilder:validation:MinLength=0
    Password string `json:"password,omitempty"`
}

// RedisClusterStatus 定义 Redis 集群的实际状态
type RedisClusterStatus struct {
    // Phase 表示集群当前阶段:Pending / Running / Failed
    // +kubebuilder:default=Pending
    Phase string `json:"phase,omitempty"`

    // Nodes 是当前实际运行的 Pod 名列表
    Nodes []string `json:"nodes,omitempty"`

    // ReadyReplicas 是就绪副本数
    ReadyReplicas int32 `json:"readyReplicas,omitempty"`
}

// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
// +kubebuilder:resource:shortName=rc
// RedisCluster 是一个由 Operator 管理的 Redis 集群
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
// RedisClusterList 包含一组 RedisCluster
type RedisClusterList struct {
    metav1.TypeMeta `json:",inline"`
    metav1.ListMeta `json:"metadata,omitempty"`
    Items           []RedisCluster `json:"items"`
}

func init() {
    SchemeBuilder.Register(&RedisCluster{}, &RedisClusterList{})
}

字段说明:

  • Size:必须 1-10 之间,默认 1。
  • Image:必填,无默认值。
  • Port:1-65535,默认 6379。
  • +kubebuilder:subresource:status:开启 status 子资源,使 status 必须通过单独的 Status().Update() 接口更新,且 status 变化不会触发 Reconcile(第五章详解)。
  • +kubebuilder:resource:shortName=rc:让 kubectl get rc 能查到 RedisCluster。

四、生成 CRD:make manifests

定义完类型后,运行:

bash
make manifests
make generate

make generate 会生成 zz_generated.deepcopy.go,让 RedisCluster 实现 runtime.Object 接口(必须有 DeepCopy 方法)。make manifests 生成 config/crd/bases/cache.example.com_redisclusters.yaml,内容大致如下:

yaml
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: redisclusters.cache.example.com
spec:
  group: cache.example.com
  names:
    kind: RedisCluster
    listKind: RedisClusterList
    plural: redisclusters
    singular: rediscluster
    shortNames:
      - rc
  scope: Namespaced
  versions:
    - name: v1
      served: true
      storage: true
      schema:
        openAPIV3Schema:
          description: RedisCluster 是一个由 Operator 管理的 Redis 集群
          properties:
            apiVersion:
              type: string
            kind:
              type: string
            metadata:
              type: object
            spec:
              properties:
                image:
                  type: string
                password:
                  minLength: 0
                  type: string
                port:
                  default: 6379
                  format: int32
                  maximum: 65535
                  minimum: 1
                  type: integer
                size:
                  default: 1
                  format: int32
                  maximum: 10
                  minimum: 1
                  type: integer
              required:
                - image
                - size
              type: object
            status:
              properties:
                nodes:
                  items:
                    type: string
                  type: array
                phase:
                  default: Pending
                  type: string
                readyReplicas:
                  format: int32
                  type: integer
              type: object
          required:
            - spec
          type: object
      subresources:
        status: {}

注意末尾的 subresources.status: {},这就是 +kubebuilder:subresource:status 标记生成的,它启用了 status 子资源。

五、实现 Reconcile 循环

现在进入核心部分——实现 Reconcile。打开 internal/controller/rediscluster_controller.go

1. 完整代码

下面是完整的 Reconcile 实现,注释里标注了每个步骤的目的:

go
package controller

import (
    "context"
    "fmt"
    "strings"

    appsv1 "k8s.io/api/apps/v1"
    corev1 "k8s.io/api/core/v1"
    "k8s.io/apimachinery/pkg/api/errors"
    "k8s.io/apimachinery/pkg/api/resource"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/runtime"
    "k8s.io/apimachinery/pkg/types"
    "k8s.io/apimachinery/pkg/util/intstr"
    ctrl "sigs.k8s.io/controller-runtime"
    "sigs.k8s.io/controller-runtime/pkg/client"
    "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
    "sigs.k8s.io/controller-runtime/pkg/log"

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

const (
    // 最终器名称,第六章会用到
    redisFinalizer = "cache.example.com/redis-finalizer"

    // 通用标签
    managedByLabel = "app.kubernetes.io/managed-by"
    instanceLabel  = "app.kubernetes.io/instance"
)

// RedisClusterReconciler 负责调谐 RedisCluster 资源
type RedisClusterReconciler struct {
    client.Client
    Scheme *runtime.Scheme
}

// +kubebuilder:rbac:groups=cache.example.com,resources=redisclusters,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=cache.example.com,resources=redisclusters/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=cache.example.com,resources=redisclusters/finalizers,verbs=update
// +kubebuilder:rbac:groups=apps,resources=deployments,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=core,resources=services,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=core,resources=pods,verbs=get;list;watch
// +kubebuilder:rbac:groups=core,resources=events,verbs=create;patch

// Reconcile 是控制循环的核心入口
func (r *RedisClusterReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
    logger := log.FromContext(ctx)

    // 步骤 1:获取 RedisCluster 实例
    var redisCluster cachev1.RedisCluster
    if err := r.Get(ctx, req.NamespacedName, &redisCluster); err != nil {
        if errors.IsNotFound(err) {
            // CR 已被删除,无需处理(级联删除会清理关联资源)
            logger.Info("RedisCluster 已删除,跳过 Reconcile", "name", req.NamespacedName)
            return ctrl.Result{}, nil
        }
        return ctrl.Result{}, err
    }

    logger.Info("开始 Reconcile", "name", redisCluster.Name, "namespace", redisCluster.Namespace,
        "size", redisCluster.Spec.Size, "image", redisCluster.Spec.Image)

    // 步骤 2:确保 Deployment 存在且符合期望
    deploymentResult, err := r.reconcileDeployment(ctx, &redisCluster)
    if err != nil {
        return ctrl.Result{}, err
    }

    // 步骤 3:确保 Service 存在且符合期望
    if err := r.reconcileService(ctx, &redisCluster); err != nil {
        return ctrl.Result{}, err
    }

    // 步骤 4:更新 Status
    if err := r.updateStatus(ctx, &redisCluster); err != nil {
        return ctrl.Result{}, err
    }

    // 步骤 5:若尚未就绪,稍后重新入队继续观察
    if deploymentResult.NeedRequeue {
        return ctrl.Result{RequeueAfter: 10}, nil
    }

    return ctrl.Result{}, nil
}

// reconcileDeployment 确保 Deployment 存在且 spec 符合期望
func (r *RedisClusterReconciler) reconcileDeployment(ctx context.Context, rc *cachev1.RedisCluster) (reconcileResult, error) {
    logger := log.FromContext(ctx)

    desired := r.buildDeployment(rc)

    var existing appsv1.Deployment
    err := r.Get(ctx, types.NamespacedName{Name: rc.Name, Namespace: rc.Namespace}, &existing)
    if err != nil && errors.IsNotFound(err) {
        // 不存在,创建
        if err := controllerutil.SetControllerReference(rc, desired, r.Scheme); err != nil {
            return reconcileResult{}, err
        }
        logger.Info("创建 Deployment", "name", desired.Name)
        if err := r.Create(ctx, desired); err != nil {
            return reconcileResult{}, err
        }
        return reconcileResult{NeedRequeue: true}, nil
    }
    if err != nil {
        return reconcileResult{}, err
    }

    // 已存在,比较关键字段,按需更新
    updated := false
    if *existing.Spec.Replicas != rc.Spec.Size {
        existing.Spec.Replicas = &rc.Spec.Size
        updated = true
    }
    if len(existing.Spec.Template.Spec.Containers) > 0 &&
        existing.Spec.Template.Spec.Containers[0].Image != rc.Spec.Image {
        existing.Spec.Template.Spec.Containers[0].Image = rc.Spec.Image
        updated = true
    }
    // 比较端口
    if len(existing.Spec.Template.Spec.Containers) > 0 {
        container := &existing.Spec.Template.Spec.Containers[0]
        if len(container.Ports) == 0 || container.Ports[0].ContainerPort != rc.Spec.Port {
            container.Ports = []corev1.ContainerPort{{ContainerPort: rc.Spec.Port, Protocol: corev1.ProtocolTCP}}
            updated = true
        }
    }

    if updated {
        logger.Info("更新 Deployment", "name", existing.Name)
        if err := r.Update(ctx, &existing); err != nil {
            return reconcileResult{}, err
        }
    }

    // 判断是否就绪
    if existing.Status.ReadyReplicas != rc.Spec.Size {
        return reconcileResult{NeedRequeue: true}, nil
    }
    return reconcileResult{}, nil
}

// reconcileService 确保 Service 存在且 spec 符合期望
func (r *RedisClusterReconciler) reconcileService(ctx context.Context, rc *cachev1.RedisCluster) error {
    logger := log.FromContext(ctx)

    desired := r.buildService(rc)

    var existing corev1.Service
    err := r.Get(ctx, types.NamespacedName{Name: rc.Name, Namespace: rc.Namespace}, &existing)
    if err != nil && errors.IsNotFound(err) {
        if err := controllerutil.SetControllerReference(rc, desired, r.Scheme); err != nil {
            return err
        }
        logger.Info("创建 Service", "name", desired.Name)
        return r.Create(ctx, desired)
    }
    if err != nil {
        return err
    }

    // Service 的端口可更新
    if len(existing.Spec.Ports) == 0 || existing.Spec.Ports[0].Port != rc.Spec.Port {
        existing.Spec.Ports = []corev1.ServicePort{
            {
                Port:       rc.Spec.Port,
                TargetPort: intstr.FromInt(int(rc.Spec.Port)),
                Protocol:   corev1.ProtocolTCP,
            },
        }
        logger.Info("更新 Service", "name", existing.Name)
        return r.Update(ctx, &existing)
    }
    return nil
}

// updateStatus 收集实际状态并更新 status 子资源
func (r *RedisClusterReconciler) updateStatus(ctx context.Context, rc *cachev1.RedisCluster) error {
    logger := log.FromContext(ctx)

    // 查询关联 Pod
    var pods corev1.PodList
    if err := r.List(ctx, &pods, client.InNamespace(rc.Namespace),
        client.MatchingLabels{instanceLabel: rc.Name}); err != nil {
        return err
    }

    var nodes []string
    var ready int32
    for _, pod := range pods.Items {
        nodes = append(nodes, pod.Name)
        if isPodReady(&pod) {
            ready++
        }
    }

    // 计算 phase
    phase := "Pending"
    if ready == rc.Spec.Size && rc.Spec.Size > 0 {
        phase = "Running"
    } else if ready > 0 {
        phase = "Scaling"
    }

    // 只有变化时才更新,避免无限循环
    if rc.Status.Phase != phase || rc.Status.ReadyReplicas != ready || !sameStringSlice(rc.Status.Nodes, nodes) {
        logger.Info("更新 Status", "phase", phase, "ready", ready, "total", rc.Spec.Size)
        rc.Status.Phase = phase
        rc.Status.ReadyReplicas = ready
        rc.Status.Nodes = nodes
        // 注意:开了 status 子资源后必须用 Status().Update()
        return r.Status().Update(ctx, rc)
    }
    return nil
}

// buildDeployment 根据期望状态构造 Deployment 对象
func (r *RedisClusterReconciler) buildDeployment(rc *cachev1.RedisCluster) *appsv1.Deployment {
    labels := labelsForRedis(rc.Name)
    size := rc.Spec.Size

    dep := &appsv1.Deployment{
        ObjectMeta: metav1.ObjectMeta{
            Name:      rc.Name,
            Namespace: rc.Namespace,
            Labels:    labels,
        },
        Spec: appsv1.DeploymentSpec{
            Replicas: &size,
            Selector: &metav1.LabelSelector{MatchLabels: labels},
            Template: corev1.PodTemplateSpec{
                ObjectMeta: metav1.ObjectMeta{Labels: labels},
                Spec: corev1.PodSpec{
                    Containers: []corev1.Container{{
                        Name:  "redis",
                        Image: rc.Spec.Image,
                        Ports: []corev1.ContainerPort{{ContainerPort: rc.Spec.Port, Protocol: corev1.ProtocolTCP}},
                        Resources: corev1.ResourceRequirements{
                            Limits: corev1.ResourceList{
                                corev1.ResourceCPU:    resource.MustParse("500m"),
                                corev1.ResourceMemory: resource.MustParse("256Mi"),
                            },
                            Requests: corev1.ResourceList{
                                corev1.ResourceCPU:    resource.MustParse("100m"),
                                corev1.ResourceMemory: resource.MustParse("128Mi"),
                            },
                        },
                        // 演示:通过命令设置密码,生产应从 Secret 注入
                        Args: []string{fmt.Sprintf("--requirepass %s", rc.Spec.Password)},
                    }},
                },
            },
        },
    }
    return dep
}

// buildService 根据期望状态构造 Service 对象
func (r *RedisClusterReconciler) buildService(rc *cachev1.RedisCluster) *corev1.Service {
    labels := labelsForRedis(rc.Name)
    return &corev1.Service{
        ObjectMeta: metav1.ObjectMeta{
            Name:      rc.Name,
            Namespace: rc.Namespace,
            Labels:    labels,
        },
        Spec: corev1.ServiceSpec{
            Type:     corev1.ServiceTypeClusterIP,
            Selector: labels,
            Ports: []corev1.ServicePort{{
                Port:       rc.Spec.Port,
                TargetPort: intstr.FromInt(int(rc.Spec.Port)),
                Protocol:   corev1.ProtocolTCP,
            }},
        },
    }
}

// labelsForRedis 生成通用标签
func labelsForRedis(name string) map[string]string {
    return map[string]string{
        "app.kubernetes.io/name":       "redis",
        "app.kubernetes.io/instance":   name,
        "app.kubernetes.io/managed-by": "redis-operator",
    }
}

// isPodReady 判断 Pod 是否就绪
func isPodReady(pod *corev1.Pod) bool {
    for _, cond := range pod.Status.Conditions {
        if cond.Type == corev1.PodReady && cond.Status == corev1.ConditionTrue {
            return true
        }
    }
    return false
}

// sameStringSlice 判断两个字符串切片内容是否一致(不考虑顺序)
func sameStringSlice(a, b []string) bool {
    if len(a) != len(b) {
        return false
    }
    set := make(map[string]struct{}, len(a))
    for _, s := range a {
        set[s] = struct{}{}
    }
    for _, s := range b {
        if _, ok := set[s]; !ok {
            return false
        }
    }
    return true
}

// reconcileResult 用于内部传递是否需要重新入队
type reconcileResult struct {
    NeedRequeue bool
}

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

// 为了避免编译器对未使用导入的报错(实际 strings 在事件上报等场景使用),保留
var _ = strings.TrimSpace

2. 关键设计解读

步骤 1:获取 CR

Reconcile 一开始用 r.Get 获取 CR 实例。如果返回 NotFound,说明 CR 已被删除,直接返回——因为关联资源设置了 ownerReferences,会被 K8s 自动级联删除,无需我们手动处理。这是 Level-Triggered 的体现:我们不关心「删除事件」,只看「现在还有没有这个 CR」。

步骤 2-3:幂等地创建/更新关联资源

每个关联资源的处理都遵循同一模式:

  1. Get 看是否存在。
  2. 不存在则 SetControllerReference(设置 owner)后 Create
  3. 存在则比较关键字段,有差异才 Update

SetControllerReference 是级联删除的关键。它把 RedisCluster 设为 Deployment/Service 的 Owner,RedisCluster 删除时它们会被自动删除。

步骤 4:更新 Status

注意必须用 r.Status().Update(ctx, rc) 而不是 r.Update,因为我们在 CRD 里开启了 status 子资源。r.Update 不能改 status 字段。

更新前先比较新旧值,避免无意义写入。Status 写入本身会改变 resourceVersion,如果不加判断可能引起无限 Reconcile(虽然开了 status 子资源后 status 变化不再触发本 CR 的 Reconcile,但仍是好习惯,减少 API 调用)。

步骤 5:重新入队

如果 Deployment 还没就绪(ReadyReplicas != Spec.Size),返回 ctrl.Result{RequeueAfter: 10} 让 Controller 10 秒后再 Reconcile 一次。这样无需依赖额外事件,靠定时轮询也能最终观察到就绪状态。这是 Operator 推进状态的常用手段。

六、main.go 入口

Kubebuilder 生成的 cmd/main.go 一般无需改动,确认它正确注册了 Scheme 和 Reconciler 即可:

go
package main

import (
    "flag"
    "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"
    "example.com/redis-operator/internal/controller"
)

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

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

func main() {
    var metricsAddr string
    var enableLeaderElection bool
    var probeAddr string
    flag.StringVar(&metricsAddr, "metrics-bind-address", "8080", "Metrics 监听地址")
    flag.StringVar(&probeAddr, "health-probe-bind-address", "8081", "健康检查监听地址")
    flag.BoolVar(&enableLeaderElection, "leader-elect", false, "启用 Leader Election")
    opts := zap.Options{Development: true}
    opts.BindFlags(flag.CommandLine)
    flag.Parse()

    ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts)))

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

    if err := (&controller.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.AddHealthzCheck("healthz", healthz.Ping); err != nil {
        os.Exit(1)
    }
    if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil {
        os.Exit(1)
    }

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

七、本地运行:make install + make run

确保 kind 集群已启动,kubectl config current-context 指向它。然后:

bash
# 1. 生成 CRD 与代码
make manifests
make generate

# 2. 安装 CRD 到集群
make install

# 3. 本地运行 Operator(前台进程,连集群)
make run

make run 后你会看到日志输出:

starting manager
starting metrics server
Starting Controller
Starting workers

Operator 正在前台监听集群事件。

八、部署 CRD 实例并验证

打开另一个终端,创建 CR 实例。先看看 config/samples/cache_v1_rediscluster.yaml 默认内容,改成:

yaml
apiVersion: cache.example.com/v1
kind: RedisCluster
metadata:
  labels:
    app.kubernetes.io/name: rediscluster
  name: my-redis
  namespace: default
spec:
  size: 3
  image: redis:7.0
  port: 6379
  password: "secret123"

应用它:

bash
kubectl apply -f config/samples/cache_v1_rediscluster.yaml

观察 Operator 日志,应该看到:

开始 Reconcile  name=my-redis namespace=default size=3 image=redis:7.0
创建 Deployment  name=my-redis
创建 Service     name=my-redis
更新 Status      phase=Pending ready=0 total=3
...
更新 Status      phase=Running ready=3 total=3

验证集群资源:

bash
# 查看 CR
kubectl get redisclusters -o wide
# NAME       PHASE     SIZE   READY
# my-redis   Running   3      3

# 查看 Deployment
kubectl get deployment my-redis
# NAME       READY   UP-TO-DATE   AVAILABLE
# my-redis   3/3     3            3

# 查看 Service
kubectl get service my-redis
# NAME       TYPE        CLUSTER-IP     PORT(S)
# my-redis   ClusterIP   10.96.10.10    6379/TCP

# 查看 Pod
kubectl get pods -l app.kubernetes.io/instance=my-redis
# NAME                       READY   STATUS
# my-redis-xxxx-aaaa         1/1     Running
# my-redis-xxxx-bbbb         1/1     Running
# my-redis-xxxx-cccc         1/1     Running

# 查看 status
kubectl get rediscluster my-redis -o jsonpath='{.status}'
# {"phase":"Running","readyReplicas":3,"nodes":["my-redis-xxxx-aaaa",...]}

测试扩缩容

spec.size 改成 5:

bash
kubectl patch rediscluster my-redis --type=merge -p '{"spec":{"size":5}}'

Operator 会感知到 spec 变化,把 Deployment 副本数扩到 5。观察:

bash
kubectl get deployment my-redis
# READY   UP-TO-DATE   AVAILABLE
# 5/5     5            5

测试自愈

手动删一个 Pod:

bash
kubectl delete pod -l app.kubernetes.io/instance=my-redis --field-selector=status.phase=Running | head -1

Deployment 控制器会自动补一个新的 Pod。这就是声明式 + 控制循环带来的自愈能力——Operator 不需要做任何事,Deployment 已经在帮我们维持期望状态了。

测试级联删除

删除 CR:

bash
kubectl delete rediscluster my-redis

由于 Deployment 和 Service 都设置了 ownerReferences 指向这个 CR,它们会被自动删除:

bash
kubectl get deployment,service,pods -l app.kubernetes.io/instance=my-redis
# 无输出,资源已被清理

九、常见问题排查

开发过程中你可能会遇到这些问题:

1. CRD 没装上

make install 报错「the server could not find the requested resource」,通常是 make manifests 没执行或 CRD 文件没生成。检查 config/crd/bases/ 下是否有 yaml 文件。

2. RBAC 权限不足

Operator 日志报「forbidden: User cannot create deployments」之类错误,是 RBAC 配置不全。检查 controller 文件里的 +kubebuilder:rbac 注释是否覆盖了所有操作的资源,然后重新 make manifestsmake install。本地 make run 时用的是 kubeconfig 里的用户,权限通常足够;部署到集群(第八章)时才会真正受 RBAC 限制。

3. Status 更新失败

报错「the body of the request was in an unknown format」,通常是 CRD 没开 status 子资源却用了 Status().Update(),或者反过来。确保 +kubebuilder:subresource:status 标记存在且重新生成了 CRD。

4. 无限 Reconcile

日志疯狂打印 Reconcile,通常是 Status 更新或 Update 引起了 resourceVersion 变化又触发 Watch。解决方法是更新前严格比较,确认确实变了才写。本篇代码的 updateStatus 已经做了这点。

十、小结

本篇从零实现了一个能用的 Redis Operator,走完了 Operator 开发的完整闭环。要点回顾:

  1. 定义 API:在 types.go 用 Go 结构体 + +kubebuilder 注释定义 Spec/Status,make manifests 生成 CRD。
  2. Reconcile 五步法:获取 CR → 调谐 Deployment → 调谐 Service → 更新 Status → 按需重新入队。这是所有 Operator 的通用骨架。
  3. 幂等是核心:每个资源先 Get 后 Create/Update,更新前比较差异,避免无意义写入。
  4. SetControllerReference 实现级联删除,让关联资源随 CR 一起被清理。
  5. Status 子资源:开了 +kubebuilder:subresource:status 后必须用 Status().Update(),且 status 变化不触发 Reconcile。
  6. RequeueAfter 让 Operator 在没有事件的情况下也能定时观察状态推进,是 Level-Triggered 的补充手段。
  7. 验证闭环:apply CR → 看 Operator 日志 → kubectl 查资源 → 测试扩缩容/自愈/级联删除。

下一篇我们会深入 Reconcile 循环的细节,包括返回值语义、Owner Reference 原理、CreateOrUpdate、Patch、Predicate、多资源 Watch 等,把本篇「能用」的代码提升到「健壮」的水平。