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.