Kubernetes Deployment Controller Source Code Analysis

Overview

Source Version: kubernetes-v1.22.3 / commit-id: c92036

The Deployment resource is one of the most commonly used native Kubernetes workload types. Most practitioners begin their Kubernetes journey by running workloads through the Deployment resource.

In the previous article, we covered the complete feature set of Deployments, focusing on "rolling updates" and "rollback" capabilities to establish a conceptual understanding of what Deployments can accomplish. Building on that foundation, we will now examine the implementation details through a source code analysis.

Prerequisite Knowledge: Understanding the Deployment source code requires familiarity with how custom controllers operate. This includes the Informer mechanism, workqueue (rate-limiting delayed queues), and ResourceEventHandler patterns. Readers lacking this background may find this content challenging. Reviewing related articles on Kubernetes controller principles is recommended before proceeding.

This analysis is divided into two parts:

Part One: Deployment Feature Overview Part Two: Source Code Implementation Walkthrough

Entry Point: startDeploymentController Function

The initialization and startup entry point for the DeploymentController is the startDeploymentController() function located at cmd/kube-controller-manager/app/apps.go:72.

func startDeploymentController(ctx ControllerContext) (http.Handler, bool, error) {
    controller, err := deployment.NewDeploymentController(
        ctx.InformerFactory.Apps().V1().Deployments(),
        ctx.InformerFactory.Apps().V1().ReplicaSets(),
        ctx.InformerFactory.Core().V1().Pods(),
        ctx.ClientBuilder.ClientOrDie("deployment-controller"),
    )
    if err != nil {
        return nil, true, fmt.Errorf("failed to create Deployment controller: %v", err)
    }
    
    go controller.Run(
        int(ctx.ComponentConfig.DeploymentController.ConcurrentDeploymentSyncs),
        ctx.Stop,
    )
    return nil, true, nil
}

This function first initializes a DeploymentController instance using the NewDeploymentController() method. The parameters include DeploymentInformer, ReplicaSetInformer, PodInformer, and a Clientset. These dependencies grant the DeploymentController the ability to receive resource change events for Deployments, ReplicaSets, and Pods, as well as perform CRUD operations on the API server for these resources.

The function then calls the Run() method to start the controller. The parameter ConcurrentDeploymentSyncs defaults to 5, meaning the controller can concurrently reconcile up to 5 Deployments by default.

DeploymentController Structure

Let's examine the DeploymentController type definition and its initialization process.

Type Definition

The controller structure is defined at pkg/controller/deployment/deployment_controller.go:68:

type DeploymentController struct {
    // ReplicaSet control interface
    rsControl     controller.RSControlInterface
    kubeClient    clientset.Interface
    eventRecorder record.EventRecorder

    syncHandler func(dKey string) error
    // Test utility function
    enqueueDeployment func(deployment *apps.Deployment)

    // Cache access for Deployments
    dLister appslisters.DeploymentLister
    // Cache access for ReplicaSets
    rsLister appslisters.ReplicaSetLister
    // Cache access for Pods
    podLister corelisters.PodLister

    dListerSynced cache.InformerSynced
    rsListerSynced cache.InformerSynced
    podListerSynced cache.InformerSynced

    // Workqueue with rate limiting capabilities
    queue workqueue.RateLimitingInterface
}

Initialization

The initialization logic is implemented at pkg/controller/deployment/deployment_controller.go:101:

func NewDeploymentController(
    dInformer appsinformers.DeploymentInformer,
    rsInformer appsinformers.ReplicaSetInformer,
    podInformer coreinformers.PodInformer,
    client clientset.Interface,
) (*DeploymentController, error) {
    // Initialize event broadcasting
    broadcaster := record.NewBroadcaster()
    broadcaster.StartStructuredLogging(0)
    broadcaster.StartRecordingToSink(&v1core.EventSinkImpl{
        Interface: client.CoreV1().Events(""),
    })

    // Create controller instance
    controller := &DeploymentController{
        kubeClient:    client,
        eventRecorder: broadcaster.NewRecorder(
            scheme.Scheme, 
            v1.EventSource{Component: "deployment-controller"},
        ),
        queue: workqueue.NewNamedRateLimitingQueue(
            workqueue.DefaultControllerRateLimiter(), 
            "deployment",
        ),
    }

    // ReplicaSet control with client and event recorder
    controller.rsControl = controller.RealRSControl{
        KubeClient: client,
        Recorder:   controller.eventRecorder,
    }

    // Register resource event handlers
    dInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc:    controller.addDeployment,
        UpdateFunc: controller.updateDeployment,
        DeleteFunc: controller.deleteDeployment,
    })
    rsInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc:    controller.addReplicaSet,
        UpdateFunc: controller.updateReplicaSet,
        DeleteFunc: controller.deleteReplicaSet,
    })
    podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        DeleteFunc: controller.deletePod,
    })

    // Set sync handler and enqueue function
    controller.syncHandler = controller.syncDeployment
    controller.enqueueDeployment = controller.enqueue

    // Cache listers
    controller.dLister = dInformer.Lister()
    controller.rsLister = rsInformer.Lister()
    controller.podLister = podInformer.Lister()
    controller.dListerSynced = dInformer.Informer().HasSynced
    controller.rsListerSynced = rsInformer.Informer().HasSynced
    controller.podListerSynced = podInformer.Informer().HasSynced

    return controller, nil
}

Resource Event Handlers

The initialization code registers several event handler callbacks that respond to changes in cluster resources. These handlers determine when items are added to the workqueue for processing.

Deployment Change Handlers

The three Deployment-related handlers are straightforward and operate similarly:

func (dc *DeploymentController) addDeployment(obj interface{}) {
    deployment := obj.(*apps.Deployment)
    klog.V(4).InfoS("Adding deployment", "deployment", klog.KObj(deployment))
    dc.enqueueDeployment(deployment)
}

func (dc *DeploymentController) updateDeployment(old, cur interface{}) {
    oldD := old.(*apps.Deployment)
    curD := cur.(*apps.Deployment)
    klog.V(4).InfoS("Updating deployment", "deployment", klog.KObj(oldD))
    // Only the current Deployment is enqueued
    dc.enqueueDeployment(curD)
}

func (dc *DeploymentController) deleteDeployment(obj interface{}) {
    deployment, ok := obj.(*apps.Deployment)
    if !ok {
        // Handle DeletedFinalStateUnknown scenario
        tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
        if !ok {
            utilruntime.HandleError(fmt.Errorf(
                "couldn't get object from tombstone %#v", obj))
            return
        }
        deployment, ok = tombstone.Obj.(*apps.Deployment)
        if !ok {
            utilruntime.HandleError(fmt.Errorf(
                "tombstone contained object that is not a Deployment %#v", obj))
            return
        }
    }
    klog.V(4).InfoS("Deleting deployment", "deployment", klog.KObj(deployment))
    dc.enqueueDeployment(deployment)
}

ReplicaSet Change Handlers

The ReplicaSet event handlers have more complex logic due to controller reference resolution and orphan handling:

Add Handler - Defined at pkg/controller/deployment/deployment_controller.go:199:

func (dc *DeploymentController) addReplicaSet(obj interface{}) {
    rs := obj.(*apps.ReplicaSet)
    
    // If marked for deletion during restart, handle as delete event
    if rs.DeletionTimestamp != nil {
        dc.deleteReplicaSet(rs)
        return
    }

    // Attempt to find controlling Deployment
    if controllerRef := metav1.GetControllerOf(rs); controllerRef != nil {
        d := dc.resolveControllerRef(rs.Namespace, controllerRef)
        if d == nil {
            return
        }
        klog.V(4).InfoS("ReplicaSet added", "replicaSet", klog.KObj(rs))
        dc.enqueueDeployment(d)
        return
    }

    // Handle orphan ReplicaSet - check if any Deployment can adopt it
    candidates := dc.getDeploymentsForReplicaSet(rs)
    if len(candidates) == 0 {
        return
    }
    klog.V(4).InfoS("Orphan ReplicaSet added", "replicaSet", klog.KObj(rs))
    
    // Enqueue all matching Deployments
    for _, deployment := range candidates {
        dc.enqueueDeployment(deployment)
    }
}

Update Handler - Defined at pkg/controller/deployment/deployment_controller.go:256:

func (dc *DeploymentController) updateReplicaSet(old, cur interface{}) {
    currentRS := cur.(*apps.ReplicaSet)
    oldRS := old.(*apps.ReplicaSet)
    
    // Skip if only resyncing with no meaningful change
    if currentRS.ResourceVersion == oldRS.ResourceVersion {
        return
    }

    currentRef := metav1.GetControllerOf(currentRS)
    oldRef := metav1.GetControllerOf(oldRS)
    controllerRefChanged := !reflect.DeepEqual(currentRef, oldRef)

    // Notify previous controller if reference changed
    if controllerRefChanged && oldRef != nil {
        if d := dc.resolveControllerRef(oldRS.Namespace, oldRef); d != nil {
            dc.enqueueDeployment(d)
        }
    }

    // Notify current controller
    if currentRef != nil {
        d := dc.resolveControllerRef(currentRS.Namespace, currentRef)
        if d == nil {
            return
        }
        klog.V(4).InfoS("ReplicaSet updated", "replicaSet", klog.KObj(currentRS))
        dc.enqueueDeployment(d)
        return
    }

    // Handle orphan ReplicaSet label changes
    labelChanged := !reflect.DeepEqual(currentRS.Labels, oldRS.Labels)
    if labelChanged || controllerRefChanged {
        candidates := dc.getDeploymentsForReplicaSet(currentRS)
        if len(candidates) == 0 {
            return
        }
        klog.V(4).InfoS("Orphan ReplicaSet updated", "replicaSet", klog.KObj(currentRS))
        for _, d := range candidates {
            dc.enqueueDeployment(d)
        }
    }
}

Delete Handler - Defined at pkg/controller/deployment/deployment_controller.go:304:

func (dc *DeploymentController) deleteReplicaSet(obj interface{}) {
    rs, ok := obj.(*apps.ReplicaSet)

    // Handle DeletedFinalStateUnknown during deletion
    if !ok {
        tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
        if !ok {
            utilruntime.HandleError(fmt.Errorf(
                "couldn't get object from tombstone %#v", obj))
            return
        }
        rs, ok = tombstone.Obj.(*apps.ReplicaSet)
        if !ok {
            utilruntime.HandleError(fmt.Errorf(
                "tombstone contained object that is not a ReplicaSet %#v", obj))
            return
        }
    }

    // Orphan ReplicaSet deletion doesn't require controller notification
    controllerRef := metav1.GetControllerOf(rs)
    if controllerRef == nil {
        return
    }
    
    d := dc.resolveControllerRef(rs.Namespace, controllerRef)
    if d == nil {
        return
    }
    
    klog.V(4).InfoS("ReplicaSet deleted", "replicaSet", klog.KObj(rs))
    dc.enqueueDeployment(d)
}

Controller Startup

After examining what events trigger workqueue additions, let's look at how these items are consumed during normal operation.

The Run Method

The Run() method is concise - it starts worker goroutines based on the specified concurrency count:

func (dc *DeploymentController) Run(workers int, stopCh <-chan struct{}) {
    defer utilruntime.HandleCrash()
    defer dc.queue.ShutDown()

    klog.InfoS("Starting controller", "controller", "deployment")
    defer klog.InfoS("Shutting down controller", "controller", "deployment")

    // Wait for all cache synchronizations
    if !cache.WaitForNamedCacheSync("deployment", stopCh,
        dc.dListerSynced, dc.rsListerSynced, dc.podListerSynced) {
        return
    }

    // Start worker goroutines
    for i := 0; i < workers; i++ {
        go wait.Until(dc.worker, time.Second, stopCh)
    }

    <-stopCh
}

The worker loop implementation retrieves and processes items from the workqueue:

func (dc *DeploymentController) worker() {
    for dc.processNextWorkItem() {
    }
}

func (dc *DeploymentController) processNextWorkItem() bool {
    key, quit := dc.queue.Get()
    if quit {
        return false
    }
    defer dc.queue.Done(key)
    
    // Invoke sync handler for the item
    err := dc.syncHandler(key.(string))
    dc.handleErr(err, key)

    return true
}

The syncHandler field points to dc.syncDeployment, which performs the actual reconciliation work for each Deployment.

The syncDeployment Method

The syncDeployment() method processes items from the workqueue by extracting namespace and name from the key, then reconciling the corresponding Deployment:

func (dc *DeploymentController) syncDeployment(key string) error {
    // Extract namespace and name from the key
    namespace, name, err := cache.SplitMetaNamespaceKey(key)
    if err != nil {
        klog.ErrorS(err, "Failed to parse cache key", "cacheKey", key)
        return err
    }

    startTime := time.Now()
    klog.V(4).InfoS("Beginning deployment sync", 
        "deployment", klog.KRef(namespace, name), 
        "startTime", startTime)
    defer func() {
        klog.V(4).InfoS("Completed deployment sync", 
            "deployment", klog.KRef(namespace, name), 
            "duration", time.Since(startTime))
    }()

    // Retrieve Deployment from cache
    deployment, err := dc.dLister.Deployments(namespace).Get(name)
    if errors.IsNotFound(err) {
        klog.V(2).InfoS("Deployment has been removed", 
            "deployment", klog.KRef(namespace, name))
        return nil
    }
    if err != nil {
        return err
    }

    // Create thread-safe copy for modifications
    target := deployment.DeepCopy()

    // Handle empty selector with warning event
    emptySelector := metav1.LabelSelector{}
    if reflect.DeepEqual(target.Spec.Selector, &emptySelector) {
        dc.eventRecorder.Eventf(target, v1.EventTypeWarning, 
            "SelectingAll", 
            "This deployment is selecting all pods. A non-empty selector is required.")
        if target.Status.ObservedGeneration < target.Generation {
            target.Status.ObservedGeneration = target.Generation
            dc.kubeClient.AppsV1().Deployments(target.Namespace).
                UpdateStatus(context.TODO(), target, metav1.UpdateOptions{})
        }
        return nil
    }

    // Collect all ReplicaSets owned by this Deployment
    rsList, err := dc.getReplicaSetsForDeployment(target)
    if err != nil {
        return err
    }

    // Build Pod map grouped by owning ReplicaSet
    podMap, err := dc.getPodMapForDeployment(target, rsList)
    if err != nil {
        return err
    }

    // Handle deletion: update status only
    if target.DeletionTimestamp != nil {
        return dc.syncStatusOnly(target, rsList)
    }

    // Update pause condition if needed
    if err = dc.checkPausedConditions(target); err != nil {
        return err
    }

    // Handle paused state
    if target.Spec.Paused {
        return dc.sync(target, rsList)
    }

    // Handle legacy rollback annotation
    if rollbackTo := getRollbackTo(target); rollbackTo != nil {
        return dc.rollback(target, rsList)
    }

    // Detect scaling events
    isScaling, err := dc.isScalingEvent(target, rsList)
    if err != nil {
        return err
    }
    if isScaling {
        return dc.sync(target, rsList)
    }

    // Execute deployment strategy
    switch target.Spec.Strategy.Type {
    case apps.RecreateDeploymentStrategyType:
        return dc.rolloutRecreate(target, rsList, podMap)
    case apps.RollingUpdateDeploymentStrategyType:
        return dc.rolloutRolling(target, rsList)
    }

    return fmt.Errorf("unknown deployment strategy: %s", target.Spec.Strategy.Type)
}

This method orchestrates the complete reconciliation process, including handling deletion, pause states, rollbacks, scaling events, and executing the appropriate deployment strategy.

Summary

Having examined the syncDeployment() method, the core logic of the Deployment controller has been covered. While some helper method implementations were described functionally rather than explored in full detail, their individual logic is straightforward and follows consistent patterns established in the controller framework.

For readers familiar with the client-go library and other controller implementations (such as the Job Controller), the patterns here will feel familiar and approachable. Those new to controller development should consider reviewing the foundasional articles on Kubernetes controller principles before diving into additional controller implementations.

Tags: kubernetes deployment-controller source-code controller-runtime Informer

Posted on Mon, 05 Oct 2026 16:40:55 +0000 by COOMERDP