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.
graph TB
subgraph "Operator Pattern Architecture"
USER[User/Admin] -->|Creates/Updates| CR[Custom Resource]
CR -->|Watched by| CTRL[Controller/Operator]
subgraph "Control Loop"
CTRL -->|Reads| CACHE[Informer Cache]
CACHE -->|Triggers| QUEUE[Work Queue]
QUEUE -->|Reconcile| RECON[Reconciliation Logic]
RECON -->|Update| STATUS[Resource Status]
RECON -->|Create/Update/Delete| K8S[Kubernetes Resources]
K8S -->|Events| CACHE
end
CTRL -->|Manages| DEP[Deployment]
CTRL -->|Manages| SVC[Service]
CTRL -->|Manages| CM[ConfigMap]
CTRL -->|Manages| SEC[Secret]
STATUS -->|Feedback| CR
end
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
flowchart TD
START([Start]) --> WATCH[Watch Custom Resource]
WATCH --> EVENT{Event Received?}
EVENT -->|No| WAIT[Wait]
WAIT --> EVENT
EVENT -->|Yes| FETCH[Fetch Current State]
FETCH --> COMPARE{Current == Desired?}
COMPARE -->|Yes| SUCCESS[Update Status: Success]
SUCCESS --> WATCH
COMPARE -->|No| ACTIONS[Determine Actions]
ACTIONS --> CREATE{Need Create?}
CREATE -->|Yes| DOCREATE[Create Resources]
CREATE -->|No| UPDATE{Need Update?}
DOCREATE --> UPDATE
UPDATE -->|Yes| DOUPDATE[Update Resources]
UPDATE -->|No| DELETE{Need Delete?}
DOUPDATE --> DELETE
DELETE -->|Yes| DODELETE[Delete Resources]
DELETE -->|No| UPDATESTATUS[Update Status]
DODELETE --> UPDATESTATUS
UPDATESTATUS --> ERROR{Error?}
ERROR -->|Yes| REQUEUE[Requeue with Backoff]
ERROR -->|No| WATCH
REQUEUE --> WATCH
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/Portwere replaced byMetrics: metricsserver.Options{BindAddress: …}andWebhookServer: webhook.NewServer(…); and from v0.15,Watches(&source.Kind{Type: &T{}}, …)becameWatches(&T{}, …). Scaffold a fresh project with the currentoperator-sdk/kubebuilderand follow its generatedmain.gofor 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