Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

Kubernetes Operators

A comprehensive guide to building and managing Kubernetes Operators using the operator pattern to automate application lifecycle management.

Kubernetes Operators

A comprehensive guide to building and managing Kubernetes Operators using the operator pattern to automate application lifecycle management.

Overview

Kubernetes Operators are software extensions that use custom resources to manage applications and their components. Operators follow Kubernetes principles, notably the control loop, to extend the Kubernetes API with domain-specific knowledge. They encode operational knowledge about managing complex, stateful applications by watching custom resources and continuously reconciling the desired state with the actual state.

Operator Pattern ArchitectureControl LoopCreates/UpdatesWatched byReadsTriggersReconcileUpdateCreate/Update/DeleteEventsManagesManagesManagesManagesFeedbackController/OperatorUser/AdminCustom ResourceInformer CacheWork QueueReconciliation LogicResource StatusKubernetes ResourcesDeploymentServiceConfigMapSecretOperator Pattern ArchitectureControl LoopCreates/UpdatesWatched byReadsTriggersReconcileUpdateCreate/Update/DeleteEventsManagesManagesManagesManagesFeedbackController/OperatorUser/AdminCustom ResourceInformer CacheWork QueueReconciliation LogicResource StatusKubernetes ResourcesDeploymentServiceConfigMapSecret

Key Concepts

Operator Pattern: Combines custom resources (data) with custom controllers (logic) to automate the entire lifecycle of complex applications.

Reconciliation: The process of observing the current state, comparing it to the desired state, and taking action to move current state towards desired state.

Level-triggered: Controllers react to the current state rather than individual events, making the system self-healing and resilient.

Controller-runtime: A set of libraries that provide the foundation for building operators with minimal boilerplate.

Custom Resource Definitions (CRDs)

CRDs define the schema for custom resources that operators manage.

Basic CRD Structure

apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: databases.db.example.com
spec:
  group: db.example.com
  versions:
    - name: v1alpha1
      served: true
      storage: true
      schema:
        openAPIV3Schema:
          type: object
          properties:
            spec:
              type: object
              required: ["type", "size"]
              properties:
                type:
                  type: string
                  enum: ["postgresql", "mysql", "mongodb"]
                size:
                  type: string
                  enum: ["small", "medium", "large"]
                version:
                  type: string
                  pattern: '^\d+\.\d+$'
                replicas:
                  type: integer
                  minimum: 1
                  maximum: 10
                  default: 3
                backup:
                  type: object
                  properties:
                    enabled:
                      type: boolean
                      default: true
                    schedule:
                      type: string
                      pattern: '^(\*|([0-9]|1[0-9]|2[0-9]|3[0-9]|4[0-9]|5[0-9])|\*\/([0-9]|1[0-9]|2[0-9]|3[0-9]|4[0-9]|5[0-9])) (\*|([0-9]|1[0-9]|2[0-3])|\*\/([0-9]|1[0-9]|2[0-3])) (\*|([1-9]|1[0-9]|2[0-9]|3[0-1])|\*\/([1-9]|1[0-9]|2[0-9]|3[0-1])) (\*|([1-9]|1[0-2])|\*\/([1-9]|1[0-2])) (\*|([0-6])|\*\/([0-6]))$'
                    retention:
                      type: integer
                      minimum: 1
                      maximum: 365
                      default: 7
                storage:
                  type: object
                  properties:
                    size:
                      type: string
                      pattern: '^\d+[MGT]i$'
                    storageClass:
                      type: string
            status:
              type: object
              properties:
                phase:
                  type: string
                  enum: ["Pending", "Creating", "Running", "Failed", "Updating"]
                ready:
                  type: boolean
                endpoints:
                  type: object
                  properties:
                    primary:
                      type: string
                    replicas:
                      type: array
                      items:
                        type: string
                lastBackup:
                  type: string
                  format: date-time
                message:
                  type: string
      subresources:
        status: {}
        scale:
          specReplicasPath: .spec.replicas
          statusReplicasPath: .status.replicas
      additionalPrinterColumns:
        - name: Type
          type: string
          jsonPath: .spec.type
        - name: Size
          type: string
          jsonPath: .spec.size
        - name: Phase
          type: string
          jsonPath: .status.phase
        - name: Ready
          type: boolean
          jsonPath: .status.ready
        - name: Age
          type: date
          jsonPath: .metadata.creationTimestamp
  scope: Namespaced
  names:
    plural: databases
    singular: database
    kind: Database
    shortNames:
      - db

CRD with Validation and Defaults

apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: webservices.app.example.com
spec:
  group: app.example.com
  versions:
    - name: v1
      served: true
      storage: true
      schema:
        openAPIV3Schema:
          type: object
          required: ["spec"]
          properties:
            spec:
              type: object
              required: ["image"]
              properties:
                image:
                  type: string
                  minLength: 1
                replicas:
                  type: integer
                  minimum: 1
                  maximum: 100
                  default: 3
                resources:
                  type: object
                  properties:
                    limits:
                      type: object
                      properties:
                        cpu:
                          type: string
                          pattern: '^\d+m?$'
                        memory:
                          type: string
                          pattern: '^\d+[MGT]i?$'
                    requests:
                      type: object
                      properties:
                        cpu:
                          type: string
                          pattern: '^\d+m?$'
                        memory:
                          type: string
                          pattern: '^\d+[MGT]i?$'
                ingress:
                  type: object
                  properties:
                    enabled:
                      type: boolean
                      default: false
                    host:
                      type: string
                      pattern: '^[a-z0-9]([-a-z0-9]*[a-z0-9])?(\.[a-z0-9]([-a-z0-9]*[a-z0-9])?)*$'
                    tls:
                      type: boolean
                      default: true
                    annotations:
                      type: object
                      x-kubernetes-preserve-unknown-fields: true
            status:
              type: object
              properties:
                availableReplicas:
                  type: integer
                conditions:
                  type: array
                  items:
                    type: object
                    properties:
                      type:
                        type: string
                      status:
                        type: string
                        enum: ["True", "False", "Unknown"]
                      lastTransitionTime:
                        type: string
                        format: date-time
                      reason:
                        type: string
                      message:
                        type: string
      subresources:
        status: {}
        scale:
          specReplicasPath: .spec.replicas
          statusReplicasPath: .status.availableReplicas
  scope: Namespaced
  names:
    plural: webservices
    singular: webservice
    kind: WebService
    shortNames:
      - ws

Managing CRDs

# Apply CRD
kubectl apply -f database-crd.yaml

# List CRDs
kubectl get crds

# Describe CRD
kubectl describe crd databases.db.example.com

# Get CRD details
kubectl get crd databases.db.example.com -o yaml

# Delete CRD (this deletes all custom resources of this type!)
kubectl delete crd databases.db.example.com

# List all custom resources of a type
kubectl get databases -A

# Create a custom resource
kubectl apply -f my-database.yaml

# Get custom resource
kubectl get database my-postgres -o yaml

# Watch custom resources
kubectl get databases -w

Custom Resource Example

apiVersion: db.example.com/v1alpha1
kind: Database
metadata:
  name: production-postgres
  namespace: production
spec:
  type: postgresql
  size: large
  version: "15.2"
  replicas: 3
  backup:
    enabled: true
    schedule: "0 2 * * *"  # Daily at 2 AM
    retention: 30
  storage:
    size: 100Gi
    storageClass: fast-ssd

Controller Logic and Reconciliation Loops

The reconciliation loop is the heart of an operator, continuously watching for changes and ensuring desired state matches actual state.

Reconciliation Loop Flow

NoYesYesNoYesNoYesNoYesNoYesNoStartWatch CustomResourceEvent Received?WaitFetch Current StateCurrent == Desired?Update Status:SuccessDetermine ActionsNeed Create?Create ResourcesNeed Update?Update ResourcesNeed Delete?Delete ResourcesUpdate StatusError?Requeue with BackoffNoYesYesNoYesNoYesNoYesNoYesNoStartWatch CustomResourceEvent Received?WaitFetch Current StateCurrent == Desired?Update Status:SuccessDetermine ActionsNeed Create?Create ResourcesNeed Update?Update ResourcesNeed Delete?Delete ResourcesUpdate StatusError?Requeue with Backoff

Reconciliation Principles

Idempotency: Running reconciliation multiple times produces the same result.

Level-triggered: React to current state, not individual events.

Error Handling: Always requeue with exponential backoff on errors.

Status Updates: Separate status updates from spec reconciliation.

Finalizers: Clean up external resources before deletion.

Operator Frameworks

Operator SDK (Go)

The Operator SDK provides tools to build, test, and package operators following best practices.

Complete Go Operator Example

// api/v1alpha1/database_types.go
package v1alpha1

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

// DatabaseSpec defines the desired state of Database
type DatabaseSpec struct {
	// Type of database (postgresql, mysql, mongodb)
	// +kubebuilder:validation:Enum=postgresql;mysql;mongodb
	Type string `json:"type"`

	// Size of the database instance
	// +kubebuilder:validation:Enum=small;medium;large
	Size string `json:"size"`

	// Version of the database
	// +kubebuilder:validation:Pattern=`^\d+\.\d+$`
	Version string `json:"version,omitempty"`

	// Number of replicas
	// +kubebuilder:validation:Minimum=1
	// +kubebuilder:validation:Maximum=10
	// +kubebuilder:default=3
	Replicas int32 `json:"replicas,omitempty"`

	// Backup configuration
	Backup *BackupConfig `json:"backup,omitempty"`

	// Storage configuration
	Storage *StorageConfig `json:"storage,omitempty"`
}

type BackupConfig struct {
	// Enable backups
	// +kubebuilder:default=true
	Enabled bool `json:"enabled,omitempty"`

	// Backup schedule in cron format
	Schedule string `json:"schedule,omitempty"`

	// Retention period in days
	// +kubebuilder:validation:Minimum=1
	// +kubebuilder:validation:Maximum=365
	// +kubebuilder:default=7
	Retention int32 `json:"retention,omitempty"`
}

type StorageConfig struct {
	// Size of persistent volume
	Size string `json:"size,omitempty"`

	// StorageClass to use
	StorageClass string `json:"storageClass,omitempty"`
}

// DatabaseStatus defines the observed state of Database
type DatabaseStatus struct {
	// Phase of the database
	// +kubebuilder:validation:Enum=Pending;Creating;Running;Failed;Updating
	Phase string `json:"phase,omitempty"`

	// Ready indicates if the database is ready
	Ready bool `json:"ready,omitempty"`

	// Endpoints for connecting to the database
	Endpoints *DatabaseEndpoints `json:"endpoints,omitempty"`

	// LastBackup timestamp
	LastBackup *metav1.Time `json:"lastBackup,omitempty"`

	// Human-readable message
	Message string `json:"message,omitempty"`

	// Conditions represent the latest observations of the object's state
	Conditions []metav1.Condition `json:"conditions,omitempty"`
}

type DatabaseEndpoints struct {
	Primary  string   `json:"primary,omitempty"`
	Replicas []string `json:"replicas,omitempty"`
}

// +kubebuilder:object:root=true
// +kubebuilder:subresource:status
// +kubebuilder:subresource:scale:specpath=.spec.replicas,statuspath=.status.replicas
// +kubebuilder:printcolumn:name="Type",type=string,JSONPath=`.spec.type`
// +kubebuilder:printcolumn:name="Size",type=string,JSONPath=`.spec.size`
// +kubebuilder:printcolumn:name="Phase",type=string,JSONPath=`.status.phase`
// +kubebuilder:printcolumn:name="Ready",type=boolean,JSONPath=`.status.ready`
// +kubebuilder:printcolumn:name="Age",type=date,JSONPath=`.metadata.creationTimestamp`

// Database is the Schema for the databases API
type Database struct {
	metav1.TypeMeta   `json:",inline"`
	metav1.ObjectMeta `json:"metadata,omitempty"`

	Spec   DatabaseSpec   `json:"spec,omitempty"`
	Status DatabaseStatus `json:"status,omitempty"`
}

// +kubebuilder:object:root=true

// DatabaseList contains a list of Database
type DatabaseList struct {
	metav1.TypeMeta `json:",inline"`
	metav1.ListMeta `json:"metadata,omitempty"`
	Items           []Database `json:"items"`
}

func init() {
	SchemeBuilder.Register(&Database{}, &DatabaseList{})
}
// controllers/database_controller.go
package controllers

import (
	"context"
	"fmt"
	"time"

	appsv1 "k8s.io/api/apps/v1"
	corev1 "k8s.io/api/core/v1"
	"k8s.io/apimachinery/pkg/api/errors"
	"k8s.io/apimachinery/pkg/api/meta"
	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"

	dbv1alpha1 "example.com/database-operator/api/v1alpha1"
)

const (
	finalizerName = "db.example.com/finalizer"

	// Condition types
	TypeAvailable = "Available"
	TypeProgressing = "Progressing"
	TypeDegraded = "Degraded"
)

// DatabaseReconciler reconciles a Database object
type DatabaseReconciler struct {
	client.Client
	Scheme *runtime.Scheme
}

// +kubebuilder:rbac:groups=db.example.com,resources=databases,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=db.example.com,resources=databases/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=db.example.com,resources=databases/finalizers,verbs=update
// +kubebuilder:rbac:groups=apps,resources=statefulsets,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=configmaps,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=core,resources=secrets,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=core,resources=persistentvolumeclaims,verbs=get;list;watch;create;update;patch;delete

// Reconcile implements the reconciliation loop
func (r *DatabaseReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
	logger := log.FromContext(ctx)
	logger.Info("Reconciling Database", "namespace", req.Namespace, "name", req.Name)

	// Fetch the Database instance
	database := &dbv1alpha1.Database{}
	if err := r.Get(ctx, req.NamespacedName, database); err != nil {
		if errors.IsNotFound(err) {
			// Resource deleted, nothing to do
			logger.Info("Database resource not found, likely deleted")
			return ctrl.Result{}, nil
		}
		logger.Error(err, "Failed to get Database")
		return ctrl.Result{}, err
	}

	// Handle deletion with finalizers
	if database.ObjectMeta.DeletionTimestamp.IsZero() {
		// Not being deleted, add finalizer if needed
		if !controllerutil.ContainsFinalizer(database, finalizerName) {
			controllerutil.AddFinalizer(database, finalizerName)
			if err := r.Update(ctx, database); err != nil {
				return ctrl.Result{}, err
			}
		}
	} else {
		// Being deleted
		if controllerutil.ContainsFinalizer(database, finalizerName) {
			// Run finalization logic
			if err := r.finalizeDatabase(ctx, database); err != nil {
				return ctrl.Result{}, err
			}

			// Remove finalizer
			controllerutil.RemoveFinalizer(database, finalizerName)
			if err := r.Update(ctx, database); err != nil {
				return ctrl.Result{}, err
			}
		}
		return ctrl.Result{}, nil
	}

	// Update phase to Creating if Pending
	if database.Status.Phase == "" || database.Status.Phase == "Pending" {
		database.Status.Phase = "Creating"
		if err := r.Status().Update(ctx, database); err != nil {
			logger.Error(err, "Failed to update Database status to Creating")
			return ctrl.Result{}, err
		}
		// Requeue immediately to continue reconciliation
		return ctrl.Result{Requeue: true}, nil
	}

	// Reconcile Secret for credentials
	secret := r.buildSecret(database)
	if err := r.reconcileResource(ctx, database, secret); err != nil {
		logger.Error(err, "Failed to reconcile Secret")
		r.updateCondition(database, TypeDegraded, metav1.ConditionTrue, "SecretFailed", err.Error())
		return ctrl.Result{}, err
	}

	// Reconcile ConfigMap
	configMap := r.buildConfigMap(database)
	if err := r.reconcileResource(ctx, database, configMap); err != nil {
		logger.Error(err, "Failed to reconcile ConfigMap")
		r.updateCondition(database, TypeDegraded, metav1.ConditionTrue, "ConfigMapFailed", err.Error())
		return ctrl.Result{}, err
	}

	// Reconcile Service
	service := r.buildService(database)
	if err := r.reconcileResource(ctx, database, service); err != nil {
		logger.Error(err, "Failed to reconcile Service")
		r.updateCondition(database, TypeDegraded, metav1.ConditionTrue, "ServiceFailed", err.Error())
		return ctrl.Result{}, err
	}

	// Reconcile StatefulSet
	statefulSet := r.buildStatefulSet(database)
	if err := r.reconcileResource(ctx, database, statefulSet); err != nil {
		logger.Error(err, "Failed to reconcile StatefulSet")
		r.updateCondition(database, TypeDegraded, metav1.ConditionTrue, "StatefulSetFailed", err.Error())
		return ctrl.Result{}, err
	}

	// Check StatefulSet status
	existingSts := &appsv1.StatefulSet{}
	if err := r.Get(ctx, types.NamespacedName{Name: statefulSet.Name, Namespace: statefulSet.Namespace}, existingSts); err != nil {
		return ctrl.Result{}, err
	}

	// Update Database status based on StatefulSet
	ready := existingSts.Status.ReadyReplicas == database.Spec.Replicas
	database.Status.Ready = ready

	if ready {
		database.Status.Phase = "Running"
		database.Status.Message = fmt.Sprintf("Database is running with %d/%d replicas ready",
			existingSts.Status.ReadyReplicas, database.Spec.Replicas)
		r.updateCondition(database, TypeAvailable, metav1.ConditionTrue, "DatabaseReady", "Database is fully operational")
		r.updateCondition(database, TypeProgressing, metav1.ConditionFalse, "ReconcileComplete", "Reconciliation completed successfully")
		r.updateCondition(database, TypeDegraded, metav1.ConditionFalse, "NoIssues", "No degradation detected")
	} else {
		database.Status.Phase = "Updating"
		database.Status.Message = fmt.Sprintf("Database is updating: %d/%d replicas ready",
			existingSts.Status.ReadyReplicas, database.Spec.Replicas)
		r.updateCondition(database, TypeProgressing, metav1.ConditionTrue, "Updating", "StatefulSet is being updated")
		r.updateCondition(database, TypeAvailable, metav1.ConditionFalse, "NotReady", "Not all replicas are ready")
	}

	// Update endpoints
	database.Status.Endpoints = &dbv1alpha1.DatabaseEndpoints{
		Primary: fmt.Sprintf("%s-0.%s.%s.svc.cluster.local:5432",
			database.Name, database.Name, database.Namespace),
		Replicas: make([]string, 0),
	}

	for i := int32(1); i < database.Spec.Replicas; i++ {
		replica := fmt.Sprintf("%s-%d.%s.%s.svc.cluster.local:5432",
			database.Name, i, database.Name, database.Namespace)
		database.Status.Endpoints.Replicas = append(database.Status.Endpoints.Replicas, replica)
	}

	// Update status
	if err := r.Status().Update(ctx, database); err != nil {
		logger.Error(err, "Failed to update Database status")
		return ctrl.Result{}, err
	}

	logger.Info("Reconciliation complete", "phase", database.Status.Phase, "ready", database.Status.Ready)

	// Requeue after 30 seconds to check health
	return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}

// reconcileResource creates or updates a Kubernetes resource
func (r *DatabaseReconciler) reconcileResource(ctx context.Context, owner *dbv1alpha1.Database, obj client.Object) error {
	logger := log.FromContext(ctx)

	// Set owner reference
	if err := controllerutil.SetControllerReference(owner, obj, r.Scheme); err != nil {
		return err
	}

	// Check if resource exists
	found := obj.DeepCopyObject().(client.Object)
	err := r.Get(ctx, types.NamespacedName{Name: obj.GetName(), Namespace: obj.GetNamespace()}, found)

	if err != nil && errors.IsNotFound(err) {
		// Create the resource
		logger.Info("Creating resource", "kind", obj.GetObjectKind().GroupVersionKind().Kind, "name", obj.GetName())
		return r.Create(ctx, obj)
	} else if err != nil {
		return err
	}

	// Update the resource
	logger.Info("Updating resource", "kind", obj.GetObjectKind().GroupVersionKind().Kind, "name", obj.GetName())
	return r.Update(ctx, obj)
}

// buildSecret creates a Secret for database credentials
func (r *DatabaseReconciler) buildSecret(db *dbv1alpha1.Database) *corev1.Secret {
	return &corev1.Secret{
		ObjectMeta: metav1.ObjectMeta{
			Name:      db.Name + "-credentials",
			Namespace: db.Namespace,
		},
		StringData: map[string]string{
			"username": "dbuser",
			"password": generatePassword(), // In production, use secure password generation
			"database": db.Name,
		},
	}
}

// buildConfigMap creates a ConfigMap for database configuration
func (r *DatabaseReconciler) buildConfigMap(db *dbv1alpha1.Database) *corev1.ConfigMap {
	config := map[string]string{
		"max_connections":    "100",
		"shared_buffers":     "256MB",
		"effective_cache_size": "1GB",
	}

	// Adjust config based on size
	switch db.Spec.Size {
	case "large":
		config["max_connections"] = "500"
		config["shared_buffers"] = "2GB"
		config["effective_cache_size"] = "8GB"
	case "medium":
		config["max_connections"] = "200"
		config["shared_buffers"] = "512MB"
		config["effective_cache_size"] = "2GB"
	}

	return &corev1.ConfigMap{
		ObjectMeta: metav1.ObjectMeta{
			Name:      db.Name + "-config",
			Namespace: db.Namespace,
		},
		Data: config,
	}
}

// buildService creates a Service for the database
func (r *DatabaseReconciler) buildService(db *dbv1alpha1.Database) *corev1.Service {
	return &corev1.Service{
		ObjectMeta: metav1.ObjectMeta{
			Name:      db.Name,
			Namespace: db.Namespace,
		},
		Spec: corev1.ServiceSpec{
			Selector: map[string]string{
				"app": db.Name,
			},
			ClusterIP: "None", // Headless service for StatefulSet
			Ports: []corev1.ServicePort{
				{
					Name:       "postgres",
					Port:       5432,
					TargetPort: intstr.FromInt(5432),
				},
			},
		},
	}
}

// buildStatefulSet creates a StatefulSet for the database
func (r *DatabaseReconciler) buildStatefulSet(db *dbv1alpha1.Database) *appsv1.StatefulSet {
	replicas := db.Spec.Replicas
	image := fmt.Sprintf("postgres:%s", db.Spec.Version)
	if db.Spec.Version == "" {
		image = "postgres:15"
	}

	return &appsv1.StatefulSet{
		ObjectMeta: metav1.ObjectMeta{
			Name:      db.Name,
			Namespace: db.Namespace,
		},
		Spec: appsv1.StatefulSetSpec{
			ServiceName: db.Name,
			Replicas:    &replicas,
			Selector: &metav1.LabelSelector{
				MatchLabels: map[string]string{
					"app": db.Name,
				},
			},
			Template: corev1.PodTemplateSpec{
				ObjectMeta: metav1.ObjectMeta{
					Labels: map[string]string{
						"app": db.Name,
					},
				},
				Spec: corev1.PodSpec{
					Containers: []corev1.Container{
						{
							Name:  "postgres",
							Image: image,
							Ports: []corev1.ContainerPort{
								{
									Name:          "postgres",
									ContainerPort: 5432,
								},
							},
							Env: []corev1.EnvVar{
								{
									Name: "POSTGRES_USER",
									ValueFrom: &corev1.EnvVarSource{
										SecretKeyRef: &corev1.SecretKeySelector{
											LocalObjectReference: corev1.LocalObjectReference{
												Name: db.Name + "-credentials",
											},
											Key: "username",
										},
									},
								},
								{
									Name: "POSTGRES_PASSWORD",
									ValueFrom: &corev1.EnvVarSource{
										SecretKeyRef: &corev1.SecretKeySelector{
											LocalObjectReference: corev1.LocalObjectReference{
												Name: db.Name + "-credentials",
											},
											Key: "password",
										},
									},
								},
								{
									Name: "POSTGRES_DB",
									ValueFrom: &corev1.EnvVarSource{
										SecretKeyRef: &corev1.SecretKeySelector{
											LocalObjectReference: corev1.LocalObjectReference{
												Name: db.Name + "-credentials",
											},
											Key: "database",
										},
									},
								},
							},
							VolumeMounts: []corev1.VolumeMount{
								{
									Name:      "data",
									MountPath: "/var/lib/postgresql/data",
								},
								{
									Name:      "config",
									MountPath: "/etc/postgresql",
								},
							},
						},
					},
					Volumes: []corev1.Volume{
						{
							Name: "config",
							VolumeSource: corev1.VolumeSource{
								ConfigMap: &corev1.ConfigMapVolumeSource{
									LocalObjectReference: corev1.LocalObjectReference{
										Name: db.Name + "-config",
									},
								},
							},
						},
					},
				},
			},
			VolumeClaimTemplates: []corev1.PersistentVolumeClaim{
				{
					ObjectMeta: metav1.ObjectMeta{
						Name: "data",
					},
					Spec: corev1.PersistentVolumeClaimSpec{
						AccessModes: []corev1.PersistentVolumeAccessMode{
							corev1.ReadWriteOnce,
						},
						Resources: corev1.ResourceRequirements{
							Requests: corev1.ResourceList{
								corev1.ResourceStorage: parseQuantity(db.Spec.Storage.Size),
							},
						},
						StorageClassName: &db.Spec.Storage.StorageClass,
					},
				},
			},
		},
	}
}

// finalizeDatabase handles cleanup before deletion
func (r *DatabaseReconciler) finalizeDatabase(ctx context.Context, db *dbv1alpha1.Database) error {
	logger := log.FromContext(ctx)
	logger.Info("Finalizing Database", "name", db.Name)

	// Perform cleanup tasks:
	// - Delete external resources (cloud databases, backups, etc.)
	// - Clean up any resources not managed by owner references

	logger.Info("Database finalized successfully")
	return nil
}

// updateCondition updates a condition in the status
func (r *DatabaseReconciler) updateCondition(db *dbv1alpha1.Database, conditionType string, status metav1.ConditionStatus, reason, message string) {
	meta.SetStatusCondition(&db.Status.Conditions, metav1.Condition{
		Type:               conditionType,
		Status:             status,
		Reason:             reason,
		Message:            message,
		ObservedGeneration: db.Generation,
		LastTransitionTime: metav1.Now(),
	})
}

// Helper functions
func generatePassword() string {
	// In production, use crypto/rand for secure password generation
	return "changeme123"
}

func parseQuantity(size string) resource.Quantity {
	q, _ := resource.ParseQuantity(size)
	return q
}

// SetupWithManager sets up the controller with the Manager
func (r *DatabaseReconciler) SetupWithManager(mgr ctrl.Manager) error {
	return ctrl.NewControllerManagedBy(mgr).
		For(&dbv1alpha1.Database{}).
		Owns(&appsv1.StatefulSet{}).
		Owns(&corev1.Service{}).
		Owns(&corev1.ConfigMap{}).
		Owns(&corev1.Secret{}).
		Complete(r)
}

controller-runtime version note: the Go examples target controller-runtime ~v0.14. Two breaking changes have landed since: from v0.16, manager Options.MetricsBindAddress/Port were replaced by Metrics: metricsserver.Options{BindAddress: …} and WebhookServer: webhook.NewServer(…); and from v0.15, Watches(&source.Kind{Type: &T{}}, …) became Watches(&T{}, …). Scaffold a fresh project with the current operator-sdk/kubebuilder and follow its generated main.go for the up-to-date API.

// main.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"

	dbv1alpha1 "example.com/database-operator/api/v1alpha1"
	"example.com/database-operator/controllers"
)

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

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

func main() {
	var metricsAddr string
	var enableLeaderElection bool
	var probeAddr string

	flag.StringVar(&metricsAddr, "metrics-bind-address", ":8080", "The address the metric endpoint binds to.")
	flag.StringVar(&probeAddr, "health-probe-bind-address", ":8081", "The address the probe endpoint binds to.")
	flag.BoolVar(&enableLeaderElection, "leader-elect", false,
		"Enable leader election for controller manager. "+
			"Enabling this will ensure there is only one active controller manager.")

	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,
		MetricsBindAddress:     metricsAddr,
		Port:                   9443,
		HealthProbeBindAddress: probeAddr,
		LeaderElection:         enableLeaderElection,
		LeaderElectionID:       "database.example.com",
	})
	if err != nil {
		setupLog.Error(err, "unable to start manager")
		os.Exit(1)
	}

	if err = (&controllers.DatabaseReconciler{
		Client: mgr.GetClient(),
		Scheme: mgr.GetScheme(),
	}).SetupWithManager(mgr); err != nil {
		setupLog.Error(err, "unable to create controller", "controller", "Database")
		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)
	}
}

Operator SDK Commands

# Initialize a new operator project
operator-sdk init --domain example.com --repo github.com/example/database-operator

# Create a new API and controller
operator-sdk create api --group db --version v1alpha1 --kind Database --resource --controller

# Generate CRD manifests
make manifests

# Generate deep copy methods
make generate

# Install CRDs into cluster
make install

# Run operator locally
make run

# Build operator image
make docker-build IMG=myregistry/database-operator:v1.0.0

# Push operator image
make docker-push IMG=myregistry/database-operator:v1.0.0

# Deploy operator to cluster
make deploy IMG=myregistry/database-operator:v1.0.0

# Uninstall CRDs
make uninstall

# Undeploy operator
make undeploy

# Run tests
make test

# Bundle operator for OLM
operator-sdk generate bundle --version 1.0.0

# Validate bundle
operator-sdk bundle validate ./bundle

Kopf (Python)

Kopf is a Python framework for building Kubernetes operators with minimal boilerplate.

Complete Python Operator Example

# database_operator.py
import kopf
import kubernetes
import logging
import hashlib
import secrets
from typing import Dict, Any, Optional
from datetime import datetime

# Configure logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)


@kopf.on.create('db.example.com', 'v1alpha1', 'databases')
def create_database(spec: Dict[str, Any], name: str, namespace: str, **kwargs):
    """
    Handle database creation.
    """
    logger.info(f"Creating database {name} in namespace {namespace}")

    # Load Kubernetes API
    api = kubernetes.client.CoreV1Api()
    apps_api = kubernetes.client.AppsV1Api()

    # Extract spec fields
    db_type = spec.get('type')
    size = spec.get('size', 'small')
    version = spec.get('version', '15')
    replicas = spec.get('replicas', 3)
    storage = spec.get('storage', {})
    storage_size = storage.get('size', '10Gi')
    storage_class = storage.get('storageClass', 'standard')

    # Create Secret for credentials
    secret = create_secret(api, name, namespace)
    logger.info(f"Created secret {secret.metadata.name}")

    # Create ConfigMap for configuration
    config_map = create_config_map(api, name, namespace, size)
    logger.info(f"Created configmap {config_map.metadata.name}")

    # Create headless Service
    service = create_service(api, name, namespace)
    logger.info(f"Created service {service.metadata.name}")

    # Create StatefulSet
    stateful_set = create_stateful_set(
        apps_api, name, namespace, db_type, version,
        replicas, storage_size, storage_class
    )
    logger.info(f"Created statefulset {stateful_set.metadata.name}")

    # Return status to be set on the resource
    return {
        'phase': 'Creating',
        'ready': False,
        'message': f'Database {name} is being created',
        'endpoints': {
            'primary': f'{name}-0.{name}.{namespace}.svc.cluster.local:5432'
        }
    }


@kopf.on.update('db.example.com', 'v1alpha1', 'databases')
def update_database(spec: Dict[str, Any], status: Dict[str, Any],
                   name: str, namespace: str, old: Dict[str, Any],
                   new: Dict[str, Any], diff: tuple, **kwargs):
    """
    Handle database updates.
    """
    logger.info(f"Updating database {name} in namespace {namespace}")
    logger.info(f"Diff: {diff}")

    api = kubernetes.client.CoreV1Api()
    apps_api = kubernetes.client.AppsV1Api()

    # Check what changed
    for operation, field, old_value, new_value in diff:
        if field == ('spec', 'replicas'):
            # Scale StatefulSet
            logger.info(f"Scaling from {old_value} to {new_value} replicas")
            stateful_set = apps_api.read_namespaced_stateful_set(name, namespace)
            stateful_set.spec.replicas = new_value
            apps_api.replace_namespaced_stateful_set(name, namespace, stateful_set)

        elif field == ('spec', 'size'):
            # Update ConfigMap with new size configuration
            logger.info(f"Changing size from {old_value} to {new_value}")
            config_map = create_config_map(api, name, namespace, new_value)
            api.replace_namespaced_config_map(f'{name}-config', namespace, config_map)

            # Restart pods to pick up new config (rolling restart)
            restart_stateful_set(apps_api, name, namespace)

    return {
        'phase': 'Updating',
        'message': f'Database {name} is being updated'
    }


@kopf.on.delete('db.example.com', 'v1alpha1', 'databases')
def delete_database(spec: Dict[str, Any], name: str, namespace: str, **kwargs):
    """
    Handle database deletion (cleanup).
    """
    logger.info(f"Deleting database {name} in namespace {namespace}")

    # Perform cleanup tasks
    # Resources with owner references will be automatically deleted
    # This handler is for external resources or custom cleanup

    # Example: Delete external backups, cloud resources, etc.
    logger.info(f"Cleaning up external resources for {name}")

    return {'message': f'Database {name} deleted successfully'}


@kopf.on.field('db.example.com', 'v1alpha1', 'databases', field='status.phase')
def phase_changed(old: Optional[str], new: Optional[str], name: str, **kwargs):
    """
    React to phase changes.
    """
    logger.info(f"Database {name} phase changed from {old} to {new}")


@kopf.timer('db.example.com', 'v1alpha1', 'databases', interval=30.0)
def check_database_health(spec: Dict[str, Any], status: Dict[str, Any],
                         name: str, namespace: str, **kwargs):
    """
    Periodically check database health and update status.
    """
    apps_api = kubernetes.client.AppsV1Api()

    try:
        # Get StatefulSet status
        stateful_set = apps_api.read_namespaced_stateful_set(name, namespace)
        replicas = spec.get('replicas', 3)
        ready_replicas = stateful_set.status.ready_replicas or 0

        # Determine phase and readiness
        if ready_replicas == replicas:
            phase = 'Running'
            ready = True
            message = f'Database is running with {ready_replicas}/{replicas} replicas ready'
        else:
            phase = 'Updating'
            ready = False
            message = f'Database is updating: {ready_replicas}/{replicas} replicas ready'

        # Build endpoints
        endpoints = {
            'primary': f'{name}-0.{name}.{namespace}.svc.cluster.local:5432',
            'replicas': [
                f'{name}-{i}.{name}.{namespace}.svc.cluster.local:5432'
                for i in range(1, replicas)
            ]
        }

        # Return updated status
        return {
            'phase': phase,
            'ready': ready,
            'message': message,
            'endpoints': endpoints
        }

    except kubernetes.client.exceptions.ApiException as e:
        logger.error(f"Failed to check database health: {e}")
        return {
            'phase': 'Failed',
            'ready': False,
            'message': f'Failed to check database health: {str(e)}'
        }


@kopf.on.event('db.example.com', 'v1alpha1', 'databases')
def log_events(event: Dict[str, Any], **kwargs):
    """
    Log all events for debugging.
    """
    event_type = event.get('type')
    name = event.get('object', {}).get('metadata', {}).get('name')
    logger.debug(f"Event {event_type} for database {name}")


# Helper functions

def create_secret(api: kubernetes.client.CoreV1Api, name: str, namespace: str):
    """Create a secret for database credentials."""
    password = secrets.token_urlsafe(16)

    secret = kubernetes.client.V1Secret(
        api_version='v1',
        kind='Secret',
        metadata=kubernetes.client.V1ObjectMeta(
            name=f'{name}-credentials',
            namespace=namespace,
            labels={'app': name, 'managed-by': 'database-operator'}
        ),
        string_data={
            'username': 'dbuser',
            'password': password,
            'database': name
        }
    )

    try:
        return api.create_namespaced_secret(namespace, secret)
    except kubernetes.client.exceptions.ApiException as e:
        if e.status == 409:  # Already exists
            return api.replace_namespaced_secret(f'{name}-credentials', namespace, secret)
        raise


def create_config_map(api: kubernetes.client.CoreV1Api, name: str,
                     namespace: str, size: str):
    """Create a ConfigMap for database configuration."""

    # Configuration based on size
    configs = {
        'small': {
            'max_connections': '100',
            'shared_buffers': '256MB',
            'effective_cache_size': '1GB',
            'work_mem': '4MB'
        },
        'medium': {
            'max_connections': '200',
            'shared_buffers': '512MB',
            'effective_cache_size': '2GB',
            'work_mem': '8MB'
        },
        'large': {
            'max_connections': '500',
            'shared_buffers': '2GB',
            'effective_cache_size': '8GB',
            'work_mem': '16MB'
        }
    }

    config = configs.get(size, configs['small'])

    config_map = kubernetes.client.V1ConfigMap(
        api_version='v1',
        kind='ConfigMap',
        metadata=kubernetes.client.V1ObjectMeta(
            name=f'{name}-config',
            namespace=namespace,
            labels={'app': name, 'managed-by': 'database-operator'}
        ),
        data=config
    )

    try:
        return api.create_namespaced_config_map(namespace, config_map)
    except kubernetes.client.exceptions.ApiException as e:
        if e.status == 409:  # Already exists
            return api.replace_namespaced_config_map(f'{name}-config', namespace, config_map)
        raise


def create_service(api: kubernetes.client.CoreV1Api, name: str, namespace: str):
    """Create a headless service for the StatefulSet."""

    service = kubernetes.client.V1Service(
        api_version='v1',
        kind='Service',
        metadata=kubernetes.client.V1ObjectMeta(
            name=name,
            namespace=namespace,
            labels={'app': name, 'managed-by': 'database-operator'}
        ),
        spec=kubernetes.client.V1ServiceSpec(
            selector={'app': name},
            cluster_ip='None',  # Headless service
            ports=[
                kubernetes.client.V1ServicePort(
                    name='postgres',
                    port=5432,
                    target_port=5432
                )
            ]
        )
    )

    try:
        return api.create_namespaced_service(namespace, service)
    except kubernetes.client.exceptions.ApiException as e:
        if e.status == 409:  # Already exists
            return api.replace_namespaced_service(name, namespace, service)
        raise


def create_stateful_set(apps_api: kubernetes.client.AppsV1Api, name: str,
                       namespace: str, db_type: str, version: str,
                       replicas: int, storage_size: str, storage_class: str):
    """Create a StatefulSet for the database."""

    image = f'postgres:{version}' if version else 'postgres:15'

    stateful_set = kubernetes.client.V1StatefulSet(
        api_version='apps/v1',
        kind='StatefulSet',
        metadata=kubernetes.client.V1ObjectMeta(
            name=name,
            namespace=namespace,
            labels={'app': name, 'managed-by': 'database-operator'}
        ),
        spec=kubernetes.client.V1StatefulSetSpec(
            service_name=name,
            replicas=replicas,
            selector=kubernetes.client.V1LabelSelector(
                match_labels={'app': name}
            ),
            template=kubernetes.client.V1PodTemplateSpec(
                metadata=kubernetes.client.V1ObjectMeta(
                    labels={'app': name}
                ),
                spec=kubernetes.client.V1PodSpec(
                    containers=[
                        kubernetes.client.V1Container(
                            name='postgres',
                            image=image,
                            ports=[
                                kubernetes.client.V1ContainerPort(
                                    name='postgres',
                                    container_port=5432
                                )
                            ],
                            env=[
                                kubernetes.client.V1EnvVar(
                                    name='POSTGRES_USER',
                                    value_from=kubernetes.client.V1EnvVarSource(
                                        secret_key_ref=kubernetes.client.V1SecretKeySelector(
                                            name=f'{name}-credentials',
                                            key='username'
                                        )
                                    )
                                ),
                                kubernetes.client.V1EnvVar(
                                    name='POSTGRES_PASSWORD',
                                    value_from=kubernetes.client.V1EnvVarSource(
                                        secret_key_ref=kubernetes.client.V1SecretKeySelector(
                                            name=f'{name}-credentials',
                                            key='password'
                                        )
                                    )
                                ),
                                kubernetes.client.V1EnvVar(
                                    name='POSTGRES_DB',
                                    value_from=kubernetes.client.V1EnvVarSource(
                                        secret_key_ref=kubernetes.client.V1SecretKeySelector(
                                            name=f'{name}-credentials',
                                            key='database'
                                        )
                                    )
                                ),
                            ],
                            volume_mounts=[
                                kubernetes.client.V1VolumeMount(
                                    name='data',
                                    mount_path='/var/lib/postgresql/data'
                                ),
                                kubernetes.client.V1VolumeMount(
                                    name='config',
                                    mount_path='/etc/postgresql'
                                )
                            ]
                        )
                    ],
                    volumes=[
                        kubernetes.client.V1Volume(
                            name='config',
                            config_map=kubernetes.client.V1ConfigMapVolumeSource(
                                name=f'{name}-config'
                            )
                        )
                    ]
                )
            ),
            volume_claim_templates=[
                kubernetes.client.V1PersistentVolumeClaim(
                    metadata=kubernetes.client.V1ObjectMeta(
                        name='data'
                    ),
                    spec=kubernetes.client.V1PersistentVolumeClaimSpec(
                        access_modes=['ReadWriteOnce'],
                        resources=kubernetes.client.V1ResourceRequirements(
                            requests={'storage': storage_size}
                        ),
                        storage_class_name=storage_class
                    )
                )
            ]
        )
    )

    try:
        return apps_api.create_namespaced_stateful_set(namespace, stateful_set)
    except kubernetes.client.exceptions.ApiException as e:
        if e.status == 409:  # Already exists
            return apps_api.replace_namespaced_stateful_set(name, namespace, stateful_set)
        raise


def restart_stateful_set(apps_api: kubernetes.client.AppsV1Api, name: str, namespace: str):
    """Trigger a rolling restart of the StatefulSet."""

    # Update annotation to trigger restart
    stateful_set = apps_api.read_namespaced_stateful_set(name, namespace)

    if stateful_set.spec.template.metadata.annotations is None:
        stateful_set.spec.template.metadata.annotations = {}

    stateful_set.spec.template.metadata.annotations['kubectl.kubernetes.io/restartedAt'] = \
        datetime.utcnow().isoformat()

    apps_api.replace_namespaced_stateful_set(name, namespace, stateful_set)
    logger.info(f"Triggered rolling restart of {name}")


# Admission webhooks (optional)

@kopf.on.validate('db.example.com', 'v1alpha1', 'databases')
def validate_database(spec: Dict[str, Any], **kwargs):
    """
    Validate database specifications.
    """
    db_type = spec.get('type')
    size = spec.get('size')
    replicas = spec.get('replicas', 3)

    # Validate database type
    if db_type not in ['postgresql', 'mysql', 'mongodb']:
        raise kopf.AdmissionError(f"Invalid database type: {db_type}. "
                                 f"Must be one of: postgresql, mysql, mongodb")

    # Validate size
    if size not in ['small', 'medium', 'large']:
        raise kopf.AdmissionError(f"Invalid size: {size}. "
                                 f"Must be one of: small, medium, large")

    # Validate replica count
    if replicas < 1 or replicas > 10:
        raise kopf.AdmissionError(f"Invalid replica count: {replicas}. "
                                 f"Must be between 1 and 10")

    # Validate storage
    storage = spec.get('storage')
    if storage:
        size_str = storage.get('size', '')
        if not size_str or not any(size_str.endswith(unit) for unit in ['Gi', 'Mi', 'Ti']):
            raise kopf.AdmissionError(f"Invalid storage size: {size_str}. "
                                     f"Must end with Gi, Mi, or Ti")

    logger.info(f"Validation passed for database spec")


@kopf.on.mutate('db.example.com', 'v1alpha1', 'databases')
def mutate_database(spec: Dict[str, Any], patch: kopf.Patch, **kwargs):
    """
    Mutate (set defaults for) database specifications.

    NOTE: a mutating webhook must emit changes via `patch` (which becomes a
    JSONPatch returned to the API server). Mutating the passed-in `spec` dict
    in place does NOT propagate back to Kubernetes.
    """
    # Set default values if not provided
    if 'replicas' not in spec:
        patch.spec['replicas'] = 3

    if 'size' not in spec:
        patch.spec['size'] = 'small'

    storage = dict(spec.get('storage') or {})
    if 'size' not in storage:
        storage['size'] = '10Gi'
    if 'storageClass' not in storage:
        storage['storageClass'] = 'standard'
    patch.spec['storage'] = storage

    logger.info("Applied defaults to database spec")

Kopf Requirements and Dockerfile

# requirements.txt
kopf>=1.36.0
kubernetes>=26.1.0
# Dockerfile
FROM python:3.11-slim

WORKDIR /app

# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# Copy operator code
COPY database_operator.py .

# Run as non-root
RUN useradd -m -u 1000 operator
USER operator

# Run the operator
CMD ["kopf", "run", "--standalone", "database_operator.py"]

Kopf Deployment Manifests

# deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: database-operator
  namespace: operators
spec:
  replicas: 1
  selector:
    matchLabels:
      app: database-operator
  template:
    metadata:
      labels:
        app: database-operator
    spec:
      serviceAccountName: database-operator
      containers:
        - name: operator
          image: myregistry/database-operator:latest
          imagePullPolicy: Always
          resources:
            limits:
              cpu: 500m
              memory: 512Mi
            requests:
              cpu: 100m
              memory: 128Mi
          env:
            - name: KOPF_NAMESPACE
              valueFrom:
                fieldRef:
                  fieldPath: metadata.namespace
---
apiVersion: v1
kind: ServiceAccount
metadata:
  name: database-operator
  namespace: operators
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
  name: database-operator
rules:
  # Watch and manage custom resources
  - apiGroups: ["db.example.com"]
    resources: ["databases"]
    verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
  - apiGroups: ["db.example.com"]
    resources: ["databases/status"]
    verbs: ["get", "update", "patch"]
  # Manage core resources
  - apiGroups: [""]
    resources: ["services", "configmaps", "secrets", "persistentvolumeclaims"]
    verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
  # Manage StatefulSets
  - apiGroups: ["apps"]
    resources: ["statefulsets"]
    verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
  # Watch events
  - apiGroups: [""]
    resources: ["events"]
    verbs: ["create", "patch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
  name: database-operator
roleRef:
  apiGroup: rbac.authorization.k8s.io
  kind: ClusterRole
  name: database-operator
subjects:
  - kind: ServiceAccount
    name: database-operator
    namespace: operators

Common Operator Patterns

Owner References

Automatically clean up resources when the parent is deleted.

// Set owner reference in Go
if err := controllerutil.SetControllerReference(parent, child, r.Scheme); err != nil {
    return err
}
# Set owner reference in Kopf
kopf.adopt(child_resource)

Finalizers

Ensure cleanup of external resources before deletion.

# Custom resource with finalizer
apiVersion: db.example.com/v1alpha1
kind: Database
metadata:
  name: my-database
  finalizers:
    - db.example.com/finalizer
spec:
  type: postgresql
// Go finalizer pattern
const finalizerName = "db.example.com/finalizer"

if database.ObjectMeta.DeletionTimestamp.IsZero() {
    // Not being deleted, add finalizer
    if !controllerutil.ContainsFinalizer(database, finalizerName) {
        controllerutil.AddFinalizer(database, finalizerName)
        return r.Update(ctx, database)
    }
} else {
    // Being deleted
    if controllerutil.ContainsFinalizer(database, finalizerName) {
        // Perform cleanup
        if err := r.cleanupExternalResources(ctx, database); err != nil {
            return ctrl.Result{}, err
        }

        // Remove finalizer
        controllerutil.RemoveFinalizer(database, finalizerName)
        return r.Update(ctx, database)
    }
}

Status Subresource

Separate status updates from spec changes to avoid reconciliation loops.

# CRD with status subresource
spec:
  versions:
    - name: v1
      subresources:
        status: {}
// Update status in Go
if err := r.Status().Update(ctx, database); err != nil {
    return ctrl.Result{}, err
}

Conditions

Track multiple aspects of resource state.

// Set condition
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
    Type:    "Available",
    Status:  metav1.ConditionTrue,
    Reason:  "DatabaseReady",
    Message: "Database is fully operational",
})

Reconciliation Strategies

Immediate Reconciliation: Reconcile as soon as event received.

return ctrl.Result{}, nil

Delayed Reconciliation: Requeue after a specific duration.

return ctrl.Result{RequeueAfter: 30 * time.Second}, nil

Error Reconciliation: Requeue with exponential backoff.

return ctrl.Result{}, err  // Built-in exponential backoff

Watching Multiple Resources

func (r *DatabaseReconciler) SetupWithManager(mgr ctrl.Manager) error {
    return ctrl.NewControllerManagedBy(mgr).
        For(&dbv1alpha1.Database{}).
        Owns(&appsv1.StatefulSet{}).
        Owns(&corev1.Service{}).
        Watches(
            &source.Kind{Type: &corev1.Secret{}},
            handler.EnqueueRequestsFromMapFunc(r.findDatabasesForSecret),
        ).
        Complete(r)
}

Predicate Filtering

Only reconcile on specific events.

import "sigs.k8s.io/controller-runtime/pkg/predicate"

func (r *DatabaseReconciler) SetupWithManager(mgr ctrl.Manager) error {
    return ctrl.NewControllerManagedBy(mgr).
        For(&dbv1alpha1.Database{}).
        WithEventFilter(predicate.GenerationChangedPredicate{}).
        Complete(r)
}

Deployment and Distribution

Operator Deployment Manifest

# operator-deployment.yaml
apiVersion: v1
kind: Namespace
metadata:
  name: database-operator-system
---
apiVersion: apps/v1
kind: Deployment
metadata:
  name: database-operator-controller-manager
  namespace: database-operator-system
  labels:
    control-plane: controller-manager
spec:
  selector:
    matchLabels:
      control-plane: controller-manager
  replicas: 1
  template:
    metadata:
      labels:
        control-plane: controller-manager
    spec:
      serviceAccountName: database-operator-controller-manager
      containers:
        - name: manager
          image: myregistry/database-operator:v1.0.0
          command:
            - /manager
          args:
            - --leader-elect
          imagePullPolicy: Always
          livenessProbe:
            httpGet:
              path: /healthz
              port: 8081
            initialDelaySeconds: 15
            periodSeconds: 20
          readinessProbe:
            httpGet:
              path: /readyz
              port: 8081
            initialDelaySeconds: 5
            periodSeconds: 10
          resources:
            limits:
              cpu: 500m
              memory: 512Mi
            requests:
              cpu: 100m
              memory: 128Mi
          securityContext:
            allowPrivilegeEscalation: false
            capabilities:
              drop:
                - ALL
            runAsNonRoot: true
      terminationGracePeriodSeconds: 10
---
apiVersion: v1
kind: ServiceAccount
metadata:
  name: database-operator-controller-manager
  namespace: database-operator-system
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
  name: database-operator-manager-role
rules:
  - apiGroups: ["db.example.com"]
    resources: ["databases"]
    verbs: ["create", "delete", "get", "list", "patch", "update", "watch"]
  - apiGroups: ["db.example.com"]
    resources: ["databases/finalizers"]
    verbs: ["update"]
  - apiGroups: ["db.example.com"]
    resources: ["databases/status"]
    verbs: ["get", "patch", "update"]
  - apiGroups: [""]
    resources: ["configmaps", "secrets", "services", "persistentvolumeclaims"]
    verbs: ["create", "delete", "get", "list", "patch", "update", "watch"]
  - apiGroups: ["apps"]
    resources: ["statefulsets"]
    verbs: ["create", "delete", "get", "list", "patch", "update", "watch"]
  - apiGroups: [""]
    resources: ["events"]
    verbs: ["create", "patch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
  name: database-operator-manager-rolebinding
roleRef:
  apiGroup: rbac.authorization.k8s.io
  kind: ClusterRole
  name: database-operator-manager-role
subjects:
  - kind: ServiceAccount
    name: database-operator-controller-manager
    namespace: database-operator-system

Operator Lifecycle Manager (OLM) Integration

# database-operator.clusterserviceversion.yaml
apiVersion: operators.coreos.com/v1alpha1
kind: ClusterServiceVersion
metadata:
  name: database-operator.v1.0.0
  namespace: placeholder
spec:
  displayName: Database Operator
  description: Manages PostgreSQL, MySQL, and MongoDB databases in Kubernetes
  version: 1.0.0
  maturity: stable
  provider:
    name: Example Inc.
  maintainers:
    - name: DevOps Team
      email: devops@example.com
  links:
    - name: Documentation
      url: https://example.com/docs
    - name: Source Code
      url: https://github.com/example/database-operator
  keywords:
    - database
    - postgresql
    - mysql
    - mongodb
  icon:
    - base64data: iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mNk+M9QDwADhgGAWjR9awAAAABJRU5ErkJggg==
      mediatype: image/png
  installModes:
    - type: OwnNamespace
      supported: true
    - type: SingleNamespace
      supported: true
    - type: MultiNamespace
      supported: false
    - type: AllNamespaces
      supported: true
  install:
    strategy: deployment
    spec:
      clusterPermissions:
        - serviceAccountName: database-operator-controller-manager
          rules:
            - apiGroups: ["db.example.com"]
              resources: ["databases"]
              verbs: ["*"]
            - apiGroups: ["apps"]
              resources: ["statefulsets"]
              verbs: ["*"]
            - apiGroups: [""]
              resources: ["services", "configmaps", "secrets"]
              verbs: ["*"]
      deployments:
        - name: database-operator-controller-manager
          spec:
            replicas: 1
            selector:
              matchLabels:
                control-plane: controller-manager
            template:
              metadata:
                labels:
                  control-plane: controller-manager
              spec:
                serviceAccountName: database-operator-controller-manager
                containers:
                  - name: manager
                    image: myregistry/database-operator:v1.0.0
                    resources:
                      limits:
                        cpu: 500m
                        memory: 512Mi
                      requests:
                        cpu: 100m
                        memory: 128Mi
  customresourcedefinitions:
    owned:
      - name: databases.db.example.com
        version: v1alpha1
        kind: Database
        displayName: Database
        description: Represents a managed database instance
        resources:
          - kind: StatefulSet
            version: v1
          - kind: Service
            version: v1
          - kind: Secret
            version: v1
        specDescriptors:
          - path: type
            description: Type of database
            displayName: Database Type
            x-descriptors:
              - 'urn:alm:descriptor:com.tectonic.ui:select:postgresql'
              - 'urn:alm:descriptor:com.tectonic.ui:select:mysql'
              - 'urn:alm:descriptor:com.tectonic.ui:select:mongodb'
          - path: size
            description: Size of database instance
            displayName: Instance Size
            x-descriptors:
              - 'urn:alm:descriptor:com.tectonic.ui:select:small'
              - 'urn:alm:descriptor:com.tectonic.ui:select:medium'
              - 'urn:alm:descriptor:com.tectonic.ui:select:large'
          - path: replicas
            description: Number of replicas
            displayName: Replicas
            x-descriptors:
              - 'urn:alm:descriptor:com.tectonic.ui:podCount'
        statusDescriptors:
          - path: phase
            description: Current phase
            displayName: Phase
            x-descriptors:
              - 'urn:alm:descriptor:io.kubernetes.phase'
          - path: ready
            description: Readiness status
            displayName: Ready
            x-descriptors:
              - 'urn:alm:descriptor:io.kubernetes.conditions'

Testing Operators

Unit Testing (Go)

// controllers/database_controller_test.go
package controllers

import (
    "context"
    "testing"

    . "github.com/onsi/ginkgo/v2"
    . "github.com/onsi/gomega"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/types"
    "sigs.k8s.io/controller-runtime/pkg/reconcile"

    dbv1alpha1 "example.com/database-operator/api/v1alpha1"
)

var _ = Describe("Database controller", func() {
    Context("When reconciling a Database", func() {
        const resourceName = "test-database"

        ctx := context.Background()

        typeNamespacedName := types.NamespacedName{
            Name:      resourceName,
            Namespace: "default",
        }
        database := &dbv1alpha1.Database{}

        BeforeEach(func() {
            By("creating the custom resource for the Kind Database")
            database = &dbv1alpha1.Database{
                ObjectMeta: metav1.ObjectMeta{
                    Name:      resourceName,
                    Namespace: "default",
                },
                Spec: dbv1alpha1.DatabaseSpec{
                    Type:     "postgresql",
                    Size:     "small",
                    Replicas: 3,
                    Storage: &dbv1alpha1.StorageConfig{
                        Size:         "10Gi",
                        StorageClass: "standard",
                    },
                },
            }
            Expect(k8sClient.Create(ctx, database)).To(Succeed())
        })

        AfterEach(func() {
            resource := &dbv1alpha1.Database{}
            err := k8sClient.Get(ctx, typeNamespacedName, resource)
            Expect(err).NotTo(HaveOccurred())

            By("Cleanup the specific resource instance Database")
            Expect(k8sClient.Delete(ctx, resource)).To(Succeed())
        })

        It("should successfully reconcile the resource", func() {
            By("Reconciling the created resource")
            controllerReconciler := &DatabaseReconciler{
                Client: k8sClient,
                Scheme: k8sClient.Scheme(),
            }

            _, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
                NamespacedName: typeNamespacedName,
            })
            Expect(err).NotTo(HaveOccurred())

            // Check that StatefulSet was created
            By("Checking if StatefulSet was created")
            // Add assertions here
        })
    })
})

Integration Testing (Python/Kopf)

# test_database_operator.py
import pytest
import kubernetes
from kubernetes import client, config
from datetime import datetime, timedelta
import time


@pytest.fixture(scope="module")
def k8s_client():
    """Load Kubernetes configuration and return API client."""
    config.load_kube_config()
    return client.CustomObjectsApi()


@pytest.fixture(scope="module")
def core_v1():
    return client.CoreV1Api()


@pytest.fixture(scope="module")
def apps_v1():
    return client.AppsV1Api()


def test_create_database(k8s_client, apps_v1, core_v1):
    """Test database creation."""
    namespace = "test"
    name = "test-postgres"

    # Create custom resource
    database = {
        "apiVersion": "db.example.com/v1alpha1",
        "kind": "Database",
        "metadata": {
            "name": name,
            "namespace": namespace
        },
        "spec": {
            "type": "postgresql",
            "size": "small",
            "replicas": 2,
            "storage": {
                "size": "5Gi",
                "storageClass": "standard"
            }
        }
    }

    try:
        # Create the resource
        k8s_client.create_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            body=database
        )

        # Wait for StatefulSet to be created
        time.sleep(10)

        # Check StatefulSet exists
        sts = apps_v1.read_namespaced_stateful_set(name, namespace)
        assert sts is not None
        assert sts.spec.replicas == 2

        # Check Service exists
        svc = core_v1.read_namespaced_service(name, namespace)
        assert svc is not None
        assert svc.spec.cluster_ip == "None"  # Headless

        # Check Secret exists
        secret = core_v1.read_namespaced_secret(f"{name}-credentials", namespace)
        assert secret is not None
        assert "username" in secret.data
        assert "password" in secret.data

    finally:
        # Cleanup
        k8s_client.delete_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            name=name
        )


def test_update_database_replicas(k8s_client, apps_v1):
    """Test scaling database replicas."""
    namespace = "test"
    name = "test-scaling"

    # Create initial database
    database = {
        "apiVersion": "db.example.com/v1alpha1",
        "kind": "Database",
        "metadata": {"name": name, "namespace": namespace},
        "spec": {
            "type": "postgresql",
            "size": "small",
            "replicas": 2,
            "storage": {"size": "5Gi", "storageClass": "standard"}
        }
    }

    try:
        k8s_client.create_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            body=database
        )

        time.sleep(10)

        # Update replicas
        database["spec"]["replicas"] = 3
        k8s_client.patch_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            name=name,
            body=database
        )

        time.sleep(10)

        # Verify StatefulSet was updated
        sts = apps_v1.read_namespaced_stateful_set(name, namespace)
        assert sts.spec.replicas == 3

    finally:
        k8s_client.delete_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            name=name
        )


def test_database_status_updates(k8s_client):
    """Test that operator updates database status."""
    namespace = "test"
    name = "test-status"

    database = {
        "apiVersion": "db.example.com/v1alpha1",
        "kind": "Database",
        "metadata": {"name": name, "namespace": namespace},
        "spec": {
            "type": "postgresql",
            "size": "small",
            "replicas": 1,
            "storage": {"size": "5Gi", "storageClass": "standard"}
        }
    }

    try:
        k8s_client.create_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            body=database
        )

        # Wait for status to be updated
        for _ in range(30):
            db = k8s_client.get_namespaced_custom_object(
                group="db.example.com",
                version="v1alpha1",
                namespace=namespace,
                plural="databases",
                name=name
            )

            if "status" in db and "phase" in db["status"]:
                assert db["status"]["phase"] in ["Creating", "Running", "Updating"]
                break

            time.sleep(2)
        else:
            pytest.fail("Status was not updated")

    finally:
        k8s_client.delete_namespaced_custom_object(
            group="db.example.com",
            version="v1alpha1",
            namespace=namespace,
            plural="databases",
            name=name
        )

Quick Reference

Operator SDK Commands

Command Description
operator-sdk init Initialise new operator project
operator-sdk create api Create new API and controller
make manifests Generate CRD manifests
make generate Generate deep copy methods
make install Install CRDs into cluster
make run Run operator locally
make deploy Deploy operator to cluster
operator-sdk generate bundle Generate OLM bundle

Kopf Decorators

Decorator Purpose
@kopf.on.create() Handle resource creation
@kopf.on.update() Handle resource updates
@kopf.on.delete() Handle resource deletion
@kopf.timer() Periodic reconciliation
@kopf.on.field() Watch specific field changes
@kopf.on.validate() Admission webhook validation
@kopf.on.mutate() Admission webhook mutation
@kopf.on.event() Watch all events

Common Reconciliation Patterns

Pattern When to Use
Immediate reconciliation Default behaviour for most cases
Periodic reconciliation Health checks, drift detection
Event-driven React to external system changes
Finalizers Clean up external resources
Owner references Automatic cleanup of child resources
Status conditions Track multiple state aspects

Common Issues and Solutions

CRD Not Found

Problem: Custom resource definitions not recognised.

# Check if CRD exists
kubectl get crd databases.db.example.com

# Describe CRD for details
kubectl describe crd databases.db.example.com

# Reinstall CRD
kubectl apply -f config/crd/bases/

Solution: Ensure CRDs are installed before creating custom resources.

Reconciliation Loop

Problem: Operator continuously reconciles without making progress.

Common Causes:

  • Updating spec during status reconciliation
  • Not using status subresource
  • Comparing resource versions incorrectly
// Wrong: Updates spec, triggers new reconciliation
database.Spec.Size = "large"
r.Update(ctx, database)

// Correct: Update status separately
database.Status.Phase = "Running"
r.Status().Update(ctx, database)

Solution: Always use status subresource and avoid modifying spec in reconciliation.

RBAC Permissions

Problem: Operator fails with "forbidden" errors.

# Check operator logs
kubectl logs -n operators deployment/database-operator

# Verify ServiceAccount
kubectl get sa -n operators

# Check ClusterRole permissions
kubectl describe clusterrole database-operator

# Check RoleBinding
kubectl describe clusterrolebinding database-operator

Solution: Ensure RBAC rules grant necessary permissions:

rules:
  - apiGroups: ["db.example.com"]
    resources: ["databases", "databases/status", "databases/finalizers"]
    verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]

Resource Conflicts

Problem: Resources owned by operator conflict with manual changes.

Symptoms:

  • "Object has been modified" errors
  • Resources reverting to operator-defined state

Solution: Use proper conflict resolution:

// Read latest version before update
current := &appsv1.StatefulSet{}
if err := r.Get(ctx, types.NamespacedName{Name: name, Namespace: namespace}, current); err != nil {
    return err
}

// Update from latest version
current.Spec.Replicas = desired.Spec.Replicas
return r.Update(ctx, current)

Finalizer Cleanup Hangs

Problem: Resources stuck in "Terminating" state.

# Check if finalizers are present
kubectl get database my-db -o yaml | grep finalizers -A 5

# Check operator logs for errors
kubectl logs -n operators deployment/database-operator

# Force remove finalizer (last resort!)
kubectl patch database my-db -p '{"metadata":{"finalizers":[]}}' --type=merge

Solution: Ensure finalizer cleanup logic handles errors properly:

if err := r.finalizeDatabase(ctx, database); err != nil {
    // Log error but don't block deletion indefinitely
    logger.Error(err, "Failed to finalize database, removing finalizer anyway")
}
controllerutil.RemoveFinalizer(database, finalizerName)

Webhook Failures

Problem: Admission webhooks timing out or failing.

# Check webhook configuration
kubectl get validatingwebhookconfigurations
kubectl get mutatingwebhookconfigurations

# Test webhook endpoint
kubectl run test --image=curlimages/curl -it --rm -- curl -k https://database-operator-webhook.operators.svc:443/validate

# Check webhook service
kubectl get svc -n operators database-operator-webhook

Solution: Ensure webhook service is running and certificates are valid:

# Regenerate webhook certificates
make webhook-cert

# Verify webhook is reachable
kubectl port-forward -n operators svc/database-operator-webhook 9443:443
curl -k https://localhost:9443/validate

Memory Leaks

Problem: Operator memory usage grows over time.

Diagnosis:

# Monitor memory usage
kubectl top pod -n operators

# Check for goroutine leaks (Go)
curl http://operator-pod:8080/debug/pprof/goroutine

Solutions:

  • Ensure informer caches are properly scoped
  • Avoid storing references to resources in memory
  • Use leader election to prevent multiple reconciliations
  • Implement proper context cancellation
// Use context with timeout for long operations
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()

Debugging Tips

Enable verbose logging:

// Go operator
ctrl.SetLogger(zap.New(zap.UseDevMode(true)))
# Python/Kopf operator
import logging
logging.basicConfig(level=logging.DEBUG)

Use kubectl describe to check events:

kubectl describe database my-db

Check operator logs with filtering:

# Filter by custom resource name
kubectl logs -n operators deployment/database-operator | grep "my-db"

# Follow logs in real-time
kubectl logs -n operators deployment/database-operator -f

# Previous container logs (if crashed)
kubectl logs -n operators deployment/database-operator --previous

Test reconciliation locally:

# Run operator outside cluster (Go)
make run

# Run operator outside cluster (Python)
kopf run --standalone database_operator.py --verbose