kube-controller-manager 超深度分析


一、模块定位

1.1 业务职责

kube-controller-manager(简称 KCM)是 Kubernetes 控制面的核心组件之一,其核心职责是运行一系列控制器(Controller),每个控制器都是一个独立的控制循环(Control Loop),通过持续观测集群的"期望状态"(Desired State)和"实际状态"(Actual State),驱使两者趋于一致。

KCM 内嵌了 30+ 种控制器,覆盖:

领域 控制器
工作负载 Deployment, ReplicaSet, StatefulSet, DaemonSet, Job, CronJob, ReplicationController
网络 Endpoint, EndpointSlice, Service
节点 NodeLifecycle, NodeIPAM, CloudNodeLifecycle
弹性 HorizontalPodAutoscaler (HPA)
存储 PersistentVolumeBinder, AttachDetach, VolumeExpand, PVC/PV Protection, EphemeralVolume
安全/准入 ServiceAccount, CSRSigning, CSRApproving, CSRCleaner, BootstrapSigner, TokenCleaner
治理 GarbageCollector, Namespace, ResourceQuota, TTLAfterFinished, PodGC, Disruption
RBAC ClusterRoleAggregation
其他 RootCACertPublisher, StorageVersionGC, Route, TTL (Node)

1.2 在系统中的位置

┌─────────────────────────────────────────────────────┐
│                  Kubernetes Control Plane            │
│                                                     │
│  ┌──────────┐  ┌──────────┐  ┌──────────────────┐  │
│  │ API Server│  │ Scheduler│  │ kube-controller- │  │
│  │           │  │          │  │    manager        │  │
│  └─────┬─────┘  └────┬─────┘  └────────┬─────────┘  │
│        │              │                 │            │
│        └──────────────┴─────────────────┘            │
│                    etcd                               │
└─────────────────────────────────────────────────────┘

KCM 是 API Server 的纯客户端——它不直接访问 etcd,所有状态读写都通过 API Server 的 REST 接口完成。KCM 通过 Informer 机制以 Watch + List 方式从 API Server 获取资源变更事件,再通过 ClientSet 向 API Server 回写变更。


二、模块整体结构

2.1 入口与命令行

入口文件cmd/kube-controller-manager/controller-manager.go

func main() {
    rand.Seed(time.Now().UnixNano())
    command := app.NewControllerManagerCommand()
    logs.InitLogs()
    defer logs.FlushLogs()
    if err := command.Execute(); err != nil {
        os.Exit(1)
    }
}

入口极为精简——创建 Cobra 命令、初始化日志、执行命令。所有实质逻辑在 app.NewControllerManagerCommand() 中。

2.2 命令构建与配置

app/controllermanager.goNewControllerManagerCommand()

  1. 创建 KubeControllerManagerOptions,包含所有子控制器配置选项
  2. 构建 cobra.Command,注册所有 FlagSet
  3. Run 函数流程:s.Config()c.Complete()Run(c.Complete(), wait.NeverStop)

2.3 核心数据结构

2.3.1 Config
// app/config/config.go
type Config struct {
    ComponentConfig kubectrlmgrconfig.KubeControllerManagerConfiguration
    SecureServing   *apiserver.SecureServingInfo
    LoopbackClientConfig *restclient.Config
    InsecureServing *apiserver.DeprecatedInsecureServingInfo
    Authentication  apiserver.AuthenticationInfo
    Authorization   apiserver.AuthorizationInfo
    Client          *clientset.Clientset
    Kubeconfig      *restclient.Config
    EventRecorder   record.EventRecorder
}

Config 经 Complete() 转换为 CompletedConfig,是一个不可变配置快照。

2.3.2 ControllerContext
type ControllerContext struct {
    ClientBuilder                   clientbuilder.ControllerClientBuilder
    InformerFactory                 informers.SharedInformerFactory
    ObjectOrMetadataInformerFactory informerfactory.InformerFactory
    ComponentConfig                 kubectrlmgrconfig.KubeControllerManagerConfiguration
    RESTMapper                      *restmapper.DeferredDiscoveryRESTMapper
    AvailableResources              map[schema.GroupVersionResource]bool
    Cloud                           cloudprovider.Interface
    LoopMode                        ControllerLoopMode
    Stop                            <-chan struct{}
    InformersStarted                chan struct{}
    ResyncPeriod                    func() time.Duration
}

ControllerContext 是所有控制器的运行时上下文,在 CreateControllerContext() 中初始化,包含:

  • InformerFactory:共享的 SharedInformerFactory,所有控制器共用同一套 Informer
  • AvailableResources:API Server 支持的资源类型映射,控制器据此判断自身是否应该启动
  • ClientBuilder:构造 kubeClient 的工厂,区分 rootClient(全权限)和普通 SA client(受限权限)
  • ResyncPeriod:带 Jitter 的 Resync 周期函数,防止所有控制器同时 List
2.3.3 InitFunc 与控制器注册
type InitFunc func(ctx ControllerContext) (debuggingHandler http.Handler, enabled bool, err error)

NewControllerInitializers() 返回 map[string]InitFunc,每个控制器名映射到启动函数:

controllers["deployment"]  = startDeploymentController
controllers["replicaset"]  = startReplicaSetController
controllers["statefulset"] = startStatefulSetController
controllers["daemonset"]  = startDaemonSetController
controllers["job"]        = startJobController
controllers["cronjob"]    = startCronJobController
controllers["endpoint"]   = startEndpointController
controllers["horizontalpodautoscaling"] = startHPAController
// ... 更多控制器

2.4 Run() 主流程详解

Run(c, stopCh)
  ├── 1. 打印版本号
  ├── 2. 注册 configz
  ├── 3. 设置健康检查(含 Leader Election 健康检查)
  ├── 4. 启动 HTTP Server(Secure + Insecure)
  │     ├── NewBaseHandler() → /healthz, /metrics, /debug/pprof, /configz
  │     └── BuildHandlerChain() → Authn → Authz → RequestInfo → PanicRecovery
  ├── 5. 创建 ClientBuilder(rootClient + SA client)
  ├── 6. 定义 run 闭包
  │     ├── CreateControllerContext()
  │     │     ├── NewSharedInformerFactory()
  │     │     ├── NewSharedInformerFactory(metadata)
  │     │     ├── WaitForAPIServer()
  │     │     ├── NewDeferredDiscoveryRESTMapper()
  │     │     ├── GetAvailableResources()
  │     │     └── createCloudProvider()
  │     ├── StartControllers()
  │     │     ├── startSATokenController() [必须最先启动!]
  │     │     ├── 遍历 controllers map
  │     │     │     ├── IsControllerEnabled() 检查
  │     │     │     ├── Jitter 间隔启动
  │     │     │     └── initFn(ctx) → 创建并运行控制器
  │     ├── InformerFactory.Start()
  │     ├── ObjectOrMetadataInformerFactory.Start()
  │     └── close(InformersStarted)
  ├── 7. 判断 Leader Election
  │     ├── 不启用 → 直接 run()
  │     └── 启用 → leaderElectAndRun()
  │           ├── 构造 ResourceLock (ConfigMap/Lease)
  │           └── leaderelection.RunOrDie()
  │                 ├── OnStartedLeading → run() 启动所有控制器
  │                 └── OnStoppedLeading → Fatal
  └── 8. Leader Migration(可选)
        ├── 主锁:启动非迁移控制器
        ├── 等待 MigrationReady
        └── 迁移锁:启动迁移控制器

2.5 Leader Election 机制

KCM 支持多实例部署,通过 Leader Election 保证同一时间只有一个实例活跃执行控制循环:

func leaderElectAndRun(c, lockIdentity, electionChecker, resourceLock, leaseName, callbacks) {
    rl := resourcelock.NewFromKubeconfig(...)  // ConfigMapLock 或 LeaseLock
    leaderelection.RunOrDie(context.TODO(), LeaderElectionConfig{
        Lock:          rl,
        LeaseDuration: 配置值,
        RenewDeadline: 配置值,
        RetryPeriod:   配置值,
        Callbacks:     callbacks,  // OnStartedLeading, OnStoppedLeading
        WatchDog:      electionChecker,
    })
}

Leader Migration:允许将部分控制器从 KCM 迁移到 cloud-controller-manager 运行,通过两把 Leader Lock 实现:

  • 主锁:运行非迁移控制器
  • 迁移锁:运行迁移控制器

2.6 控制器通用架构模式

所有控制器共享同一架构模式:

Informer → EventHandler → WorkQueue → Worker → syncHandler → API Server
  1. Informer:通过 SharedInformerFactory 获取特定资源的 Informer
  2. EventHandler:注册 Add/Update/Delete 回调,将资源 Key 入队
  3. WorkQueue:RateLimitingInterface,带去重和退避重试
  4. Worker:从队列取出 Key,调用 syncHandler
  5. syncHandler:执行核心协调逻辑

2.7 Expectations 机制

controller/controller_utils.go 中定义的 Expectations 是控制器协调的关键优化:

type ControlleeExpectations struct {
    add int32  // 预期创建数
    del int32  // 预期删除数
    key string // 控制器 key
}
  • 控制器在创建/删除 Pod 前,先设置 Expectations(如 ExpectCreations(key, 5)
  • 每当 Informer 观察到 Pod 创建/删除事件,对应 Expectation 被减少(CreationObserved
  • SatisfiedExpectations(key) 判断所有预期是否已满足或超时(5分钟)
  • 未满足则跳过同步,避免重复操作

2.8 ControllerRef 管理

controller/controller_ref_manager.go 实现了通用的 Adopt/Release 机制:

type BaseControllerRefManager struct {
    Controller   metav1.Object
    Selector     labels.Selector
    CanAdoptFunc func() error  // 延迟检查,防止已被删除的控制器收养
}

func (m *BaseControllerRefManager) ClaimObject(obj, match, adopt, release) (bool, error) {
    // 1. 有 OwnerRef 且 UID 匹配 → 已拥有
    // 2. 有 OwnerRef 但不匹配 → 忽略(别人的)
    // 3. 有 OwnerRef 且匹配但 Selector 不匹配 → Release
    // 4. 无 OwnerRef 且匹配 → Adopt
    // 5. 无 OwnerRef 且不匹配 → 忽略
}

衍生类型:PodControllerRefManagerReplicaSetControllerRefManagerControllerRevisionControllerRefManager,分别管理 Pod、ReplicaSet、ControllerRevision 的 OwnerReference。


三、核心业务逻辑深度解析

3.1 Controller Manager 启动流程

No

Yes

No

Yes

Yes

main

NewControllerManagerCommand

s.Config 生成 Config

c.Complete 完成化

Run

LeaderElect?

直接 run

leaderElectAndRun

构造 ResourceLock

RunOrDie 选举

OnStartedLeading

CreateControllerContext

StartControllers

startSATokenController 最先

遍历 controllers map

IsControllerEnabled?

跳过

Jitter 间隔

initFn 创建控制器

InformerFactory.Start

close InformersStarted

select 阻塞

Leader Migration?

主锁: 非迁移控制器

等待 MigrationReady

迁移锁: 迁移控制器

3.2 Informer 机制架构

Controllers

SharedInformerFactory

API Server

Watch + List

distribute

Register Handler

Register Handler

Register Handler

enqueue key

dequeue

dequeue

dequeue

REST API

Reflector

DeltaFIFO

Store/Indexer

Controller Loop

Event Handlers

DeploymentController

ReplicaSetController

StatefulSetController

WorkQueue

核心流程

  1. Reflector 通过 List 获取全量数据,再通过 Watch 持续监听变更
  2. 变更事件封装为 Delta 对象,写入 DeltaFIFO
  3. Controller Loop 从 DeltaFIFO Pop,更新 Store(本地缓存)
  4. 同时触发注册在 Informer 上的 EventHandler
  5. EventHandler 将资源 Key 放入控制器的 WorkQueue
  6. Worker 从 WorkQueue 取出 Key,执行 syncHandler

SharedInformerFactory 保证同一资源类型只创建一个 Informer,多个控制器共享同一份缓存。

3.3 Informer-Reflector-Store 数据流

Worker WorkQueue EventHandler Store DeltaFIFO API Server Reflector Worker WorkQueue EventHandler Store DeltaFIFO API Server Reflector loop [Watch 循环] loop [Pop 循环] loop [Worker 循环] List 获取全量 资源列表 + ResourceVersion 批量入队 Sync/Delta Watch(ResourceVersion) ADDED/MODIFIED/DELETED 事件 入队对应 Delta Pop Delta Add/Update/Delete 触发 OnAdd/OnUpdate/OnDelete AddRateLimited(key) Get(key) syncHandler(key) Lister.Get 读取缓存 Create/Update/Delete Done(key)

3.4 Controller 总体协调架构

API Server

Kube-Controller-Manager

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

Read/Write

SharedInformerFactory

DeploymentController

ReplicaSetController

StatefulSetController

DaemonSetController

JobController

CronJobController

EndpointController

HorizontalController

NodeLifecycleController

NodeIPAMController

NamespaceController

ServiceAccountController

GarbageCollector

ResourceQuotaController

REST API + etcd


3.5 Deployment Controller 深度解析

3.5.1 类结构
type DeploymentController struct {
    rsControl       controller.RSControlInterface
    client          clientset.Interface
    eventRecorder   record.EventRecorder
    syncHandler     func(dKey string) error
    enqueueDeployment func(deployment *apps.Deployment)
    dLister         appslisters.DeploymentLister
    rsLister        appslisters.ReplicaSetLister
    podLister       corelisters.PodLister
    dListerSynced   cache.InformerSynced
    rsListerSynced  cache.InformerSynced
    podListerSynced cache.InformerSynced
    queue           workqueue.RateLimitingInterface
}
3.5.2 事件注册
dInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
    AddFunc:    dc.addDeployment,     // Deployment 创建 → 入队
    UpdateFunc: dc.updateDeployment,  // Deployment 更新 → 入队
    DeleteFunc: dc.deleteDeployment,  // Deployment 删除 → 入队
})
rsInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
    AddFunc:    dc.addReplicaSet,     // RS 创建 → 反查所属 Deployment 入队
    UpdateFunc: dc.updateReplicaSet,  // RS 更新 → 反查所属 Deployment 入队
    DeleteFunc: dc.deleteReplicaSet,  // RS 删除 → 反查所属 Deployment 入队
})
podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
    DeleteFunc: dc.deletePod,         // Pod 删除(仅 Recreate 策略) → 入队
})
3.5.3 syncDeployment 完整流程
func (dc *DeploymentController) syncDeployment(key string) error {
    // 1. 解析 namespace/name
    // 2. 从 Lister 获取 Deployment
    // 3. DeepCopy 避免修改缓存
    // 4. 检查空 Selector → 警告
    // 5. getReplicaSetsForDeployment() → 通过 ControllerRefManager 收养/释放 RS
    // 6. getPodMapForDeployment() → 获取所有 Pod,按 RS UID 分组
    // 7. 检查 DeletionTimestamp → 仅同步状态
    // 8. checkPausedConditions() → 检查暂停条件
    // 9. if Paused → sync()
    // 10. if RollbackTo → rollback()
    // 11. if ScalingEvent → sync()
    // 12. 根据 Strategy 分发:
    //     - Recreate → rolloutRecreate()
    //     - RollingUpdate → rolloutRolling()
}
3.5.4 Deployment Controller 流程图

Yes

No

Yes

No

Yes

No

Yes

No

Yes

No

Recreate

RollingUpdate

syncDeployment key

获取 Deployment

DeepCopy

Selector 为空?

警告 + 更新 Status

getReplicaSetsForDeployment

getPodMapForDeployment

DeletionTimestamp != nil?

syncStatusOnly

checkPausedConditions

Paused?

sync 比例扩缩

RollbackTo != nil?

rollback

isScalingEvent?

Strategy Type

rolloutRecreate

rolloutRolling

3.5.5 Rolling Update 流程

rolloutRolling() 的核心逻辑(rolling.go):

  1. 获取 newRS 和 allOldRSs
  2. 如果 newRS 不存在 → 创建 newRS(带 pod-template-hash)
  3. 缩容 old RS(按 MaxUnavailable 计算)
  4. 扩容 new RS(按 MaxSurge 计算)
  5. 清理旧 RS(超过 RevisionHistoryLimit 则删除)
  6. 同步 Deployment Status
3.5.6 Recreate 流程

rolloutRecreate() 的核心逻辑(recreate.go):

  1. 缩容所有 old RS 到 0
  2. 等待所有 Pod 终止
  3. 扩容 new RS 到期望副本数
  4. 与 Rolling 不同,Recreate 是先全部停再全部启,有中断窗口
3.5.7 比例扩缩(Proportional Scaling)

scale() 方法在存在多个活跃 RS 时,按比例分配副本数:

// 计算 allowedSize = replicas + MaxSurge
// deploymentReplicasToAdd = allowedSize - 当前总副本
// 按比例分配到各 RS:新 RS 优先扩,旧 RS 优先缩
proportion := deploymentutil.GetProportion(rs, deployment, replicasToAdd, replicasAdded)

3.6 ReplicaSet Controller 深度解析

3.6.1 类结构
type ReplicaSetController struct {
    schema.GroupVersionKind
    kubeClient    clientset.Interface
    podControl    controller.PodControlInterface
    burstReplicas int    // 突发上限,默认 500
    syncHandler   func(rsKey string) error
    expectations  *controller.UIDTrackingControllerExpectations
    rsLister      appslisters.ReplicaSetLister
    podLister     corelisters.PodLister
    queue         workqueue.RateLimitingInterface
}
3.6.2 syncReplicaSet 完整流程
func (rsc *ReplicaSetController) syncReplicaSet(key string) error {
    // 1. 获取 RS
    // 2. SatisfiedExpectations() 检查
    // 3. List 所有 Pod → FilterActivePods → claimPods (Adopt/Release)
    // 4. if rsNeedsSync && !deleting → manageReplicas()
    // 5. calculateStatus() 计算新状态
    // 6. updateReplicaSetStatus() 更新状态
    // 7. MinReadySeconds 后重新入队(处理 Available 延迟)
}
3.6.3 manageReplicas 核心逻辑
func (rsc *ReplicaSetController) manageReplicas(filteredPods []*v1.Pod, rs *apps.ReplicaSet) error {
    diff := len(filteredPods) - int(*(rs.Spec.Replicas))
    
    if diff < 0 {
        // Pod 不够 → 创建
        // 1. 限制 burstReplicas
        // 2. ExpectCreations()
        // 3. slowStartBatch() 创建 Pod
        //    批次大小: 1, 2, 4, 8, ... 指数增长
        //    任一批次失败则停止后续批次
    } else if diff > 0 {
        // Pod 过多 → 删除
        // 1. 限制 burstReplicas
        // 2. ExpectDeletions()
        // 3. getPodsToDelete() → 排序选择待删 Pod
        //    排序策略: 同节点相关 Pod 数少的优先删
        // 4. 并发删除
    }
}
3.6.4 slowStartBatch
func slowStartBatch(count int, initialBatchSize int, fn func() error) (int, error) {
    // 批次大小从 initialBatchSize 开始,每批成功后翻倍
    // 1, 2, 4, 8, 16, ...
    // 任一批次有失败则停止
    // 返回成功数
}

这是关键的防雪崩机制:避免一次性创建大量 Pod 导致 API Server 过载或配额耗尽后大量失败事件。

3.6.5 ReplicaSet Controller 流程图

Yes

No

Yes

No

Yes

No

Yes

No

Yes

No

syncReplicaSet key

获取 ReplicaSet

RS 已删除?

DeleteExpectations + return

SatisfiedExpectations?

List Pods

FilterActivePods

claimPods: Adopt/Release

rsNeedsSync && not deleting?

manageReplicas

跳过 Pod 管理

diff < 0?

slowStartBatch 创建 Pod

diff > 0?

选择并删除 Pod

无需操作

calculateStatus

updateReplicaSetStatus

MinReadySeconds > 0?

延时重新入队

完成


3.7 StatefulSet Controller 深度解析

3.7.1 类结构
type StatefulSetController struct {
    kubeClient    clientset.Interface
    control       StatefulSetControlInterface    // 核心控制逻辑抽象
    podControl    controller.PodControlInterface
    podLister     corelisters.PodLister
    setLister     appslisters.StatefulSetLister
    pvcListerSynced cache.InformerSynced
    revListerSynced cache.InformerSynced
    queue         workqueue.RateLimitingInterface
}

StatefulSet 的控制逻辑被抽象到 StatefulSetControlInterface,默认实现为 DefaultStatefulSetControl,组合了:

  • StatefulPodControlInterface:Pod/PVC 创建删除
  • StatefulSetStatusUpdaterInterface:状态更新
  • HistoryInterface:ControllerRevision 管理(用于滚动更新)
3.7.2 StatefulSet 的核心特征

与 Deployment/ReplicaSet 的关键区别:

  1. 有序性:Pod 按序号(0, 1, 2, …)创建/删除,扩容顺序创建,缩容逆序删除
  2. 稳定的网络标识:Pod 名为 <statefulset-name>-<ordinal>,DNS 可预测
  3. 稳定存储:每个 Pod 绑定独立 PVC,Pod 重建后自动重新挂载同一 PVC
  4. 滚动更新:逆序逐个更新(从最大序号开始),支持 Partition 更新
3.7.3 StatefulSet Controller 流程图

Yes

No

Yes

No

Yes

No

sync StatefulSet

获取 StatefulSet

DeletionTimestamp?

return

获取 ControllerRevisions

获取 Pods + PVCs

按序号排序 Pods

Replicas 不匹配?

按序号创建/删除 Pod

创建对应 PVC

等待 Pod Ready

继续下一个序号

RollingUpdate?

逆序更新 Pod

删除旧 Pod → 创建新 Pod

等待 Ready → 继续

更新 Status

清理旧 ControllerRevisions


3.8 DaemonSet Controller 深度解析

3.8.1 类结构
type DaemonSetsController struct {
    kubeClient    clientset.Interface
    eventRecorder record.EventRecorder
    podControl    controller.PodControlInterface
    crControl     controller.ControllerRevisionControlInterface
    burstReplicas int   // 250
    expectations  controller.ControllerExpectationsInterface
    dsLister      appslisters.DaemonSetLister
    historyLister appslisters.ControllerRevisionLister
    podLister     corelisters.PodLister
    podNodeIndex  cache.Indexer   // 按 nodeName 索引
    nodeLister    corelisters.NodeLister
    queue         workqueue.RateLimitingInterface
    failedPodsBackoff *flowcontrol.Backoff
}
3.8.2 核心特征
  • 每个 Node 最多一个 Pod:与 Deployment/RS 的"总副本数"模式完全不同
  • Node 亲和性 + Taint/Toleration:DaemonSet Pod 只调度到满足条件的 Node
  • 滚动更新:按 Node 逐个更新,支持 maxUnavailable
3.8.3 核心协调流程
  1. 获取所有 Node
  2. 对每个 Node,检查是否应该运行 Daemon Pod(NodeSelector、Taint/Toleration 等)
  3. 如果应该运行但没有 Pod → 创建
  4. 如果不应该运行但有 Pod → 删除
  5. 如果 Pod 存在但模板不匹配 → 根据 UpdateStrategy 决定是否更新
3.8.4 DaemonSet Controller 流程图

No

Yes

No

Yes

No

Yes

No

Yes

No

OnDelete

RollingUpdate

Yes

syncDaemonSet key

获取 DaemonSet

获取所有 Nodes

获取现有 Daemon Pods

遍历每个 Node

Node 满足调度条件?

Node 上有 Daemon Pod?

删除 Pod

跳过

Node 上有 Daemon Pod?

SatisfiedExpectations?

创建 Pod

跳过等待

Pod 模板匹配?

UpdateStrategy

等待手动删除

按 maxUnavailable 更新

无需操作

计算 Status

cleanupHistory

updateDaemonSetStatus


3.9 Job Controller 深度解析

3.9.1 类结构
type Controller struct {
    kubeClient    clientset.Interface
    podControl    controller.PodControlInterface
    updateHandler func(job *batch.Job) error
    syncHandler   func(jobKey string) (bool, error)
    expectations  controller.ControllerExpectationsInterface
    jobLister     batchv1listers.JobLister
    podStore      corelisters.PodLister
    queue         workqueue.RateLimitingInterface
    recorder      record.EventRecorder
}
3.9.2 syncJob 核心流程
func (jm *Controller) syncJob(key string) (bool, error) {
    // 1. 获取 Job
    // 2. 如果已完成 → return
    // 3. 检查 FeatureGate (IndexedJob)
    // 4. SatisfiedExpectations()
    // 5. getPodsForJob() → Adopt/Release
    // 6. 统计 active/succeeded/failed
    // 7. 设置 StartTime(首次)
    // 8. 判断失败条件:
    //    - BackoffLimitExceeded(失败次数超限)
    //    - DeadlineExceeded(ActiveDeadlineSeconds 超时)
    // 9. 如果失败 → 删除所有 Active Pod,设置 Failed Condition
    // 10. 否则 → manageJob()
    // 11. 判断完成条件:
    //     - completions == nil: 任一 Pod 成功即完成
    //     - completions != nil: 成功数 >= completions
    // 12. 更新 Job Status
}
3.9.3 manageJob
func (jm *Controller) manageJob(job, activePods, succeeded, pods) (int32, error) {
    // 根据 Parallelism 和 Completions 计算需要的 Active Pod 数
    // active < wantActive → 创建 Pod(slowStartBatch)
    // active > wantActive → 删除多余 Pod(优先删除 Failed/Pending)
}
3.9.4 Job/CronJob Controller 流程图

Yes

No

Yes

No

Yes

No

Yes

No

Yes

Yes

No

Yes

syncJob key

获取 Job

Job 已完成?

return forget

getPodsForJob

统计 active/succeeded/failed

设置 StartTime

BackoffLimit 超限?

标记 Failed + 删除 Active Pods

ActiveDeadline 超时?

SatisfiedExpectations?

manageJob

跳过

Completions == nil?

succeeded > 0 && active == 0?

标记 Complete

succeeded >= completions?

更新 Status

3.9.5 CronJob Controller

CronJob Controller 采用不同于其他控制器的轮询模式(不使用 Informer + WorkQueue):

func (jm *Controller) Run(stopCh <-chan struct{}) {
    go wait.Until(jm.syncAll, 10*time.Second, stopCh)
}

每 10 秒全量 List CronJob 和 Job,计算哪些 Job 应该被创建/清理。V2 版本 (NewControllerV2) 切换为标准的 Informer + WorkQueue 模式。


3.10 HPA (Autoscaler) Controller 深度解析

3.10.1 类结构
type HorizontalController struct {
    scaleNamespacer  scaleclient.ScalesGetter
    hpaNamespacer    autoscalingclient.HorizontalPodAutoscalersGetter
    mapper           apimeta.RESTMapper
    replicaCalc      *ReplicaCalculator
    eventRecorder    record.EventRecorder
    downscaleStabilisationWindow time.Duration
    hpaLister        autoscalinglisters.HorizontalPodAutoscalerLister
    podLister        corelisters.PodLister
    queue            workqueue.RateLimitingInterface
    recommendations  map[string][]timestampedRecommendation
    scaleUpEvents    map[string][]timestampedScaleEvent
    scaleDownEvents  map[string][]timestampedScaleEvent
}
3.10.2 HPA 协调流程
reconcileAutoscaler(hpa)
  ├── 1. 获取 Scale Sub-resource(Deployment/RS 的 Scale API)
  ├── 2. 计算当前 Pod 指标
  │     ├── Resource 指标(CPU/Memory)
  │     ├── Pods 指标(自定义 Pod 指标)
  │     ├── Object 指标(其他对象指标)
  │     ├── External 指标(外部系统指标)
  │     └── ContainerResource 指标
  ├── 3. 计算期望副本数
  │     desiredReplicas = ceil[currentUtilization / targetUtilization * currentReplicas]
  ├── 4. 应用缩容稳定窗口(downscaleStabilisationWindow)
  │     取窗口内最大推荐值,防止震荡
  ├── 5. 应用缩放限制
  │     scaleUpLimit = max(currentReplicas * scaleUpLimitFactor, scaleUpLimitMinimum)
  │     不能一次扩容超过 scaleUpLimit
  │     缩容不能低于 minReplicas
  ├── 6. 更新 Scale Sub-resource
  └── 7. 更新 HPA Status
3.10.3 ReplicaCalculator
type ReplicaCalculator struct {
    metricsClient         metrics.MetricsClient
    podLister             corelisters.PodLister
    tolerance             float64    // 默认 0.1,即 10% 容差
    cpuInitializationPeriod time.Duration
    delayOfInitialReadinessStatus time.Duration
}

核心公式:

desiredReplicas = ceil[ metricValue / targetMetricValue * currentReplicas ]

如果 |metricValue/targetMetricValue - 1.0| < tolerance,则不进行缩放。

3.10.4 HPA Controller 流程图

Resource

Pods

Object

External

ContainerResource

Yes

No

reconcileAutoscaler hpa

获取 Scale 对象

computeReplicasForMetrics

Metric 类型

计算 CPU/Memory 利用率

计算自定义 Pod 指标

计算对象指标

计算外部指标

计算容器级指标

取所有指标的最大推荐副本数

应用缩容稳定窗口

应用扩容速率限制

应用 minReplicas/maxReplicas 限制

desired != current?

更新 Scale Sub-resource

无需缩放

更新 HPA Status + Conditions


3.11 Service/Endpoint Controller 深度解析

3.11.1 类结构
type Controller struct {
    client           clientset.Interface
    serviceLister    corelisters.ServiceLister
    podLister        corelisters.PodLister
    endpointsLister  corelisters.EndpointsLister
    queue            workqueue.RateLimitingInterface
    triggerTimeTracker *endpointutil.TriggerTimeTracker
    endpointUpdatesBatchPeriod time.Duration
    serviceSelectorCache *endpointutil.ServiceSelectorCache
}
3.11.2 核心逻辑

Endpoint Controller 维护 Service → Endpoints 的映射:

  1. 监听 Service Add/Update/Delete → 入队
  2. 监听 Pod Add/Update/Delete → 查找匹配的 Service → 入队
  3. 监听 Endpoints Delete → 入队(重建)
  4. 同步时:
    • 获取 Service 的 Selector
    • List 匹配的 Ready Pod
    • 构建 Endpoints 对象(IP:Port 列表)
    • Create 或 Update Endpoints

关键细节:

  • Pod 必须是 Ready 才会加入 Endpoints(除非 Service 设置了 publishNotReadyAddresses
  • 支持 TolerateUnreadyEndpointsAnnotation(StatefulSet 使用)
  • 批量更新:endpointUpdatesBatchPeriod 减少 API Server 写入压力
  • triggerTimeTracker:记录 Endpoints 最后变更时间,用于 EndpointsLastChangeTriggerTime 注解
3.11.3 Service/Endpoint Controller 流程图

No

Yes

No

Yes

Yes

No

Service 变更事件

入队 Service Key

Pod 变更事件

查找匹配 Service

Endpoints 被删除

syncService Key

获取 Service

Service 有 Selector?

清理 Endpoints

List 匹配 Pod

过滤 Ready Pod

构建 Endpoints Subsets

Endpoints 已存在?

Create Endpoints

内容有变化?

Update Endpoints

无需更新


3.12 Node Controller 深度解析

3.12.1 类结构
type Controller struct {
    taintManager    *scheduler.NoExecuteTaintManager
    podLister       corelisters.PodLister
    kubeClient      clientset.Interface
    knownNodeSet    map[string]*v1.Node
    nodeHealthMap   *nodeHealthMap
    nodeEvictionMap *nodeEvictionMap
    zonePodEvictor  map[string]*scheduler.RateLimitedTimedQueue
    zoneNoExecuteTainter map[string]*scheduler.RateLimitedTimedQueue
    zoneStates      map[string]ZoneState
    leaseLister     coordlisters.LeaseLister
    nodeLister      corelisters.NodeLister
    daemonSetStore  appsv1listers.DaemonSetLister
    nodeMonitorPeriod        time.Duration
    nodeStartupGracePeriod   time.Duration
    nodeMonitorGracePeriod   time.Duration
    podEvictionTimeout       time.Duration
    evictionLimiterQPS       float32
    secondaryEvictionLimiterQPS float32
    largeClusterThreshold    int32
    unhealthyZoneThreshold   float32
    runTaintManager          bool
    nodeUpdateQueue          workqueue.Interface
    podUpdateQueue           workqueue.RateLimitingInterface
}
3.12.2 核心职责
  1. 节点健康监控:持续监控 Node Status 和 Node Lease
  2. Taint 管理:根据节点状态添加/移除 Taint
    • NodeReady=Falsenode-not-ready Taint
    • NodeReady=Unknownnode-unreachable Taint
    • MemoryPressure=Truememory-pressure Taint
    • DiskPressure=Truedisk-pressure Taint
    • NetworkUnavailable=Truenetwork-unavailable Taint
    • PIDPressure=Truepid-pressure Taint
  3. Pod 驱逐:对 NotReady 超时的节点驱逐 Pod
  4. Zone 级别降级保护
    • 全部 Node NotReady(FullDisruption)→ 降低驱逐速率
    • 部分 Node NotReady(PartialDisruption)→ 使用二级驱逐速率
    • 正常状态 → 使用正常驱逐速率
  5. Label 协调:同步 beta/GA 标签
3.12.3 节点健康判断
probeTimestamp = max(LastTransitionTime, Lease.RenewTime)
now - probeTimestamp > nodeMonitorGracePeriod → Unknown

支持两种健康信号:

  • NodeStatus:kubelet 定期上报 Node Status
  • NodeLease:kubelet 定期续约 Lease(更轻量)

取两者最新的时间戳作为健康探测时间。

3.12.4 Node Controller 流程图

Yes

No

Yes

No

Normal

PartialDisruption

FullDisruption

monitorNodeHealth 循环

遍历所有 Node

获取 nodeHealthData

probeTimestamp 过期?

设置 NodeReady=Unknown

保持现有状态

添加 Unreachable Taint

NodeReady=False?

添加 NotReady Taint

移除 Taint

计算 Zone 状态

Zone 状态

正常驱逐速率

二级驱逐速率

0 速率暂停驱逐

zonePodEvictor

不驱逐

驱逐 NotReady 节点上的 Pod

tryUpdateNodeStatus

更新 Node Status Conditions


3.13 控制器启动注册表

NewControllerInitializers() 完整映射:

控制器名 启动函数 配置来源 并发数配置
endpoint startEndpointController core.go ConcurrentEndpointSyncs
endpointslice startEndpointSliceController discovery.go -
endpointslicemirroring startEndpointSliceMirroringController discovery.go -
replicationcontroller startReplicationController core.go ConcurrentRCSyncs
podgc startPodGCController core.go TerminatePodGCThreshold
resourcequota startResourceQuotaController core.go ConcurrentResourceQuotaSyncs
namespace startNamespaceController core.go ConcurrentNamespaceSyncs
serviceaccount startServiceAccountController core.go -
garbagecollector startGarbageCollectorController core.go ConcurrentGCSyncs
daemonset startDaemonSetController apps.go ConcurrentDaemonSetSyncs
job startJobController batch.go ConcurrentJobSyncs
deployment startDeploymentController apps.go ConcurrentDeploymentSyncs
replicaset startReplicaSetController apps.go ConcurrentRSSyncs
horizontalpodautoscaling startHPAController autoscaling.go -
disruption startDisruptionController policy.go -
statefulset startStatefulSetController apps.go ConcurrentStatefulSetSyncs
cronjob startCronJobController batch.go ConcurrentCronJobSyncs
nodelifecycle startNodeLifecycleController core.go -
nodeipam startNodeIpamController core.go -
persistentvolume-binder startPersistentVolumeBinderController core.go -
attachdetach startAttachDetachController core.go -
persistentvolume-expander startVolumeExpandController core.go -
clusterrole-aggregation startClusterRoleAggregrationController rbac.go -
pvc-protection startPVCProtectionController core.go -
pv-protection startPVProtectionController core.go -
ttl-after-finished startTTLAfterFinishedController core.go ConcurrentTTLSyncs
root-ca-cert-publisher startRootCACertPublisher core.go -
ephemeral-volume startEphemeralVolumeController core.go -
bootstrapsigner startBootstrapSignerController bootstrap.go (默认禁用)
tokencleaner startTokenCleanerController bootstrap.go (默认禁用)
service startServiceController core.go (cloud) ConcurrentServiceSyncs
route startRouteController core.go (cloud) -
cloud-node-lifecycle startCloudNodeLifecycleController core.go (cloud) -
storage-version-gc startStorageVersionGCController core.go -

四、关键设计模式与机制深度解析

4.1 SharedInformerFactory

// 所有控制器共享同一个 InformerFactory
sharedInformers := informers.NewSharedInformerFactory(versionedClient, ResyncPeriod(s)())
  • 同一资源类型只创建一个 Informer,多个控制器共享
  • ResyncPeriod 带 Jitter:factor := rand.Float64() + 1,防止所有控制器同时 List
  • InformerFactory.Start() 启动所有 Informer 的 Reflector
  • InformersStarted channel 在所有控制器初始化后关闭,某些控制器(如 GarbageCollector)需要等待此信号

4.2 RateLimitingQueue

所有控制器使用 workqueue.NewNamedRateLimitingQueue

queue := workqueue.NewNamedRateLimitingQueue(
    workqueue.DefaultControllerRateLimiter(),  // 5ms * 2^(retry-1), max 1000s
    "controller-name"
)

特性:

  • 去重:同一 Key 在处理中时,新入队请求合并
  • 限速重试:指数退避 + 随机 Jitter
  • Forget:成功处理后清除重试计数
  • ShutDown:优雅关闭,等待进行中的工作完成

4.3 Expectations 的深层设计

Expectations 解决的核心问题是避免控制器的热循环(Hot Loop)

场景:RS 期望 3 个 Pod,当前只有 1 个,控制器创建了 2 个 Pod。

  • 无 Expectations:创建后立即再次 sync → 发现 Pod 数仍为 1(Informer 还没收到事件)→ 又创建 2 个 → 重复 → 最终可能创建大量多余 Pod
  • 有 Expectations:创建后期望 add=2,sync 时发现 Expectations 未满足 → 跳过 → Informer 收到创建事件 → CreationObservedadd 减为 0 → 下次 sync 才继续

超时机制:5 分钟(ExpectationsTimeout)后自动视为满足,防止 Watch 丢事件导致控制器永久卡住。

4.4 Slow Start Batch

func slowStartBatch(count, initialBatchSize int, fn func() error) (int, error) {
    remaining := count
    successes := 0
    for batchSize := min(remaining, initialBatchSize); batchSize > 0; batchSize = min(2*batchSize, remaining) {
        // 并发执行 batchSize 次 fn
        // 任一失败 → 返回已成功数 + 错误
        successes += curSuccesses
        remaining -= batchSize
    }
}

为什么初始 batchSize=1?

  • 如果配额不足,第 1 个就会失败
  • 避免一次性提交大量请求后全部失败,产生大量 Event 垃圾
  • 成功后指数增长:1 → 2 → 4 → 8 → …,快速追赶

4.5 ControllerRef (OwnerReference) 管理

type OwnerReference struct {
    APIVersion         string
    Kind               string
    Name               string
    UID                types.UID
    Controller         *bool  // 标记为控制器引用
    BlockOwnerDeletion *bool  // 是否阻止 GC 删除被控对象
}

Adopt/Release 流程:

  1. Adopt:通过 Patch 添加 OwnerReference 到孤儿对象
  2. Release:通过 StrategicMergePatch 移除 OwnerReference
  3. CanAdoptFunc:延迟检查控制器是否仍然存在且未被删除,防止竞争条件

RecheckDeletionTimestamp 确保:

  • List RS 后、Adopt 前,如果 Deployment 已被删除,则不收养
  • 防止已删除的控制器收养新对象,导致 GC 无法清理

4.6 Hash Collision 处理

Deployment 创建 RS 时使用 pod-template-hash 标签:

podTemplateSpecHash := controller.ComputeHash(&newRSTemplate, d.Status.CollisionCount)

如果哈希冲突(不同模板得到相同哈希):

  1. 尝试 Create → AlreadyExists
  2. 检查已有 RS 是否属于同一 Deployment 且模板相同
  3. 如果不同 → CollisionCount++ → 重新入队 → 下次用新 CollisionCount 计算不同哈希

4.7 Leader Migration

type LeaderMigrator struct {
    MigrationReady chan struct{}
    FilterFunc     FilterFunc
}

FilterFunc 返回三种结果:

  • ControllerNonMigrated:在主锁中运行
  • ControllerMigrated:在迁移锁中运行
  • ControllerUnowned:不在本实例中运行

迁移流程:

  1. 主锁获取 Leader → 启动非迁移控制器 + SA Token Controller
  2. SA Token Controller 启动后关闭 MigrationReady
  3. 迁移锁获取 Leader → 启动迁移控制器

五、Mermaid 图总览

图1:Controller Manager 启动流程图

(见 3.1 节)

图2:Informer 机制架构图

(见 3.2 节)

图3:Controller 总体协调架构图

(见 3.4 节)

图4:Deployment Controller 流程图

(见 3.5.4 节)

图5:ReplicaSet Controller 流程图

(见 3.6.5 节)

图6:StatefulSet Controller 流程图

(见 3.7.3 节)

图7:DaemonSet Controller 流程图

(见 3.8.4 节)

图8:Job/CronJob Controller 流程图

(见 3.9.4 节)

图9:HPA Controller 流程图

(见 3.10.4 节)

图10:Service/Endpoint Controller 流程图

(见 3.11.3 节)

图11:Node Controller 流程图

(见 3.12.4 节)

图12:Informer-Reflector-Store 数据流图

(见 3.3 节)


六、数据流入流出方式

6.1 数据流入

来源 方式 目标
API Server Informer (List+Watch) 本地 Store (Indexer)
Cloud Provider Cloud Interface Node/Route/Service 控制器
Metrics Server REST Client HPA Controller

6.2 数据流出

目标 方式 来源
API Server ClientSet (Create/Update/Delete/Patch) 所有控制器
Event EventRecorder → API Server 所有控制器
Metrics Prometheus workqueue, ratelimiter
Healthz HTTP healthz checker

6.3 控制器间数据流

CronJob → Job → Pod
Deployment → ReplicaSet → Pod
StatefulSet → Pod + PVC
DaemonSet → Pod
HPA → Scale (Deployment/RS) → ReplicaSet → Pod
Service → Endpoints ← Pod
NodeLifecycle → Taint → Pod Eviction

七、关键配置参数

参数 默认值 说明
–min-resync-period 12h Informer 最小 Resync 间隔
–controller-start-interval 0s 控制器启动间隔
–leader-elect true 是否启用 Leader Election
–leader-elect-lease-duration 15s Leader 租约时长
–leader-elect-renew-deadline 10s Leader 续约截止时间
–leader-elect-retry-period 2s Leader 选举重试间隔
–concurrent-deployment-syncs 5 Deployment 并发同步数
–concurrent-replicaset-syncs 5 ReplicaSet 并发同步数
–concurrent-statefulset-syncs 5 StatefulSet 并发同步数
–concurrent-daemonset-syncs 2 DaemonSet 并发同步数
–concurrent-job-syncs 5 Job 并发同步数
–concurrent-endpoint-syncs 5 Endpoint 并发同步数
–node-monitor-grace-period 40s Node 健康监测宽限期
–pod-eviction-timeout 5m Pod 驱逐超时
–horizontal-pod-autoscaler-sync-period 15s HPA 同步周期
–use-service-account-credentials false 是否使用 SA 凭证

八、错误处理与容错

8.1 WorkQueue 重试

所有控制器通过 WorkQueue 实现自动重试:

  • 成功 → Forget(key)
  • 失败 → AddRateLimited(key) → 指数退避重试
  • 超过 maxRetries → 放弃并记录错误

8.2 Expectations 超时

5 分钟未满足的 Expectations 自动视为满足,防止 Watch 丢事件导致控制器永久卡住。

8.3 Status Update 重试

for i := 0; ; i++ {
    rs.Status = newStatus
    updatedRS, err = c.UpdateStatus(...)
    if err == nil { return }
    if i >= statusUpdateRetries { break }
    rs, getErr = c.Get(...)  // 重新获取最新版本
}

8.4 Panic Recovery

defer utilruntime.HandleCrash()

每个 Worker 启动时注册 panic 恢复,防止单个 sync 失败导致整个控制器崩溃。

8.5 Graceful Shutdown

defer dc.queue.ShutDown()

控制器关闭时:

  1. 停止接受新工作
  2. 等待进行中的工作完成
  3. 关闭 Informer

九、性能优化设计

9.1 共享 Informer

避免每个控制器独立 List/Watch 同一资源类型,减少 API Server 负载。

9.2 Indexer

DaemonSet Controller 使用 podNodeIndex 按 NodeName 索引 Pod,避免全量扫描:

podInformer.Informer().GetIndexer().AddIndexers(cache.Indexers{
    "nodeName": indexByPodNodeName,
})

9.3 Jitter

  • ResyncPeriod 带随机 Jitter,防止所有控制器同时 Resync
  • Controller 启动带 Jitter 间隔,防止同时向 API Server 发起 List

9.4 Batch Update

Endpoint Controller 支持 endpointUpdatesBatchPeriod,将短时间内的多次更新合并为一次。

9.5 ServiceSelectorCache

serviceSelectorCache *endpointutil.ServiceSelectorCache

缓存 Service 的 Selector 解析结果,避免每次 Pod 事件都重新解析。

9.6 Rate Limiting

所有控制器的 ClientSet RESTClient 配置了 RateLimiter,并注册了 Prometheus 指标:

ratelimiter.RegisterMetricAndTrackRateLimiterUsage("controller_name", restClient.GetRateLimiter())

十、总结

kube-controller-manager 是 Kubernetes 控制面的"大脑",它内嵌 30+ 种控制器,通过 Informer + WorkQueue + syncHandler 的统一架构模式,持续驱动集群状态向期望状态收敛。

核心设计精髓:

  1. 声明式协调:控制器不关心"如何到达",只关心"当前 ≠ 期望 → 调整"
  2. 共享 Informer:一次 Watch,多处消费,最小化 API Server 负载
  3. Expectations:防止控制器的热循环,在 Informer 延迟期间避免重复操作
  4. Slow Start:渐进式创建,防止配额不足时的大面积失败
  5. OwnerReference:清晰的层级关系,支持级联删除和 GC
  6. Leader Election:多实例高可用,同一时间只有一个活跃
  7. WorkQueue:去重 + 限速重试 + 优雅关闭

每个控制器虽然具体逻辑不同,但都遵循相同的架构范式,使得 KCM 成为一个可扩展、可维护、高可靠的控制面核心组件。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐