【kubernetes v1.21】(三)kube-controller-manager 超深度分析
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.go → NewControllerManagerCommand():
- 创建
KubeControllerManagerOptions,包含所有子控制器配置选项 - 构建
cobra.Command,注册所有 FlagSet - 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
- Informer:通过 SharedInformerFactory 获取特定资源的 Informer
- EventHandler:注册 Add/Update/Delete 回调,将资源 Key 入队
- WorkQueue:RateLimitingInterface,带去重和退避重试
- Worker:从队列取出 Key,调用 syncHandler
- 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 且不匹配 → 忽略
}
衍生类型:PodControllerRefManager、ReplicaSetControllerRefManager、ControllerRevisionControllerRefManager,分别管理 Pod、ReplicaSet、ControllerRevision 的 OwnerReference。
三、核心业务逻辑深度解析
3.1 Controller Manager 启动流程
3.2 Informer 机制架构
核心流程:
- Reflector 通过
List获取全量数据,再通过Watch持续监听变更 - 变更事件封装为 Delta 对象,写入 DeltaFIFO
- Controller Loop 从 DeltaFIFO Pop,更新 Store(本地缓存)
- 同时触发注册在 Informer 上的 EventHandler
- EventHandler 将资源 Key 放入控制器的 WorkQueue
- Worker 从 WorkQueue 取出 Key,执行 syncHandler
SharedInformerFactory 保证同一资源类型只创建一个 Informer,多个控制器共享同一份缓存。
3.3 Informer-Reflector-Store 数据流
3.4 Controller 总体协调架构
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 流程图
3.5.5 Rolling Update 流程
rolloutRolling() 的核心逻辑(rolling.go):
- 获取 newRS 和 allOldRSs
- 如果 newRS 不存在 → 创建 newRS(带 pod-template-hash)
- 缩容 old RS(按 MaxUnavailable 计算)
- 扩容 new RS(按 MaxSurge 计算)
- 清理旧 RS(超过 RevisionHistoryLimit 则删除)
- 同步 Deployment Status
3.5.6 Recreate 流程
rolloutRecreate() 的核心逻辑(recreate.go):
- 缩容所有 old RS 到 0
- 等待所有 Pod 终止
- 扩容 new RS 到期望副本数
- 与 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 流程图
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 的关键区别:
- 有序性:Pod 按序号(0, 1, 2, …)创建/删除,扩容顺序创建,缩容逆序删除
- 稳定的网络标识:Pod 名为
<statefulset-name>-<ordinal>,DNS 可预测 - 稳定存储:每个 Pod 绑定独立 PVC,Pod 重建后自动重新挂载同一 PVC
- 滚动更新:逆序逐个更新(从最大序号开始),支持 Partition 更新
3.7.3 StatefulSet Controller 流程图
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 核心协调流程
- 获取所有 Node
- 对每个 Node,检查是否应该运行 Daemon Pod(NodeSelector、Taint/Toleration 等)
- 如果应该运行但没有 Pod → 创建
- 如果不应该运行但有 Pod → 删除
- 如果 Pod 存在但模板不匹配 → 根据 UpdateStrategy 决定是否更新
3.8.4 DaemonSet Controller 流程图
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 流程图
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 流程图
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 的映射:
- 监听 Service Add/Update/Delete → 入队
- 监听 Pod Add/Update/Delete → 查找匹配的 Service → 入队
- 监听 Endpoints Delete → 入队(重建)
- 同步时:
- 获取 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 流程图
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 核心职责
- 节点健康监控:持续监控 Node Status 和 Node Lease
- Taint 管理:根据节点状态添加/移除 Taint
NodeReady=False→node-not-readyTaintNodeReady=Unknown→node-unreachableTaintMemoryPressure=True→memory-pressureTaintDiskPressure=True→disk-pressureTaintNetworkUnavailable=True→network-unavailableTaintPIDPressure=True→pid-pressureTaint
- Pod 驱逐:对 NotReady 超时的节点驱逐 Pod
- Zone 级别降级保护:
- 全部 Node NotReady(FullDisruption)→ 降低驱逐速率
- 部分 Node NotReady(PartialDisruption)→ 使用二级驱逐速率
- 正常状态 → 使用正常驱逐速率
- 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 流程图
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 的 ReflectorInformersStartedchannel 在所有控制器初始化后关闭,某些控制器(如 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 收到创建事件 →CreationObserved→add减为 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 流程:
- Adopt:通过 Patch 添加 OwnerReference 到孤儿对象
- Release:通过 StrategicMergePatch 移除 OwnerReference
- CanAdoptFunc:延迟检查控制器是否仍然存在且未被删除,防止竞争条件
RecheckDeletionTimestamp 确保:
- List RS 后、Adopt 前,如果 Deployment 已被删除,则不收养
- 防止已删除的控制器收养新对象,导致 GC 无法清理
4.6 Hash Collision 处理
Deployment 创建 RS 时使用 pod-template-hash 标签:
podTemplateSpecHash := controller.ComputeHash(&newRSTemplate, d.Status.CollisionCount)
如果哈希冲突(不同模板得到相同哈希):
- 尝试 Create → AlreadyExists
- 检查已有 RS 是否属于同一 Deployment 且模板相同
- 如果不同 →
CollisionCount++→ 重新入队 → 下次用新 CollisionCount 计算不同哈希
4.7 Leader Migration
type LeaderMigrator struct {
MigrationReady chan struct{}
FilterFunc FilterFunc
}
FilterFunc 返回三种结果:
ControllerNonMigrated:在主锁中运行ControllerMigrated:在迁移锁中运行ControllerUnowned:不在本实例中运行
迁移流程:
- 主锁获取 Leader → 启动非迁移控制器 + SA Token Controller
- SA Token Controller 启动后关闭 MigrationReady
- 迁移锁获取 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()
控制器关闭时:
- 停止接受新工作
- 等待进行中的工作完成
- 关闭 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 的统一架构模式,持续驱动集群状态向期望状态收敛。
核心设计精髓:
- 声明式协调:控制器不关心"如何到达",只关心"当前 ≠ 期望 → 调整"
- 共享 Informer:一次 Watch,多处消费,最小化 API Server 负载
- Expectations:防止控制器的热循环,在 Informer 延迟期间避免重复操作
- Slow Start:渐进式创建,防止配额不足时的大面积失败
- OwnerReference:清晰的层级关系,支持级联删除和 GC
- Leader Election:多实例高可用,同一时间只有一个活跃
- WorkQueue:去重 + 限速重试 + 优雅关闭
每个控制器虽然具体逻辑不同,但都遵循相同的架构范式,使得 KCM 成为一个可扩展、可维护、高可靠的控制面核心组件。
更多推荐



所有评论(0)