基于 Kubernetes v1.37 源码(pkg/controller/deployment/pkg/controller/replicaset/

一、什么是控制器? @

在 Kubernetes 中,控制器(Controller) 是一个永不休止的控制循环,它监控集群中资源对象的实际状态,并驱动其向用户声明的期望状态收敛。

用户提交期望状态(Spec)
       ↓
控制器读取 Spec
       ↓
对比实际状态(Status)
       ↓
执行操作(创建/删除/更新资源)
       ↓
更新 Status
       ↓
循环

这就是著名的调谐循环(Reconcile Loop),也称为控制循环(Control Loop)。k8s 中所有控制器都遵循这一模式:Deployment、StatefulSet、DaemonSet、Job、HPA、Namespace……


二、控制器的三大支柱 @

在深入 Deployment 控制器之前,先理解构成控制器的三个核心组件:

2.1 Informer(监听器) @

Informer 监听 API Server 中资源的变化,并将事件推送到工作队列。

// deployment_controller.go
dInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
    AddFunc: func(obj interface{}) {
        dc.addDeployment(logger, obj)
    },
    UpdateFunc: func(oldObj, newObj interface{}) {
        dc.updateDeployment(logger, oldObj, newObj)
    },
    DeleteFunc: func(obj interface{}) {
        dc.deleteDeployment(logger, obj)
    },
})

Informer 内部维护了本地缓存(Store),避免每次操作都查询 API Server,大幅降低 etcd 压力。

2.2 Work Queue(工作队列) @

工作队列是控制器的"待办事项列表",具有以下特性:

  • 去重:同一资源在队列中只出现一次
  • 限速:失败重试时指数退避(5ms → 10ms → 20ms → …)
  • 异步:多个 worker 并发处理
// deployment_controller.go
queue: workqueue.NewTypedRateLimitingQueueWithConfig(
    workqueue.DefaultTypedControllerRateLimiter[string](),
    workqueue.TypedRateLimitingQueueConfig[string]{
        Name: "deployment",
    },
),

2.3 Lister(列表器) @

Lister 从 Informer 的本地缓存中读取数据,提供快速查询能力。

// deployment_controller.go
dLister appslisters.DeploymentLister   // 查询 Deployment
rsLister appslisters.ReplicaSetLister  // 查询 ReplicaSet
podLister corelisters.PodLister        // 查询 Pod

三、Deployment 控制器全景 @

Deployment 控制器负责管理 Deployment 资源,其核心职责是:

  1. 管理 ReplicaSet:创建新的 RS、淘汰旧的 RS
  2. 执行滚动更新:按策略逐步替换 Pod
  3. 维护状态:更新 Deployment 的 Status 字段

3.1 结构定义 @

type DeploymentController struct {
    rsControl controller.RSControlInterface    // RS 管理(创建/删除/更新)
    client    clientset.Interface              // API Server 客户端
    
    dLister   appslisters.DeploymentLister     // Deployment 缓存
    rsLister  appslisters.ReplicaSetLister     // ReplicaSet 缓存
    podLister corelisters.PodLister            // Pod 缓存
    
    dListerSynced  cache.InformerSynced        // 缓存是否同步完成
    rsListerSynced cache.InformerSynced
    podListerSynced cache.InformerSynced
    
    queue workqueue.TypedRateLimitingInterface[string]  // 工作队列
}

3.2 事件监听 @

Deployment 控制器监听三种资源的变化:

资源 监听事件 触发动作
Deployment Add/Update/Delete 将对应 Deployment 加入队列
ReplicaSet Add/Update/Delete 找到所属的 Deployment,加入队列
Pod Delete 如果是 Recreate 策略,触发重新同步

为什么监听 Pod 的 Delete?

因为 Recreate 策略需要先删除所有旧 Pod,再创建新 Pod。当旧 Pod 被删除时,控制器需要知道"旧 Pod 已全部清除",才能继续创建新的 ReplicaSet。


四、调谐循环的核心逻辑 @

Deployment 控制器的入口是 syncDeployment 方法,每次从队列取出一个 Deployment key 时调用:

func (dc *DeploymentController) syncDeployment(ctx context.Context, key string) error {
    // 1. 从缓存读取 Deployment
    deployment, err := dc.dLister.Deployments(namespace).Get(name)
    
    // 2. 深拷贝,避免修改缓存
    d := deployment.DeepCopy()
    
    // 3. 获取该 Deployment 管理的所有 ReplicaSet
    rsList, err := dc.getReplicaSetsForDeployment(ctx, d)
    
    // 4. 根据状态分发到不同处理逻辑
    if d.DeletionTimestamp != nil {
        return dc.syncStatusOnly(ctx, d, rsList)
    }
    // 暂停/恢复时更新 Deployment 条件,避免恢复后 progressDeadline 误判
    if err = dc.checkPausedConditions(ctx, d); err != nil {
        return err
    }
    if d.Spec.Paused {
        return dc.sync(ctx, d, rsList)
    }
    if getRollbackTo(d) != nil {
        return dc.rollback(ctx, d, rsList)
    }
    
    // 5. 纯扩容/缩容事件(仅副本数变化)走 dc.sync
    scalingEvent, err := dc.isScalingEvent(ctx, d, rsList)
    if scalingEvent {
        return dc.sync(ctx, d, rsList)
    }
    
    // 6. 根据策略类型执行不同的发布逻辑
    switch d.Spec.Strategy.Type {
    case apps.RecreateDeploymentStrategyType:
        return dc.rolloutRecreate(ctx, d, rsList, podMap)
    case apps.RollingUpdateDeploymentStrategyType:
        return dc.rolloutRolling(ctx, d, rsList)
    }
}

4.1 两种发布策略 @

RollingUpdate(滚动更新,默认) @

逐步替换 Pod,确保服务不中断:

func (dc *DeploymentController) rolloutRolling(ctx context.Context, d, rsList) error {
    newRS, oldRSs, _ := dc.getAllReplicaSetsAndSyncRevision(ctx, d, rsList, true)
    
    // 1. 扩容新 RS
    scaledUp, _ := dc.reconcileNewReplicaSet(ctx, allRSs, newRS, d)
    if scaledUp { return dc.syncRolloutStatus(...) }
    
    // 2. 缩容旧 RS
    scaledDown, _ := dc.reconcileOldReplicaSets(ctx, allRSs, oldRSs, newRS, d)
    if scaledDown { return dc.syncRolloutStatus(...) }
    
    // 3. 清理旧 RS
    if deploymentutil.DeploymentComplete(d, &d.Status) {
        dc.cleanupDeployment(ctx, oldRSs, d)
    }
    
    return dc.syncRolloutStatus(ctx, allRSs, newRS, d)
}

关键参数:

strategy:
  type: RollingUpdate
  rollingUpdate:
    maxUnavailable: 25%   # 最大不可用 Pod 数
    maxSurge: 25%         # 最大可超出的 Pod 数

Recreate(重建) @

先删除所有旧 Pod,再创建新 Pod:

func (dc *DeploymentController) rolloutRecreate(ctx context.Context, d, rsList, podMap) error {
    newRS, oldRSs, _ := dc.getAllReplicaSetsAndSyncRevision(ctx, d, rsList, false)
    
    // 1. 先缩容旧 RS
    scaledDown, _ := dc.scaleDownOldReplicaSetsForRecreate(ctx, activeOldRSs, d)
    if scaledDown { return dc.syncRolloutStatus(...) }
    
    // 2. 等待旧 Pod 全部终止
    if oldPodsRunning(newRS, oldRSs, podMap) {
        return dc.syncRolloutStatus(...)
    }
    
    // 3. 创建新 RS
    if newRS == nil {
        newRS, oldRSs, _ = dc.getAllReplicaSetsAndSyncRevision(ctx, d, rsList, true)
    }
    
    // 4. 扩容新 RS
    if _, err := dc.scaleUpNewReplicaSetForRecreate(ctx, newRS, d); err != nil {
        return err
    }
    
    // 5. 部署完成后清理旧 RS
    if deploymentutil.DeploymentComplete(d, &d.Status) {
        dc.cleanupDeployment(ctx, oldRSs, d)
    }
    return dc.syncRolloutStatus(...)
}

五、期望机制(Expectations) @

控制器如何知道"我的操作已经完成"?答案是期望机制

5.1 核心思想 @

每次控制器执行创建/删除操作时,它会记录一个"期望":

// controller_utils.go
func (r *ControllerExpectations) SetExpectations(logger, controllerKey, add, del int) {
    exp := &ControlleeExpectations{
        add: int64(add),  // 期望创建的数量
        del: int64(del),  // 期望删除的数量
        key: controllerKey,
        timestamp: clock.RealClock{}.Now(),
    }
    r.Add(exp)
}

当 Informer 观察到对应的资源变化时,递减期望计数:

func (r *ControllerExpectations) CreationObserved(logger, controllerKey string) {
    r.LowerExpectations(logger, controllerKey, 1, 0)
}

func (r *ControllerExpectations) DeletionObserved(logger, controllerKey string) {
    r.LowerExpectations(logger, controllerKey, 0, 1)
}

adddel 都 ≤ 0 时,表示期望已满足:

func (e *ControlleeExpectations) Fulfilled() bool {
    return atomic.LoadInt64(&e.add) <= 0 && atomic.LoadInt64(&e.del) <= 0
}

5.2 为什么需要期望机制? @

问题场景: 控制器创建一个 Pod,但还没收到 Informer 的创建事件(网络延迟),它可能会误以为创建失败,再次创建同一个 Pod。

解决方案: 创建 Pod 前先记录期望。在期望被满足之前(即 Informer 的创建事件尚未全部到达),SatisfiedExpectations 返回 false,控制器会跳过本轮的 manageReplicas,避免重复创建;待事件陆续到达、期望计数递减完毕后,才恢复正常的副本管理。

5.3 UID 追踪期望 @

对于删除操作,普通期望只记录数量。但存在一个问题:Pod 的删除会先以 Update 事件(设置 DeletionTimestamp)、再以 Delete 事件 被 Informer 观察到,如果期望只按数量递减,同一次删除可能被重复计数(double counting)。

UID 追踪期望 解决了这个问题,它记录的是具体要删除的 Pod key,而不是数量,同一次删除只会被计数一次:

type UIDTrackingControllerExpectations struct {
    ControllerExpectationsInterface
    uidStore cache.Store  // 存储待删除 Pod 的 key 集合
}

六、ReplicaSet 控制器 @

Deployment 控制器本身不直接操作 Pod,它通过管理 ReplicaSet 来间接管理 Pod。ReplicaSet 控制器才是真正创建/删除 Pod 的角色。

6.1 核心结构 @

type ReplicaSetController struct {
    kubeClient clientset.Interface
    podControl controller.PodControlInterface  // Pod 管理
    
    rsLister  appslisters.ReplicaSetLister
    podLister corelisters.PodLister
    
    expectations *controller.UIDTrackingControllerExpectations  // 期望机制
    
    queue workqueue.TypedRateLimitingInterface[string]
}

6.2 调谐循环 @

func (rsc *ReplicaSetController) syncReplicaSet(ctx context.Context, key string) error {
    // 1. 获取 ReplicaSet
    rs, err := rsc.rsLister.ReplicaSets(namespace).Get(name)
    
    // 2. 检查期望是否满足
    rsNeedsSync := rsc.expectations.SatisfiedExpectations(logger, key)
    
    // 3. 获取所有管理的 Pod(通过 Pod 索引按 owner 过滤,并认领孤儿 Pod)
    allRSPods, err := controller.FilterPodsByOwner(rsc.podIndexer, &rs.ObjectMeta, rsc.Kind, true)
    activePods, err := rsc.claimPods(ctx, rs, selector, controller.FilterActivePods(logger, allRSPods))
    
    // 4. 期望未满足时跳过 manageReplicas,但仍会继续更新 RS Status
    if rsNeedsSync && rs.DeletionTimestamp == nil {
        manageReplicasErr = rsc.manageReplicas(ctx, activePods, rs)
    }
    
    // 5. 无论期望是否满足,都更新 RS Status
    return updateReplicaSetStatus(...)
}

manageReplicas先设置期望,再执行操作

// replica_set.go manageReplicas
diff := len(activePods) - int(*(rs.Spec.Replicas))
if diff < 0 {
    // 副本不足:先记录期望创建数,再批量创建
    rsc.expectations.ExpectCreations(logger, rsKey, diff)
    successfulCreations, err := slowStartBatch(diff, controller.SlowStartInitialBatchSize, func() error {
        return rsc.podControl.CreatePods(...)
    })
} else if diff > 0 {
    // 副本过多:先记录期望删除的 Pod key 列表,再并发删除
    podsToDelete := getPodsToDelete(activePods, relatedPods, diff)
    rsc.expectations.ExpectDeletions(logger, rsKey, getPodKeys(podsToDelete))
    // 并发调用 podControl.DeletePod 删除这些 Pod
}

6.3 慢启动(Slow Start) @

当需要一次性创建大量 Pod 时,ReplicaSet 控制器采用慢启动策略,避免瞬间压垮 API Server:

// replica_set.go(初始批量常量 SlowStartInitialBatchSize = 1 定义在 controller_utils.go)
// 批量大小:1, 2, 4, 8, 16, ...
func slowStartBatch(count int, initialBatchSize int, fn func() error) (int, error) {
    remaining := count
    successes := 0
    for batchSize := min(remaining, initialBatchSize); batchSize > 0; batchSize = min(2*batchSize, remaining) {
        var wg sync.WaitGroup
        for i := 0; i < batchSize; i++ {
            go func() { // 同一批内并发调用 fn
                defer wg.Done()
                if err := fn(); err != nil { /* 记录错误 */ }
            }()
        }
        wg.Wait()
        // 整批成功后批量翻倍;批内有失败则剩余批次全部跳过
    }
    return successes, err
}

七、并发控制 @

7.1 工作队列 @

Deployment 控制器使用限速队列,支持:

  • 去重:同一 key 在队列中只出现一次
  • 退避:失败重试时指数退避
  • 限速:限制每秒操作次数
queue: workqueue.NewTypedRateLimitingQueueWithConfig(
    workqueue.DefaultTypedControllerRateLimiter[string](),
    workqueue.TypedRateLimitingQueueConfig[string]{
        Name: "deployment",
    },
),

7.2 并发 Worker @

控制器启动多个 worker 协程,并发处理队列中的任务:

func (dc *DeploymentController) Run(ctx context.Context, workers int) {
    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Go(func() {
            wait.UntilWithContext(ctx, dc.worker, time.Second)
        })
    }
    <-ctx.Done()
}

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

func (dc *DeploymentController) processNextWorkItem(ctx context.Context) bool {
    key, quit := dc.queue.Get()
    if quit {
        return false
    }
    defer dc.queue.Done(key)
    
    err := dc.syncHandler(ctx, key)
    dc.handleErr(ctx, err, key)
    
    return true
}

7.3 同一 Key 的串行化 @

虽然多个 worker 并发运行,但同一个 Deployment key 的调谐循环是串行的:队列保证同一 key 同时只有一个 worker 在处理(syncDeployment 的注释也明确说明"不允许同 key 并发调用")。

另外,syncDeployment 内部对缓存对象进行了深拷贝,避免修改 Informer 共享缓存中的对象。


八、OwnerReference 机制 @

k8s 资源之间通过 .ownerReferences 字段建立父子关系:

metadata:
  ownerReferences:
  - apiVersion: apps/v1
    kind: Deployment
    name: nginx-deployment
    uid: 1234-5678-...
    controller: true

8.1 作用 @

  • 级联删除:删除父资源时,自动删除子资源(Garbage Collector)
  • 事件传播:子资源变化时,找到父资源并触发同步
  • 孤儿管理:没有父资源的资源可以被其他父资源"收养"

8.2 ControllerRefManager @

Deployment 控制器使用 ControllerRefManager 管理 ReplicaSet 的归属:

func (dc *DeploymentController) getReplicaSetsForDeployment(ctx context.Context, d *apps.Deployment) ([]*apps.ReplicaSet, error) {
    // 1. 列出所有 RS
    rsList, _ := dc.rsLister.ReplicaSets(d.Namespace).List(labels.Everything())
    
    // 2. 使用标签选择器筛选
    cm := controller.NewReplicaSetControllerRefManager(dc.rsControl, d, deploymentSelector, controllerKind, canAdoptFunc)
    
    // 3. 认领/释放 RS
    return cm.ClaimReplicaSets(ctx, rsList)
}

ClaimReplicaSets 的逻辑:

  • 认领孤儿:没有 controller(无 controllerRef)且标签匹配的,自动收养
  • 释放不合群的:已由自己控制但标签不再匹配的,释放
  • 跳过别人的:已有其他 controller 的,不碰

九、状态管理 @

Deployment 的 Status 字段反映了当前的运行状态:

type DeploymentStatus struct {
    ObservedGeneration int64    // 已观测到的 Generation
    Replicas           int32    // 总副本数
    UpdatedReplicas    int32    // 已更新的副本数
    ReadyReplicas      int32    // 就绪副本数
    AvailableReplicas  int32    // 可用副本数
    Conditions         []DeploymentCondition  // 条件列表
}

9.1 条件(Conditions) @

条件类型 含义
Progressing 部署是否进行中
Available 是否达到最小可用副本数
ReplicaFailure 是否有副本失败
type DeploymentCondition struct {
    Type   DeploymentConditionType
    Status ConditionStatus  // True / False / Unknown
    Reason string           // 原因
    Message string          // 人类可读消息
}

9.2 进度检测 @

如果设置了 progressDeadlineSeconds,Deployment 控制器会监控部署进度:

进度超时检测不在单独的方法中,而是在 syncRolloutStatus(progress.go)更新状态时完成;requeueStuckDeployment 负责在到达 deadline 时把 Deployment 重新入队,以便再次检查:

// progress.go syncRolloutStatus
newStatus := calculateStatus(allRSs, newRS, d)

// 设置了 progressDeadlineSeconds 且本轮发布未完成时,评估进度
if util.HasProgressDeadline(d) && !isCompleteDeployment {
    switch {
    case util.DeploymentComplete(d, &newStatus):
        // 发布完成:Progressing=True, Reason=NewRSAvailable
    case util.DeploymentProgressing(d, &newStatus):
        // 有进展:Progressing=True, Reason=ReplicaSetUpdated
    case util.DeploymentTimedOut(ctx, d, &newStatus):
        // 超过 progressDeadlineSeconds 仍无进展
        msg := fmt.Sprintf("Deployment %q has timed out progressing.", d.Name)
        condition := util.NewDeploymentCondition(apps.DeploymentProgressing, v1.ConditionFalse, util.TimedOutReason, msg)
        util.SetDeploymentCondition(&newStatus, *condition)
    }
}

十、完整流程示例 @

假设用户执行 kubectl set image deployment/nginx nginx=nginx:1.20

1. API Server 更新 Deployment 的 Pod Template
   ↓
2. Deployment Informer 检测到 Update 事件
   ↓
3. 将 Deployment 加入工作队列
   ↓
4. Worker 取出 Deployment,执行 syncDeployment
   ↓
5. 发现新的 Pod Template,创建新的 ReplicaSet(RS-new)
   ↓
6. ReplicaSet 检测到 RS-new,创建新 Pod
   ↓
7. 新 Pod 就绪后,逐步缩容旧的 ReplicaSet(RS-old)
   ↓
8. RS-old 缩容到 0 后保留,用于回滚(kubectl rollout undo);cleanupDeployment 只删除超出 revisionHistoryLimit(默认 10)且 spec.replicas==0 的多余旧 RS
   ↓
9. 更新 Deployment 的 Status

十一、设计亮点总结 @

  • 声明式 API:用户只需声明期望状态,控制器负责实现,无需关心"怎么做"。
  • 异步解耦:Informer、工作队列、Worker 三层解耦,保证高并发下的稳定性。
  • 期望机制:通过记录"期望操作",避免重复创建/删除资源,保证幂等性。
  • OwnerReference:建立资源间的父子关系,实现级联删除和事件传播。
  • 缓存优先:Informer 的本地缓存大幅减少 API Server 压力,提高响应速度。
  • 限速与退避:工作队列的限速和退避机制,避免雪崩效应。

十二、其他控制器速览 @

控制器 核心逻辑
ReplicaSet 确保 Pod 副本数等于 spec.replicas
StatefulSet 为每个 Pod 分配稳定标识(如 web-0, web-1)
DaemonSet 每个节点运行一个 Pod
Job 创建 Pod 直到成功完成
CronJob 定时创建 Job
HPA 根据指标自动调整副本数
Namespace 清理命名空间下的所有资源
Node Lifecycle 监控节点状态,执行驱逐

所有控制器都遵循相同的调谐循环模式,只是业务逻辑不同。


十三、总结 @

Kubernetes 的控制器模式是整个系统的"大脑":

  • Informer 是感官,监听变化
  • Work Queue 是记忆,存储待办
  • Worker 是手脚,执行操作
  • Expectations 是反馈,确认完成

理解控制器模式,就掌握了 k8s 的运行原理。后续的调度器、网络插件、存储插件,都是与控制器协同工作的。