返回首页

8.1 Operator 模式与 CRD 开发(Go + kubebuilder)

本章目标

这一章进入 Kubernetes 进阶开发领域。学完之后,你应该能回答下面几个问题:

  • 为什么内置资源还不够,企业还会大量使用 CRD 和 Operator?
  • 一个自定义资源应该怎么设计字段、版本和校验规则?
  • kubebuilder 初始化项目之后,真正需要自己写的代码在哪?
  • Reconcile 循环到底在“调和”什么?
  • 如果你想从零学会写一个简单 Operator,推荐路径是什么?

1. 为什么 Kubernetes 需要 CRD 与 Operator

Kubernetes 内置了 Pod、Deployment、Service、StatefulSet 等一批通用资源,但真实生产环境里,团队往往还需要描述“数据库实例”“消息队列集群”“缓存集群”“备份任务”“应用发布策略”等业务语义。CRD(CustomResourceDefinition)允许我们把这些领域对象扩展进 Kubernetes API;Operator 则把“人肉运维经验”写成控制器逻辑,让集群持续把实际状态拉回到期望状态。Kubernetes 官方把 CRD 视为扩展 API 的标准方式,而控制器 / Reconciler 是让这些扩展真正自动化运行的核心机制::cite[254].

可以把它理解成:

  • CRD 解决“怎么表达业务对象”
  • Controller / Operator 解决“怎么自动把业务对象执行出来”

Kubebuilder 文档对控制器的描述非常直接:控制器的职责就是持续比较“期望状态”和“实际状态”,并执行调和动作;这正是 Operator 模式的底层本质::cite[195].

1.1 没有 Operator 时的问题

很多团队的早期做法是:

  1. 用 Helm 安装一堆 YAML
  2. 靠运维手册处理扩缩容、升级、备份、故障恢复
  3. 出问题时人工排查资源状态

这种模式的问题在于:

问题 表现
业务语义分散 “数据库实例”被拆成 StatefulSet、Secret、PVC、Service 等多个对象,认知成本高
运维动作不可复用 扩容、备份、恢复全靠手工执行,流程容易漂移
状态不闭环 某个依赖资源被删除后,系统不会自动修复
知识难沉淀 经验停留在文档和人脑里,无法持续执行

Operator 的价值,就是把“手册”变成“控制循环”。

1.2 什么场景特别适合 Operator

如果你的对象具有下面任意特征,就很适合用 Operator:

  • 有明显的生命周期管理:创建、升级、暂停、恢复、删除
  • 需要跨多个 Kubernetes 资源协同
  • 有状态系统较多,如 MySQL、Redis、Kafka、Elasticsearch
  • 需要自动化巡检、自愈、备份、回滚
  • 业务方希望通过一个高层对象直接声明目标状态

一个简单判断标准是:如果你已经写了很多脚本来维护同一种系统,那么这类脚本往往就值得演进成 Operator。


2. 自定义资源的定义方式与版本管理

CRD 的核心作用,是把一个新的 Kind 注册到 Kubernetes API Server 中。注册完成后,你就可以像使用 Deployment 一样使用自己的资源,比如 CacheDBClusterAppRelease 等::cite[254].

2.1 CRD 的基本结构

一个 CRD 通常需要定义这些关键部分:

  • group:API 组,例如 apps.mycompany.io
  • versions:版本列表,例如 v1alpha1v1beta1v1
  • scope:作用域,通常是 NamespacedCluster
  • names:资源名称、复数名、Kind、缩写
  • schema:OpenAPI v3 Schema,用来描述字段结构

下面是一个简化版 CRD 示例:

apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: caches.apps.example.com
spec:
  group: apps.example.com
  scope: Namespaced
  names:
    plural: caches
    singular: cache
    kind: Cache
    shortNames:
      - cc
  versions:
    - name: v1alpha1
      served: true
      storage: true
      schema:
        openAPIV3Schema:
          type: object
          properties:
            spec:
              type: object
              required:
                - image
                - replicas
              properties:
                image:
                  type: string
                replicas:
                  type: integer
                  minimum: 1
                  maximum: 10
                port:
                  type: integer
                  minimum: 1
                  maximum: 65535
            status:
              type: object
              properties:
                readyReplicas:
                  type: integer
                phase:
                  type: string
      subresources:
        status: {}

2.2 版本管理怎么设计

真实项目里,自定义资源几乎一定会演进,因此版本设计非常重要。一般遵循下面原则:

阶段 典型版本 适用场景
试验期 v1alpha1 字段可能频繁变动,适合快速验证
预发布 v1beta1 结构逐渐稳定,开始小范围推广
稳定期 v1 对外承诺兼容性,慎改字段语义

建议你在早期就明确:

  • 哪个版本是 storage: true(存储版本)
  • 哪些版本对外提供 served: true
  • 字段升级是否需要 conversion webhook

2.3 CEL 校验为什么值得重视

Kubernetes 支持使用 CEL(Common Expression Language)直接在 API Server 内做声明式校验。官方文档明确指出,CEL 可以用于验证规则、策略规则等场景,而且它是在 API Server 内直接执行,因此对很多扩展需求来说,可以替代一部分额外 webhook 的复杂度::cite[255].

CRD Validation Rules 从 Kubernetes 1.25 升到 Beta,而在 Kubernetes 1.29 达到 GA。也就是说,在 v1.32 里,基于 CEL 的 CRD 校验已经是成熟能力,可以放心纳入自定义资源设计基线::cite[316].

比如,我们要求:

  • minReplicas 不能大于 maxReplicas
  • 如果开启备份,必须提供备份计划

可以写成:

apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: appreleases.platform.example.com
spec:
  group: platform.example.com
  scope: Namespaced
  names:
    plural: appreleases
    singular: apprelease
    kind: AppRelease
  versions:
    - name: v1alpha1
      served: true
      storage: true
      schema:
        openAPIV3Schema:
          type: object
          properties:
            spec:
              type: object
              properties:
                minReplicas:
                  type: integer
                maxReplicas:
                  type: integer
                backup:
                  type: object
                  properties:
                    enabled:
                      type: boolean
                    schedule:
                      type: string
              x-kubernetes-validations:
                - rule: "self.minReplicas <= self.maxReplicas"
                  message: "minReplicas 不能大于 maxReplicas"
                - rule: "!has(self.backup) || !self.backup.enabled || has(self.backup.schedule)"
                  message: "开启备份后必须提供 schedule"

2.4 设计 CRD 时的经验建议

实践里,specstatus 一定要分清:

  • spec:用户声明的期望状态
  • status:控制器观测到的实际状态

你最好避免下面几类设计错误:

  1. 把运行时信息写进 spec
  2. 一个字段承担多种语义
  3. 直接暴露过多底层实现细节
  4. 在 alpha 阶段就承诺长期兼容

一个常用模板是:

spec:
  image / version / replicas / resources / config / policy
status:
  observedGeneration / conditions / phase / readyReplicas / endpoints

3. 使用 kubebuilder 初始化一个 Operator 项目

Kubebuilder 是官方生态里最主流的 Go Operator 开发脚手架之一。它基于 controller-runtime 和 controller-tools,提供项目结构、CRD 生成、RBAC 生成、Webhook、测试框架等能力。Quick Start 文档给出的典型流程就是:初始化项目、创建 API、实现控制器、安装 CRD、运行控制器并创建样例资源::cite[194].

3.1 环境准备

建议环境如下:

组件 建议版本
Go 1.22+
Kubebuilder 4.x
controller-runtime 与脚手架生成版本保持一致
Kubernetes 集群 v1.32
kind / minikube / 任意测试集群 均可

3.2 初始化项目

mkdir cache-operator && cd cache-operator
kubebuilder init --domain example.com --repo example.com/cache-operator
kubebuilder create api --group apps --version v1alpha1 --kind Cache

执行完之后,你会看到类似结构:

cache-operator/
├── api/v1alpha1/
├── cmd/
├── config/
├── internal/controller/
├── Makefile
├── go.mod
└── main.go

3.3 自定义资源类型:api/v1alpha1/cache_types.go

下面给出一个可以直接学习的完整示例。这个资源描述一个简单缓存服务,控制器会为它创建 Deployment 和 Service。

package v1alpha1

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

// CacheSpec 定义用户的期望状态
type CacheSpec struct {
	// +kubebuilder:validation:MinLength=1
	Image string `json:"image"`

	// +kubebuilder:validation:Minimum=1
	// +kubebuilder:validation:Maximum=10
	Replicas *int32 `json:"replicas,omitempty"`

	// +kubebuilder:validation:Minimum=1
	// +kubebuilder:validation:Maximum=65535
	// +kubebuilder:default=6379
	Port int32 `json:"port,omitempty"`
}

// CacheStatus 定义控制器观测到的状态
type CacheStatus struct {
	ReadyReplicas int32              `json:"readyReplicas,omitempty"`
	Phase         string             `json:"phase,omitempty"`
	Conditions    []metav1.Condition `json:"conditions,omitempty"`
}

// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
// +kubebuilder:resource:scope=Namespaced,shortName=cc
// +kubebuilder:printcolumn:name="Ready",type="integer",JSONPath=".status.readyReplicas"
// +kubebuilder:printcolumn:name="Phase",type="string",JSONPath=".status.phase"
// +kubebuilder:validation:XValidation:rule="self.spec.replicas == null || self.spec.replicas >= 1",message="replicas 必须大于等于 1"

type Cache struct {
	metav1.TypeMeta   `json:",inline"`
	metav1.ObjectMeta `json:"metadata,omitempty"`

	Spec   CacheSpec   `json:"spec,omitempty"`
	Status CacheStatus `json:"status,omitempty"`
}

// +kubebuilder:object:root=true

type CacheList struct {
	metav1.TypeMeta `json:",inline"`
	metav1.ListMeta `json:"metadata,omitempty"`
	Items           []Cache `json:"items"`
}

func init() {
	SchemeBuilder.Register(&Cache{}, &CacheList{})
}

3.4 控制器实现:internal/controller/cache_controller.go

控制器核心任务:

  1. 读取 Cache 对象
  2. 确保对应 Deployment 存在且副本数正确
  3. 确保对应 Service 存在
  4. 根据 Deployment 状态回写 status
package controller

import (
	context "context"
	fmt "fmt"

	appsv1 "k8s.io/api/apps/v1"
	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/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"

	appsv1alpha1 "example.com/cache-operator/api/v1alpha1"
)

// CacheReconciler 调和 Cache 资源
type CacheReconciler struct {
	client.Client
	Scheme *runtime.Scheme
}

// +kubebuilder:rbac:groups=apps.example.com,resources=caches,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=apps.example.com,resources=caches/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=apps.example.com,resources=caches/finalizers,verbs=update
// +kubebuilder:rbac:groups=apps,resources=deployments,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups="",resources=services,verbs=get;list;watch;create;update;patch;delete

func (r *CacheReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
	logger := log.FromContext(ctx)

	var cache appsv1alpha1.Cache
	if err := r.Get(ctx, req.NamespacedName, &cache); err != nil {
		if apierrors.IsNotFound(err) {
			return ctrl.Result{}, nil
		}
		return ctrl.Result{}, err
	}

	replicas := int32(1)
	if cache.Spec.Replicas != nil {
		replicas = *cache.Spec.Replicas
	}

	deployment := r.desiredDeployment(&cache, replicas)
	if err := controllerutil.SetControllerReference(&cache, deployment, r.Scheme); err != nil {
		return ctrl.Result{}, err
	}

	var currentDeploy appsv1.Deployment
	err := r.Get(ctx, types.NamespacedName{Name: deployment.Name, Namespace: deployment.Namespace}, &currentDeploy)
	if err != nil && apierrors.IsNotFound(err) {
		logger.Info("creating deployment", "name", deployment.Name)
		if err := r.Create(ctx, deployment); err != nil {
			return ctrl.Result{}, err
		}
	} else if err != nil {
		return ctrl.Result{}, err
	} else {
		updated := false
		if *currentDeploy.Spec.Replicas != replicas {
			currentDeploy.Spec.Replicas = &replicas
			updated = true
		}
		if len(currentDeploy.Spec.Template.Spec.Containers) > 0 && currentDeploy.Spec.Template.Spec.Containers[0].Image != cache.Spec.Image {
			currentDeploy.Spec.Template.Spec.Containers[0].Image = cache.Spec.Image
			updated = true
		}
		if updated {
			logger.Info("updating deployment", "name", currentDeploy.Name)
			if err := r.Update(ctx, &currentDeploy); err != nil {
				return ctrl.Result{}, err
			}
		}
	}

	service := r.desiredService(&cache)
	if err := controllerutil.SetControllerReference(&cache, service, r.Scheme); err != nil {
		return ctrl.Result{}, err
	}

	var currentSvc corev1.Service
	err = r.Get(ctx, types.NamespacedName{Name: service.Name, Namespace: service.Namespace}, &currentSvc)
	if err != nil && apierrors.IsNotFound(err) {
		logger.Info("creating service", "name", service.Name)
		if err := r.Create(ctx, service); err != nil {
			return ctrl.Result{}, err
		}
	} else if err != nil {
		return ctrl.Result{}, err
	}

	var latestDeploy appsv1.Deployment
	if err := r.Get(ctx, types.NamespacedName{Name: deployment.Name, Namespace: deployment.Namespace}, &latestDeploy); err != nil {
		return ctrl.Result{}, err
	}

	cache.Status.ReadyReplicas = latestDeploy.Status.ReadyReplicas
	if latestDeploy.Status.ReadyReplicas == replicas {
		cache.Status.Phase = "Ready"
	} else {
		cache.Status.Phase = "Progressing"
	}

	metaCondition := metav1.Condition{
		Type:               "Available",
		Status:             metav1.ConditionTrue,
		Reason:             cache.Status.Phase,
		Message:            fmt.Sprintf("ready replicas: %d", latestDeploy.Status.ReadyReplicas),
		ObservedGeneration: cache.Generation,
		LastTransitionTime: metav1.Now(),
	}
	cache.Status.Conditions = []metav1.Condition{metaCondition}

	if err := r.Status().Update(ctx, &cache); err != nil {
		return ctrl.Result{}, err
	}

	return ctrl.Result{}, nil
}

func (r *CacheReconciler) desiredDeployment(cache *appsv1alpha1.Cache, replicas int32) *appsv1.Deployment {
	labels := map[string]string{
		"app":   "cache",
		"cache": cache.Name,
	}

	return &appsv1.Deployment{
		ObjectMeta: metav1.ObjectMeta{
			Name:      cache.Name,
			Namespace: cache.Namespace,
		},
		Spec: appsv1.DeploymentSpec{
			Replicas: &replicas,
			Selector: &metav1.LabelSelector{MatchLabels: labels},
			Template: corev1.PodTemplateSpec{
				ObjectMeta: metav1.ObjectMeta{Labels: labels},
				Spec: corev1.PodSpec{
					Containers: []corev1.Container{
						{
							Name:  "cache",
							Image: cache.Spec.Image,
							Ports: []corev1.ContainerPort{
								{
									ContainerPort: cache.Spec.Port,
									Name:          "tcp",
								},
							},
						},
					},
				},
			},
		},
	}
}

func (r *CacheReconciler) desiredService(cache *appsv1alpha1.Cache) *corev1.Service {
	labels := map[string]string{
		"app":   "cache",
		"cache": cache.Name,
	}

	return &corev1.Service{
		ObjectMeta: metav1.ObjectMeta{
			Name:      cache.Name,
			Namespace: cache.Namespace,
		},
		Spec: corev1.ServiceSpec{
			Selector: labels,
			Ports: []corev1.ServicePort{
				{
					Name:       "tcp",
					Port:       cache.Spec.Port,
					TargetPort: intstrFromInt32(cache.Spec.Port),
				},
			},
		},
	}
}

func intstrFromInt32(v int32) intstr.IntOrString {
	return intstr.FromInt32(v)
}

func (r *CacheReconciler) SetupWithManager(mgr ctrl.Manager) error {
	return ctrl.NewControllerManagedBy(mgr).
		For(&appsv1alpha1.Cache{}).
		Owns(&appsv1.Deployment{}).
		Owns(&corev1.Service{}).
		Complete(r)
}

说明:上面代码用于教学演示,重点是理解 Reconcile 模式、owner reference、状态回写和资源对齐。实际项目里,你还应补充错误分类、重试节奏、条件状态细化、finalizer、事件记录、Webhook 校验和测试用例。

3.5 样例资源:config/samples/apps_v1alpha1_cache.yaml

apiVersion: apps.example.com/v1alpha1
kind: Cache
metadata:
  name: demo-cache
spec:
  image: redis:7.2
  replicas: 2
  port: 6379

3.6 生成 CRD 与运行控制器

make manifests
make generate
make install
make run
kubectl apply -f config/samples/apps_v1alpha1_cache.yaml
kubectl get cache
kubectl get deploy,svc

你会观察到:

  • Cache 对象被创建
  • 控制器自动创建同名 Deployment 与 Service
  • Deployment Ready 后,status.readyReplicas 会被更新

这就是一个最小可运行 Operator 的基本闭环。


4. Reconcile 循环与控制器开发基本模式

Kubebuilder 文档强调:控制器的本质就是 Reconciler,不停地把世界修正到资源声明的目标状态。这种模式并不是“一次性脚本”,而是一个持续收敛过程::cite[195].

4.1 Reconcile 的典型流程

可以把 Reconcile 理解为下面的伪代码:

读取当前对象
  -> 校验是否被删除 / 是否需要 finalizer
  -> 获取依赖资源当前状态
  -> 计算目标状态
  -> 创建 / 更新 / 删除依赖资源
  -> 更新 status
  -> 返回是否需要重试或延迟重试

在真实项目里,常见模式基本都围绕这几步展开。

4.2 写控制器时最重要的 6 个原则

原则 1:Reconcile 必须幂等

同一个对象可能会被反复调和很多次,因此你不能假设代码只执行一次。最常见写法是:

  • Get
  • 不存在就 Create
  • 存在就比较关键字段后 Update

原则 2:spec 驱动,status 反馈

不要把业务状态机全塞进 annotations 或 labels;status.conditionsphaseobservedGeneration 才是更规范的反馈渠道。

原则 3:只管理自己负责的资源

通过 owner reference 绑定子资源,并在 SetupWithManager 里声明 Owns(...),这样事件链路更清晰,也方便垃圾回收。

原则 4:删除流程要走 finalizer

如果你的 Operator 创建了外部资源(云盘、DNS、数据库账号、对象存储桶),删除时必须先清理外部资源,再移除 finalizer,否则会留下“集群里删掉了,外部还在”的脏状态。

原则 5:显式处理错误类型

  • 临时错误:返回 error,让 controller-runtime 重试
  • 需要延后观察:返回 ctrl.Result{RequeueAfter: ...}
  • 对象不存在:通常直接结束

原则 6:让状态可观测

至少补齐这些输出:

  • status.conditions
  • 关键事件 Event
  • 结构化日志
  • 指标(调和次数、失败率、耗时)

4.3 一个常见的 Reconcile 骨架

func (r *AppReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
    var obj apiv1.App
    if err := r.Get(ctx, req.NamespacedName, &obj); err != nil {
        return ctrl.Result{}, client.IgnoreNotFound(err)
    }

    if !obj.DeletionTimestamp.IsZero() {
        return r.reconcileDelete(ctx, &obj)
    }

    if err := r.ensureFinalizer(ctx, &obj); err != nil {
        return ctrl.Result{}, err
    }

    desired := r.buildDesiredState(&obj)
    if err := r.applyDesiredState(ctx, &obj, desired); err != nil {
        return ctrl.Result{}, err
    }

    if err := r.updateStatus(ctx, &obj); err != nil {
        return ctrl.Result{}, err
    }

    return ctrl.Result{}, nil
}

这套骨架非常常见。你以后写数据库 Operator、发布 Operator、配置同步 Operator,思路都大同小异。


5. 一个简单 Operator 的实战学习路线

如果你是第一次做 Operator,不建议一上来就写复杂数据库系统。更好的方式是逐步升级:

第 1 阶段:只学 CRD

目标:先会定义资源,不急着写控制器。

练习建议:

  • 定义 WebApp / Cache / JobPolicy 这类资源
  • 熟悉 spec / status / printcolumn
  • 用 CEL 做字段关系约束

第 2 阶段:单资源控制器

目标:一个 CR 只控制一个 Deployment。

练习建议:

  • Cache -> Deployment
  • 副本数和镜像字段能自动同步
  • 状态能回写 readyReplicas

第 3 阶段:多资源协同

目标:一个 CR 控制 Deployment + Service + ConfigMap + Secret。

练习建议:

  • 配置变更触发滚动更新
  • Service 自动暴露端口
  • Secret 缺失时给出明确状态

第 4 阶段:生命周期管理

目标:加入删除处理、暂停、升级、回滚。

练习建议:

  • finalizer 清理外部资源
  • spec.paused 暂停调和
  • status.conditions 表达升级过程

第 5 阶段:生产级增强

目标:补齐可靠性和工程化能力。

练习建议:

  • webhook 默认值与高级校验
  • 指标、日志、事件
  • envtest / 集成测试
  • leader election、HA 部署、RBAC 最小权限

5.1 学习时的避坑清单

常见误区 更好的做法
一开始就设计大而全 CRD 先做最小字段集,逐步演进
把业务逻辑全写进 Reconcile 拆分为 build / ensure / updateStatus 等小函数
用 if-else 直接堆状态机 引入条件、阶段、显式 helper 函数
不写 status 让用户只能靠 kubectl describe 猜测状态
不做删除清理 使用 finalizer 管理外部资源

5.2 推荐的最小实战项目

如果你想用 1~2 周真正入门,可以按这个顺序练:

  1. 写一个 Cache Operator,自动创建 Deployment + Service
  2. 增加 ConfigMap 挂载能力
  3. 增加 finalizer,在删除时清理外部注册记录
  4. 增加 webhook 默认值和 CEL 校验
  5. 加入 status.conditions 和 Prometheus 指标

做到这里,你就已经不是“会看 Operator 代码”,而是“能独立做一个简单 Operator”。


6. 本章小结

这一章最重要的不是记住命令,而是建立一套心智模型:

  • CRD 用来表达新的领域对象
  • Operator 用控制循环持续执行运维逻辑
  • Kubebuilder 帮你快速搭建工程骨架
  • Reconcile 是整个模式的核心
  • CEL 校验 让很多规则可以直接在 API Server 里声明化完成::cite[255]

如果你后面真的要开始做生产级 Operator,建议先把“最小闭环”跑通:

一个 CR → 一个控制器 → 两三个子资源 → 一套明确的状态字段。

先跑通,再扩展,这是最稳的学习方式。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

7.3 Tracing:OpenTelemetry

下一篇

9.1 排查方法论与工具