基于 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 资源,其核心职责是:
- 管理 ReplicaSet:创建新的 RS、淘汰旧的 RS
- 执行滚动更新:按策略逐步替换 Pod
- 维护状态:更新 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)
}
当 add 和 del 都 ≤ 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 的运行原理。后续的调度器、网络插件、存储插件,都是与控制器协同工作的。