

























kube-scheduler 是 Kubernetes 集群的核心组件之一,负责将未调度的 Pod 绑定到合适的 Node 上。它的启动流程涉及到许多关键步骤:
理解这些环节不仅能帮助我们掌握 Kubernetes 控制平面的工作原理,也能为我们开发_scheduler plugins_和 Operator 提供坚实的地基。
读完本篇,你应该能回答:
Kubernetes Go kube-scheduler Cobra Informer Leader Election Scheduler Framework k8s v1.36.1
学习重点提示 — 建议先通读全文,再重点回顾标注内容
重点掌握(必须)
- Cobra 命令架构:NewSchedulerCommand、PersistentPreRunE、RunE、Args 的协作关系
- Setup + Run 分离:Setup 负责创建 Scheduler 实例,Run 负责启动所有运行时组件
- Leader Election 回调模式:OnStartedLeading/OnStoppedLeading 驱动调度循环的启停
- Informer 缓存同步:Start、WaitForCacheSync、WaitForHandlersSync 的三段式同步
- ScheduleOne 主循环:findNodesThatFitPod + prioritizeNodes + 绑定 的三阶段调度
次重点(了解即可)
- buildHandlerChain 的过滤器链路设计
- CoordinatedLeaderElection 的 LeaseCandidate 机制
- EventBroadcaster 的 StartRecordingToSink 和 Shutdown 生命周期
- SchedulingQueue 的三队列(activeQ/backoffQ/unschedulablePods)设计
文章目录
思考记忆提示 — FAQ 是全篇的"临考前速背"模块,先过一遍问题确认自己是否掌握,再回读薄弱章节
入口文件是 cmd/kube-scheduler/main.go,只有 3 行代码。整个 main 函数没有直接调用任何逻辑,而是先调用 app.NewSchedulerCommand() 创建一个 Cobra 命令对象,然后把这个命令交给 cli.Run(command) 去执行。cli.Run 负责处理信号、初始化日志等通用逻辑,然后最终触发我们在 NewSchedulerCommand 中注册的 RunE 回调。
// cmd/kube-scheduler/main.go
func main() {
command := app.NewSchedulerCommand()
code := cli.Run(command)
os.Exit(code)
}
Cobra 命令包含 Use、Long、PersistentPreRunE、RunE、Args 五个关键配置。cmd/kube-scheduler/app/server.go 中 NewSchedulerCommand 创建命令时:Use 设置命令名称 "kube-scheduler";Long 是长描述文本,解释了调度器的工作原理;PersistentPreRunE 在 RunE 之前执行,用来初始化 Feature Gates;RunE 是实际运行函数,指向 runCommand;Args 是参数校验函数,拒绝任何传入参数。
PersistentPreRunE 先于 RunE 执行,分别负责初始化特性门控和执行调度器主逻辑。PersistentPreRunE 调用 opts.ComponentGlobalsRegistry.Set() 确保 Feature Gates 在任何业务逻辑运行前就已正确设置。RunE 则执行 runCommand,完成日志初始化、配置加载、Scheduler 实例化和启动。这种分离设计保证了特性门控错误可以尽早被捕获,而不是等到业务逻辑运行后才报错。
因为 kube-scheduler 不接受任何位置参数,所有配置都通过 --config 文件或标准 CLI 标志传递。Args 函数的实现遍历所有 args,只要发现任何非空字符串就返回错误。这是一种防御性设计,确保调度器不会被意外传入的额外参数误导。配置通过 YAML 文件(--config)或独立标志(--kubeconfig、--leader-elect 等)传入,这种方式比命令行参数更易于版本管理和审计。
Setup 负责根据命令行选项创建完整的调度器配置和 Scheduler 实例。cmd/kube-scheduler/app/server.go 中 Setup 依次完成:加载默认配置 latest.Default()、校验选项 opts.Validate()、生成运行时配置 opts.Config(ctx)、完成配置 c.Complete()、注册树外插件(out-of-tree registry),最后调用 scheduler.New() 创建 Scheduler 实例。返回的是 CompletedConfig 和 Scheduler 两个对象。
第一个 client 用于调度器的核心操作(带 UserAgent),第二个 eventClient 用于事件发布(不带 UserAgent)。cmd/kube-scheduler/app/options/options.go 中 createClients 的源码:
func createClients(kubeConfig *restclient.Config) (clientset.Interface, clientset.Interface, error) {
client, err := clientset.NewForConfig(
restclient.AddUserAgent(kubeConfig, "scheduler")) // 调度器主操作
if err != nil {
return nil, nil, err
}
eventClient, err := clientset.NewForConfig(kubeConfig) // 事件发布
if err != nil {
return nil, nil, err
}
return client, eventClient, nil
}
第一个 client 被调度器用来 list/watch Pod/Node、更新 Binding 对象等核心操作,restclient.AddUserAgent 在请求头中追加 "scheduler",方便 API Server 端做审计和限流。第二个 eventClient 被 eventBroadcaster 用来将调度事件写入 API Server,不加 UserAgent 是因为事件记录的发起方应该是集群本身。
resyncPeriod 为 0 表示关闭周期性全量同步,只依赖 Watch 增量推送来更新本地缓存。cmd/kube-scheduler/app/options/options.go 第 337 行:
c.InformerFactory = scheduler.NewInformerFactory(client, 0)
在 Kubernetes Informer 架构中,resyncPeriod 控制多久触发一次全量 resync。当值为 0 时,Reflector 不会启动定时 resync 循环,完全依赖 API Server 的 Watch 机制推送变更。对于调度器这种对实时性要求极高的组件,Watch 推送已经是足够快的增量同步机制,关闭 resync 可以避免不必要的 CPU 和内存开销。
InformerFactory 监听 Kubernetes 内置资源,DynInformerFactory 监听动态资源(如 CRD 自定义资源)。cmd/kube-scheduler/app/options/options.go 第 340-341 行:
dynClient := dynamic.NewForConfigOrDie(c.KubeConfig)
c.DynInformerFactory = dynamicinformer.NewFilteredDynamicSharedInformerFactory(
dynClient, 0, corev1.NamespaceAll, nil)
DynInformerFactory 使用 dynamic client,可以监听任意 GroupVersionResource,而不需要预先定义 CRD 类型。这为调度器的扩展性提供了基础——例如 scheduler plugins 可能需要监听自定义资源,CoordinatedLeaderElection 也依赖 DynInformerFactory 来监听 LeaseCandidate 资源。
StartRecordingToSink 启动事件记录的异步处理管道,Shutdown 优雅关闭该管道。cmd/kube-scheduler/app/server.go 第 196-197 行:
cc.EventBroadcaster.StartRecordingToSink(ctx.Done())
defer cc.EventBroadcaster.Shutdown()
StartRecordingToSink 会启动一个 goroutine,从内部 channel 中消费事件并异步写入 API Server 的 events 资源。Shutdown 会等待所有正在处理的事件完成写入,然后关闭内部 channel。defer Shutdown() 确保无论调度器以何种方式退出,事件记录管道都会被优雅关闭,不会遗漏或丢失正在处理的事件。
NewEventBroadcasterAdapterWithContext 是 v1.36.1 中的新 API,DeprecatedNewLegacyRecorder 是旧的兼容接口。cmd/kube-scheduler/app/options/options.go 第 329 行创建 EventBroadcaster,第 346 行创建 LegacyRecorder:
// 在 Options.Config 中
c.EventBroadcaster = events.NewEventBroadcasterAdapterWithContext(ctx, eventClient)
// 在 makeLeaderElectionConfig 中(用于 leader election)
coreRecorder := c.EventBroadcaster.DeprecatedNewLegacyRecorder(schedulerName)
EventBroadcasterAdapterWithContext 是 Kubernetes 1.36 中引入的新设计,提供了更好的 context 支持和事件路由能力。DeprecatedNewLegacyRecorder 保留了旧的 recorder API,主要用于 leader election 资源锁记录事件(因为 leader election 库本身使用旧的 recorder 接口)。两者底层最终都使用 eventClient 将事件写入 API Server。
从外到内依次是:PanicRecovery → HTTPLogging → CacheControl → RequestInfo → Authentication → Authorization。cmd/kube-scheduler/app/server.go 第 350-363 行:
func buildHandlerChain(handler http.Handler, authn authenticator.Request,
authz authorizer.Authorizer) http.Handler {
requestInfoResolver := &apirequest.RequestInfoFactory{}
failedHandler := genericapifilters.Unauthorized(scheme.Codecs)
handler = genericapifilters.WithAuthorization(handler, authz, scheme.Codecs)
handler = genericapifilters.WithAuthentication(handler, authn, failedHandler, nil, nil)
handler = genericapifilters.WithRequestInfo(handler, requestInfoResolver)
handler = genericapifilters.WithCacheControl(handler)
handler = genericfilters.WithHTTPLogging(handler)
handler = genericfilters.WithPanicRecovery(handler, requestInfoResolver)
return handler
}
由于每个 handler = WithXxx(handler, ...) 都是包装,所以最终请求进来时第一个遇到的是 PanicRecovery,最后遇到的是 Authorization(最核心的安全检查)。这种洋葱模型和 Go 的 HTTP Middleware 模式完全一致。
kube-scheduler 暴露 /healthz、/readyz、/metrics、/configz 等核心端点,并通过 installMetricHandler 注册。cmd/kube-scheduler/app/server.go 中 installMetricHandler(第 365-377 行)注册了:/configz 返回调度器配置、/metrics 返回 Prometheus 指标、/metrics/resources 返回 Pod 资源指标(仅 leader 可访问)。newEndpointsHandler 在 staging/src/k8s.io/apiserver/pkg/server/routes.go 中定义,注册了 /healthz 和 /readyz 系列健康检查端点。
internalStopCh 控制 HTTPS Server 立即停止接收新连接,gracefulShutdownSecureServer 等待已有连接优雅关闭。cmd/kube-scheduler/app/server.go 第 261-275 行:
internalStopCh := make(chan struct{})
shutdownTimeout := 5 * time.Second
stoppedCh, listenerStoppedCh, err := cc.SecureServing.Serve(
handler, shutdownTimeout, internalStopCh)
// 延迟执行的优雅关闭函数
gracefulShutdownSecureServer = func() {
close(internalStopCh) // 1. 通知 Server 停止接收新请求
这种两阶段关闭模式是 Kubernetes 控制面组件的标准优雅关闭模式:先通知停止接收(internalStopCh),等监听器退出(listenerStoppedCh),再等现有请求处理完(stoppedCh)。gracefulShutdownSecureServer 在 leader election 丢失时由 OnStoppedLeading 回调调用。
Start 启动所有 Informer 的 Reflector goroutine,WaitForCacheSync 等待初始 List 同步完成。cmd/kube-scheduler/app/server.go 第 278-291 行:
startInformersAndWaitForSync := func(ctx context.Context) {
cc.InformerFactory.Start(ctx.Done()) // 启动所有 Informer
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.Start(ctx.Done())
}
cc.InformerFactory.WaitForCacheSync(ctx.Done()) // 等待初始同步完成
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.WaitForCacheSync(ctx.Done())
}
if err := sched.WaitForHandlersSync(ctx); err != nil { // 等待事件处理器同步
logger.Error(err, "handlers are not fully synchronized")
}
close(handlerSyncReadyCh)
logger.V(3).Info("Handlers synced")
}
Start 启动的是"启动"——它只是让每个 Informer 的 Reflector goroutine 开始运行。Reflector 启动后第一步是执行 List(从 API Server 获取全量数据)来初始化本地缓存。在 List 完成之前,如果调度器就开始处理 Pod,可能会看到不完整的 Node 信息导致调度决策错误。WaitForCacheSync 就是用来"等"这一步完成的。
WaitForCacheSync 等待 Informer 的本地缓存初始化完成,WaitForHandlersSync 等待所有 ResourceEventHandler 收到初始的 Add 事件。pkg/scheduler/scheduler.go 中 WaitForHandlersSync 检查 registeredHandlers 中每个 handler 是否都收到了初始的列表数据。只有 handlerSyncReadyCh 被关闭后,/readyz/sched-handler-sync 健康检查才会通过,确保调度器在完全准备好之前不会被 Kubernetes 认为已就绪。
DelayCacheUntilActive=true 时,非 leader 实例不启动 Informer,以节省资源。cmd/kube-scheduler/app/server.go 第 301 行:
if !cc.ComponentConfig.DelayCacheUntilActive || cc.LeaderElection == nil {
startInformersAndWaitForSync(ctx)
}
当配置了 Leader Election 且 DelayCacheUntilActive=true 时,只有获得 leader 资格(OnStartedLeading 被触发)的实例才会启动 Informer。非 leader 实例会跳过 startInformersAndWaitForSync,直接进入 leader election 等待循环。这在高可用多副本部署中可以显著减少非 leader 实例的资源消耗,因为它们知道自己的 Informer 数据可能不完整,不会参与调度。
说明 Informer 的启动时机取决于 DelayCacheUntilActive 配置:可能提前到 Run 函数中,也可能推迟到获 leader 身份之后。第 302-311 行的完整逻辑:
// 条件一:在 Run 中启动 Informer(DelayCacheUntilActive=false)
if !cc.ComponentConfig.DelayCacheUntilActive || cc.LeaderElection == nil {
startInformersAndWaitForSync(ctx)
}
// 条件二:在 leader election 回调中启动(DelayCacheUntilActive=true)
OnStartedLeading: func(ctx context.Context) {
if cc.ComponentConfig.DelayCacheUntilActive {
logger.Info("Starting informers and waiting for sync...")
startInformersAndWaitForSync(ctx)
logger.Info("Sync completed")
}
sched.Run(ctx)
},
这是一种延迟初始化的优化策略。当集群规模很大时,所有副本同时启动 Informer 会给 API Server 带来巨大的初始 list 压力。DelayCacheUntilActive=true 让非 leader 实例跳过 Informer 启动,只有一个 leader 启动 Informer,从而保护了 API Server。
LeaseCandidate 是一种协同 leader election 机制,通过 lease.candidateleases.coordination.k8s.io 资源协调多个候选者。cmd/kube-scheduler/app/server.go 第 241-254 行:
leaseCandidate, waitForSync, err := leaderelection.NewCandidate(
cc.Client,
metav1.NamespaceSystem,
cc.LeaderElection.Lock.Identity(),
kubeScheduler,
binaryVersion.FinalizeVersion(),
emulationVersion.FinalizeVersion(),
coordinationv1.OldestEmulationVersion,
)
readyzChecks = append(readyzChecks, healthz.NewInformerSyncHealthz(waitForSync))
go leaseCandidate.Run(ctx)
传统的 leader election 使用 configmapsleases 锁,同一时间只有一个候选者持有锁。CoordinatedLeaderElection 引入了 Candidate 机制,每个候选者都持有自己的 lease candidate 记录,并通过版本号(binaryVersion + emulationVersion)确定优先级。这种方式避免了单一锁的竞争瓶颈,提高了高可用场景下的选举效率。waitForSync 会被注册到 /readyz 健康检查中,确保在候选人同步完成前不认为调度器就绪。
OnStoppedLeading 会优雅关闭 HTTPS Server,然后根据退出原因决定是正常退出还是异常退出。cmd/kube-scheduler/app/server.go 第 319-330 行:
OnStoppedLeading: func() {
gracefulShutdownSecureServer() // 关闭 HTTPS Server
select {
case <-ctx.Done():
// 被外部信号终止(SIGTERM/SIGINT),正常退出
logger.Info("Requested to terminate, exiting")
os.Exit(0)
default:
// 丢失 leader 锁,调度器不再负责调度,异常退出
logger.Error(nil, "Leaderelection lost")
klog.FlushAndExit(klog.ExitFlushTimeout, 1)
}
},
gracefulShutdownSecureServer 关闭 HTTPS Server 是必要的,因为非 leader 实例不应该再接受任何来自 API Server 的调度请求。区分 ctx.Done() 和 default 分支是因为:收到 SIGTERM/SIGINT 时上下文被取消(ctx 被 cancel),这是正常的优雅关闭流程;如果上下文未取消而 OnStoppedLeading 被调用,说明是 leader election 竞争失败,此时 os.Exit(1) 让进程以错误码退出。
因为 ScheduleOne 在队列为空时会阻塞在 NextPod() 上,如果放在主 goroutine 中会阻塞后续的队列关闭逻辑。pkg/scheduler/scheduler.go 第 545-573 行:
func (sched *Scheduler) Run(ctx context.Context) {
sched.SchedulingQueue.Run(logger)
if sched.APIDispatcher != nil {
sched.APIDispatcher.Run(logger)
}
// ScheduleOne 在独立 goroutine 中运行,避免阻塞
go wait.UntilWithContext(ctx, sched.ScheduleOne, 0)
<-ctx.Done()
// 只有在上下文取消后,才会执行到这里
if sched.APIDispatcher != nil {
sched.APIDispatcher.Close()
}
sched.SchedulingQueue.Close()
sched.Profiles.Close()
}
wait.UntilWithContext 会反复调用 ScheduleOne 直到上下文取消。如果把 ScheduleOne 放在主 goroutine 中调用,当调度队列为空时 NextPod() 会一直阻塞,<-ctx.Done() 就永远不会被执行,调度队列也就无法被关闭,导致死锁。独立 goroutine 确保了主 goroutine 可以监听 ctx.Done() 信号来触发优雅关闭。
activeQ 存放待调度的就绪 Pod,backoffQ 存放因调度失败进入退避等待的 Pod,unschedulablePods 存放长期无法调度的 Pod。pkg/scheduler/backend/queue/scheduling_queue.go 中 PriorityQueue 的设计:新建或更新后的 Pod 进入 activeQ 被立即处理;调度失败的 Pod 根据失败原因进入 backoffQ 或 unschedulablePods。backoffQ 中的 Pod 在退避超时后会被移回 activeQ;unschedulablePods 中的 Pod 在集群资源状态变化(节点增加、PVC 绑定等)时,通过 QueueingHint 机制被重新评估,如果变为可调度则移回 activeQ。
意味着调度器可以同时运行多个调度配置文件(profiles),每个 profile 拥有独立的 Plugin 组合和调度行为。pkg/scheduler/scheduler.go 第 356-366 行:
profiles, err := profile.NewMap(ctx, options.profiles, registry, recorderFactory,
frameworkruntime.WithComponentConfigVersion(options.componentConfigVersion),
frameworkruntime.WithClientSet(client),
frameworkruntime.WithInformerFactory(informerFactory),
// ... 更多配置选项 ...
)
if len(profiles) == 0 {
return nil, errors.New("at least one profile is required")
}
多 profile 设计支持在同一调度器中运行多个调度策略。例如可以为不同的命名空间配置不同的调度插件组合:一个 profile 使用默认插件,另一个 profile 专门为 GPU 节点使用自定义插件。Pod 通过 spec.schedulerName 字段选择对应的 profile。
因为绑定操作(更新 Pod 的 nodeName、写入 Binding 对象)需要调用 API Server,是网络 I/O 操作,异步执行可以避免阻塞下一个 Pod 的调度。pkg/scheduler/schedule_one.go 第 141-148 行:
// 同步执行调度算法(过滤+评分)
scheduleResult, assumedPodInfo, status := sched.schedulingCycle(
schedulingCycleCtx, state, fwk, podInfo, start, podsToActivate)
if !status.IsSuccess() {
sched.FailureHandler(...)
return
}
// 绑定操作在独立 goroutine 中异步执行
go sched.runBindingCycle(ctx, state, fwk, scheduleResult,
assumedPodInfo, start, podsToActivate)
调度算法(findNodesThatFitPod + prioritizeNodes)执行很快,但绑定操作需要向 API Server 发起 HTTP 请求。如果同步等待绑定完成再处理下一个 Pod,会导致调度器在高并发场景下吞吐量极低。异步绑定让调度器可以在等待 API Server 响应期间继续处理其他 Pod,大幅提升调度效率。
因为 PreFilter 的结果(如 Requestset)是全局的,所有节点共享同一个过滤上下文。pkg/scheduler/schedule_one.go 第 670-675 行:
if !s.IsSuccess() {
if !s.IsRejected() {
return nil, diagnosis, "", s.AsError()
}
// 所有节点在 NodeToStatus 中设置相同状态,便于后续 preemption 处理
diagnosis.NodeToStatus.SetAbsentNodesStatus(s)
return nil, diagnosis, "", nil
}
PreFilter 插件(如 NodeResourcesFit)计算的是 Pod 的资源请求,这些请求对所有节点都是相同的。如果 PreFilter 返回不可调度,说明 Pod 本身的资源请求就已经超过了集群总容量,所有节点都会失败。将相同状态设置到所有节点是为了让 preemption 逻辑能够统一处理这种情况。
调度器注册了 Pod、Node、StorageClass、PV、PVC、CSINode、VolumeAttachment 等资源的 Informer 事件处理器。pkg/scheduler/eventhandlers.go 第 479-706 行 addAllEventHandlers 中:Pod 事件处理器直接操作调度器缓存(addPod 更新 Cache、deletePod 从 Cache 移除);Node 事件处理器更新本地节点缓存;StorageClass、PV、PVC 等资源的变化会触发 MoveAllToActiveOrBackoffQueue,将受影响的未调度 Pod 移回 activeQ 以重新评估。
关注 "Starting Kubernetes Scheduler"、"Golang settings"、handlers synced、"Attempting to schedule pod" 等日志。cmd/kube-scheduler/app/server.go 中关键日志输出:第 178 行的 "Starting Kubernetes Scheduler" 标志 Run 函数开始执行;第 299 行的 "Handlers synced" 标志 Informer 和事件处理器初始化完成;pkg/scheduler/schedule_one.go 中的 "Attempting to schedule pod" 标志调度循环正式启动。查看这些日志可以帮助定位启动失败的具体阶段。
YAML 格式,关键字段包括 LeaderElection、ClientConnection、Profiles、PercentageOfNodesToScore。典型配置结构:
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
clientConnection:
kubeconfig: "/path/to/kubeconfig"
burst: 100
leaderElection:
leaderElect: true
leaseDuration: 15s
renewDeadline: 10s
retryPeriod: 2s
resourceLock: leases
resourceName: kube-scheduler
resourceNamespace: kube-system
profiles:
- schedulerName: default-scheduler
percentageOfNodesToScore: 50
pluginConfig:
- name: NodeResourcesFit
args:
scoringStrategy:
type: LeastAllocated
APIDispatcher 是异步 API 调用调度器,启用后调度器的 list/watch 操作可以在独立 goroutine 中并发执行。pkg/scheduler/scheduler.go 第 352-355 行:
var apiDispatcher *apidispatcher.APIDispatcher
if feature.DefaultFeatureGate.Enabled(features.SchedulerAsyncAPICalls) {
apiDispatcher = apidispatcher.New(client, int(options.parallelism), apicalls.Relevances)
}
当启用 SchedulerAsyncAPICalls 后,调度器的 SchedulingQueue 会通过 APIDispatcher 与 API Server 交互,而不是直接的同步 client 调用。这使得多个调度周期可以并发地执行 API 操作,提高了高吞吐量场景下的调度效率。APIDispatcher 通过工作池模式管理并发请求,避免过载。
通过 Leader Election 机制保证同一时间只有一个 leader 实例运行调度循环,其他实例处于等待状态。整个流程:多个 scheduler 实例同时启动,都尝试获取 kube-scheduler 这个 identity 的 Lease/ConfigMap 锁。只有获得锁的实例(leader)会触发 OnStartedLeading 回调,进而调用 sched.Run() 启动调度循环。非 leader 实例跳过 startInformersAndWaitForSync(如果 DelayCacheUntilActive=true),直接阻塞在 leaderElector.Run(ctx) 上等待 leader 身份变化。当 leader 丢失锁时,OnStoppedLeading 被调用,调度循环停止,另一个等待中的实例会获得 leader 身份并接管调度。
schedulerCache 是调度器的本地节点信息缓存,Assumed 是 Pod 在缓存中被"假设绑定"到某个节点的状态。pkg/scheduler/scheduler.go 第 341-342 行创建缓存:
schedulerCache := internalcache.New(ctx, apiDispatcher,
feature.DefaultFeatureGate.Enabled(features.GenericWorkload))
当调度算法选中一个节点后,调度器会调用 cache.AssumePod() 将 Pod "假设绑定"到该节点——在本地缓存中标记该节点已使用了这部分资源,而不需要立即向 API Server 发起绑定请求。如果后续绑定失败(API Server 返回错误),调度器会调用 cache.ForgetPod() 从缓存中撤销这个假设。这种"先假设再确认"的策略避免了调度算法看到脏数据,同时减少了不必要的 API Server 压力。nodeInfoSnapshot 是调度时对缓存的一个快照,保证整个调度周期内节点信息一致。
全篇必记总纲
kube-scheduler 启动流程的闭环是:main() → NewSchedulerCommand(Cobra) → runCommand → Setup(Scheduler 实例化) → Run(组件启动) → LeaderElection(leader 获取) → Scheduler.Run → ScheduleOne(调度循环),中间穿插 EventBroadcaster(事件)、HTTPS Server(健康检查)、Informer(资源同步)、Cache(本地缓存)四大支撑组件。
思考记忆提示 — 本节是全篇的"入口"——找到入口才能顺藤摸瓜理解后面的每一步
让我们从最简单的文件开始——cmd/kube-scheduler/main.go:
// cmd/kube-scheduler/main.go
package main
import (
"os"
"k8s.io/component-base/cli"
_ "k8s.io/component-base/logs/json/register"
_ "k8s.io/component-base/metrics/prometheus/clientgo"
_ "k8s.io/component-base/metrics/prometheus/version"
"k8s.io/kubernetes/cmd/kube-scheduler/app"
)
func main() {
command := app.NewSchedulerCommand()
code := cli.Run(command)
os.Exit(code)
}
整个 main 函数干净利落,没有一行多余的代码。这里有几点值得注意:
第一,三个 import _ 语句是"仅导入执行 init() 函数"的惯用手法。json/register 将日志输出格式注册为 JSON;两个 metrics 相关的 import 分别注册了 client-go 的 Prometheus 指标和 kube-scheduler 版本的指标端点。这些注册动作在包的 init() 函数中自动完成。
第二,cli.Run 来自 staging/src/k8s.io/component-base/cli/cli.go,它是 Kubernetes 提供的一个 CLI 框架封装。cli.Run 内部会处理一些通用逻辑(如信号初始化、日志初始化、版本打印),然后调用 Cobra 框架执行命令。
我的理解的意思是说
可以把这个 main() 函数想象成一个剧院的领座员——它只做一件事:把票(command)交给剧院入口(cli.Run),然后让观众(调度器进程)自己进去。领座员不关心里面演什么戏、坐多少人,它只负责把人带进门。
真正的演出(调度器的业务逻辑)在 NewSchedulerCommand() 返回的 command 里已经编排好了。
必记闭环逻辑(核心考点)
main 函数是整个进程的起点,它创建 Cobra 命令并交给 CLI 框架执行。所有后续的启动逻辑都注册在命令的 RunE 回调中。
思考记忆提示 — NewSchedulerCommand 是整个命令定义的"宪法"——理解它的结构,后面的源码就顺了
所有启动逻辑的"剧本"都定义在 cmd/kube-scheduler/app/server.go 的 NewSchedulerCommand 函数中:
// cmd/kube-scheduler/app/server.go
func NewSchedulerCommand(registryOptions ...Option) *cobra.Command {
opts := options.NewOptions()
cmd := &cobra.Command{
Use: "kube-scheduler",
Long: `The Kubernetes scheduler is a control plane process which assigns
Pods to Nodes. The scheduler determines which Nodes are valid placements for
each Pod in the scheduling queue according to constraints and available
resources. The scheduler then ranks each valid Node and binds the Pod to a
suitable Node. Multiple different schedulers may be used within a cluster;
kube-scheduler is the reference implementation.`,
PersistentPreRunE: func(*cobra.Command, []string) error {
return opts.ComponentGlobalsRegistry.Set()
},
RunE: func(cmd *cobra.Command, args []string) error {
return runCommand(cmd, opts, registryOptions...)
},
Args: func(cmd *cobra.Command, args []string) error {
for _, arg := range args {
if len(arg) > 0 {
return fmt.Errorf("%q does not take any arguments, got %q",
cmd.CommandPath(), args)
}
}
return nil
},
}
nfs := opts.Flags
verflag.AddFlags(nfs.FlagSet("global"))
globalflag.AddGlobalFlags(nfs.FlagSet("global"), cmd.Name(),
logs.SkipLoggingConfigurationFlags())
fs := cmd.Flags()
for _, f := range nfs.FlagSets {
fs.AddFlagSet(f)
}
cols, _, _ := term.TerminalSize(cmd.OutOrStdout())
cliflag.SetUsageAndHelpFunc(cmd, *nfs, cols)
if err := cmd.MarkFlagFilename("config", "yaml", "yml", "json"); err != nil {
klog.Background().Error(err, "Failed to mark flag filename")
}
return cmd
}
我们来逐一拆解每个关键配置:
Use 设置命令名称为 "kube-scheduler",这个名称出现在 kubectl get pods -n kube-system 中 Pod 的 ownerReference 里。
Long 是长描述文本,解释了调度器的工作原理。kubectl 客户端默认不显示这个文本,但通过 --help 或 API 文档会用到它。
PersistentPreRunE 是一个钩子函数,在 RunE 之前执行。注意这里用的是 PersistentPreRunE 而不是 PersistentPreRun——带后缀 E 的版本返回 error,如果返回非 nil 会阻止 RunE 执行。这里它只做一件事:调用 opts.ComponentGlobalsRegistry.Set() 初始化 Feature Gates。
RunE 是实际运行函数,指向 runCommand。这是启动流程的核心入口。
Args 是参数校验函数。这里它遍历所有位置参数(args),只要发现任何非空字符串就返回错误——这意味着 kube-scheduler 不接受任何命令行位置参数,所有配置都通过 YAML 文件或标准 CLI 标志传递。
最后一段代码负责注册所有 CLI 标志:nfs 包含 options 中定义的所有 FlagSets(secure serving、authentication、authorization、deprecated、leader election、feature gate、metrics、logs),这些标志通过 fs.AddFlagSet 挂载到 Cobra 命令上。
设计精髓
NewSchedulerCommand 的设计体现了"配置即代码"的理念。Cobra 命令本身是一个配置对象,opts.NewOptions() 返回的 Options 包含所有可配置项。这种设计让 kube-scheduler 既可以通过 YAML 配置文件启动,也可以通过命令行标志启动,提供了极大的灵活性。
思考记忆提示 — 本节深入 Cobra 参数解析的细节,是理解启动配置的钥匙
cmd/kube-scheduler/app/options/options.go 中的 Options 结构体定义了调度器所有的可配置项:
// cmd/kube-scheduler/app/options/options.go
type Options struct {
ComponentConfig *kubeschedulerconfig.KubeSchedulerConfiguration
SecureServing *apiserveroptions.SecureServingOptions
Authentication *apiserveroptions.DelegatingAuthenticationOptions
Authorization *apiserveroptions.DelegatingAuthorizationOptions
Metrics *metrics.Options
Logs *logs.Options
Deprecated *DeprecatedOptions
LeaderElection *componentbaseconfig.LeaderElectionConfiguration
ConfigFile string
WriteConfigTo string
Master string
ComponentGlobalsRegistry basecompatibility.ComponentGlobalsRegistry
Flags *cliflag.NamedFlagSets
}
这些字段分别对应不同的配置组:
| 配置组 | 对应字段 | 功能 |
|---|---|---|
| 调度器主配置 | ComponentConfig | Profiles、Extenders、PercentageOfNodesToScore 等 |
| 安全服务 | SecureServing | HTTPS 端口、证书、TLS 配置 |
| 认证授权 | Authentication / Authorization | API Server 请求的认证授权配置 |
| 指标 | Metrics | Prometheus 指标端点配置 |
| 日志 | Logs | 日志级别、格式配置 |
| 废弃选项 | Deprecated | PodMaxInUnschedulablePodsDuration 等 |
| Leader 选举 | LeaderElection | LeaseDuration、RenewDeadline、ResourceLock 等 |
cmd/kube-scheduler/app/options/options.go 中的 initFlags 函数(第 197-218 行)将每个配置组注册为命名的 FlagSet:
func (o *Options) initFlags() {
nfs := cliflag.NamedFlagSets{}
fs := nfs.FlagSet("misc")
fs.StringVar(&o.ConfigFile, "config", o.ConfigFile, "...")
fs.StringVar(&o.WriteConfigTo, "write-config-to", o.WriteConfigTo, "...")
fs.StringVar(&o.Master, "master", o.Master, "...")
o.SecureServing.AddFlags(nfs.FlagSet("secure serving"))
o.Authentication.AddFlags(nfs.FlagSet("authentication"))
o.Authorization.AddFlags(nfs.FlagSet("authorization"))
o.Deprecated.AddFlags(nfs.FlagSet("deprecated"))
options.BindLeaderElectionFlags(o.LeaderElection, nfs.FlagSet("leader election"))
o.ComponentGlobalsRegistry.AddFlags(nfs.FlagSet("feature gate"))
o.Metrics.AddFlags(nfs.FlagSet("metrics"))
logsapi.AddFlags(o.Logs, nfs.FlagSet("logs"))
o.Flags = &nfs
}
小贴士
NamedFlagSets 是 Kubernetes 对标准 Cobra FlagSet 的增强,它给每个 FlagSet 起了一个名字(如 "secure serving"、"authentication")。当用户运行 kube-scheduler --help 时,这些分组名字让帮助文本更有条理。
Validate 方法(第 281-298 行)负责在启动前校验所有配置的有效性:
func (o *Options) Validate() []error {
var errs []error
if err := o.ComponentGlobalsRegistry.SetFallback(); err != nil {
errs = append(errs, err)
} else {
errs = append(errs, o.ComponentGlobalsRegistry.Validate()...)
}
if err := validation.ValidateKubeSchedulerConfiguration(o.ComponentConfig); err != nil {
errs = append(errs, err.Errors()...)
}
errs = append(errs, o.SecureServing.Validate()...)
errs = append(errs, o.Authentication.Validate()...)
errs = append(errs, o.Authorization.Validate()...)
errs = append(errs, o.Metrics.Validate()...)
return errs
}
注意
SetFallback() 在 Validate() 的最前面被调用。SetFallback 的作用是为那些没有被配置文件或 CLI 标志显式设置的选项填充默认值。如果不先调用 SetFallback,某些选项可能是零值,导致 Validate 失败。
必记闭环逻辑(核心考点)
Cobra 的 Flags 系统和 Options 结构体是一一对应的:通过 initFlags 注册标志,通过 Validate 校验合法性,通过 ApplyTo 将标志值应用到运行时配置对象。
思考记忆提示 — runCommand 是第一个"真正做事的函数"——它串联了 Setup 和 Run
cmd/kube-scheduler/app/server.go 第 141-171 行的 runCommand 是整个启动流程的真正入口:
// cmd/kube-scheduler/app/server.go
func runCommand(cmd *cobra.Command, opts *options.Options,
registryOptions ...Option) error {
verflag.PrintAndExitIfRequested()
fg := opts.ComponentGlobalsRegistry.FeatureGateFor(
basecompatibility.DefaultKubeComponent)
if err := logsapi.ValidateAndApply(opts.Logs, fg); err != nil {
fmt.Fprintf(os.Stderr, "%v\n", err)
os.Exit(1)
}
cliflag.PrintFlags(cmd.Flags())
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go func() {
stopCh := server.SetupSignalHandler()
这段代码的执行顺序非常清晰:
第一步:版本标志检查(verflag.PrintAndExitIfRequested())。如果用户传了 --version 或 --v=2 这类版本/日志级别标志,在这里就处理掉了,后面的代码不会执行。
第二步:日志初始化(logsapi.ValidateAndApply)。这里用命令行标志和 Feature Gates 共同决定日志格式和级别。如果用户传了 --logging-format=json,日志会以 JSON 格式输出。
第三步:上下文 + 信号处理。context.WithCancel 创建一个可取消的上下文,然后一个 goroutine 专门监听系统信号(SIGTERM/SIGINT)。staging/src/k8s.io/apiserver/pkg/server/signal.go 中的 SetupSignalHandler 的实现逻辑是:
// staging/src/k8s.io/apiserver/pkg/server/signal.go
func SetupSignalHandler()
这是一种经典的优雅关闭模式:收到第一个信号(比如 SIGTERM)时,cancel 上下文触发所有监听 ctx.Done() 的组件优雅关闭;收到第二个信号时,直接 os.Exit(1) 强制退出,防止第一个关闭流程卡死。
第四步:调用 Setup 创建调度器实例。
第五步:调用 Run 启动调度器。
设计精髓
Setup 和 Run 分离是 Kubernetes 控制面组件的通用模式。Setup 负责创建组件(只执行一次),Run 负责启动组件(可能长期运行)。这种分离使得 Setup 可以在 Run 之外独立测试,也支持"先创建后启动"的更复杂生命周期管理。
思考记忆提示 — Setup 是整个启动流程中代码量最大的函数——它把所有组件组装在一起
cmd/kube-scheduler/app/server.go 第 424-482 行的 Setup 函数:
// cmd/kube-scheduler/app/server.go
func Setup(ctx context.Context, opts *options.Options,
outOfTreeRegistryOptions ...Option) (
*schedulerserverconfig.CompletedConfig,
*scheduler.Scheduler, error) {
// 1. 加载默认配置
if cfg, err := latest.Default(); err != nil {
return nil, nil, err
} else {
opts.ComponentConfig = cfg
}
// 2. 校验选项
if errs := opts.Validate(); len(errs) > 0 {
return nil, nil, utilerrors.NewAggregate(errs)
}
// 3. 生成运行时配置
c, err := opts.Config(ctx)
if err != nil {
return nil, nil, err
}
// 4. 完成配置(合并默认配置和命令行覆盖)
cc := c.Complete()
// 5. 创建 out-of-tree 插件注册表
outOfTreeRegistry := make(runtime.Registry)
for _, option := range outOfTreeRegistryOptions {
if err := option(outOfTreeRegistry); err != nil {
return nil, nil, err
}
}
// 6. 获取事件记录器工厂
recorderFactory := getRecorderFactory(&cc)
// 7. 创建 Scheduler 实例
sched, err := scheduler.New(ctx,
cc.Client,
cc.InformerFactory,
cc.DynInformerFactory,
recorderFactory,
scheduler.WithComponentConfigVersion(cc.ComponentConfig.TypeMeta.APIVersion),
scheduler.WithKubeConfig(cc.KubeConfig),
scheduler.WithProfiles(cc.ComponentConfig.Profiles...),
scheduler.WithPercentageOfNodesToScore(
cc.ComponentConfig.PercentageOfNodesToScore),
scheduler.WithFrameworkOutOfTreeRegistry(outOfTreeRegistry),
scheduler.WithPodMaxBackoffSeconds(cc.ComponentConfig.PodMaxBackoffSeconds),
scheduler.WithPodInitialBackoffSeconds(
cc.ComponentConfig.PodInitialBackoffSeconds),
scheduler.WithPodMaxInUnschedulablePodsDuration(
cc.PodMaxInUnschedulablePodsDuration),
scheduler.WithExtenders(cc.ComponentConfig.Extenders...),
scheduler.WithParallelism(cc.ComponentConfig.Parallelism),
scheduler.WithBuildFrameworkCapturer(
func(profile kubeschedulerconfig.KubeSchedulerProfile) {
completedProfiles = append(completedProfiles, profile)
}),
)
if err != nil {
return nil, nil, err
}
return &cc, sched, nil
}
整个 Setup 流程是线性的,每一步都建立在前一步的结果之上。scheduler.New 是最核心的实例化步骤,我们下一节专门展开。
getRecorderFactory 是一个简单的工厂函数,返回一个闭包用于创建事件记录器:
// cmd/kube-scheduler/app/server.go
func getRecorderFactory(cc *schedulerserverconfig.CompletedConfig,
) profile.RecorderFactory {
return func(name string) events.EventRecorderLogger {
return cc.EventBroadcaster.NewRecorder(name)
}
}
每次调用这个工厂函数,都会从 EventBroadcaster 创建一个新的事件记录器,每个记录器带一个名称参数(通常是调度器名称或插件名称),方便在 Kubernetes 事件中标识事件来源。
小贴士
out-of-tree registry 是一种扩展机制,允许在调度器编译之后加载额外的调度插件(如自定义的过滤或评分插件)。在 Setup 中,所有 out-of-tree 插件通过 Option 函数注册到 registry 中,然后在 scheduler.New 时通过 WithFrameworkOutOfTreeRegistry 传入,与树内插件(in-tree)合并。
pkg/scheduler/scheduler.go 中的 Scheduler 结构体是整个调度器的核心数据结构:
// pkg/scheduler/scheduler.go
type Scheduler struct {
Cache internalcache.Cache // 本地节点/Pod 缓存
Extenders []fwk.Extender // 扩展调度器列表
NextPod func(logger klog.Logger) (*framework.QueuedPodInfo, error)
FailureHandler FailureHandlerFn // 调度失败时的回调
SchedulePod func(ctx, fwk, state, podInfo) (ScheduleResult, error)
StopEverything
这些字段的初始化贯穿整个 scheduler.New 函数(第 276-468 行),其中最关键的初始化包括:Cache(内部缓存)、SchedulingQueue(调度队列)、Profiles(调度框架)和注册事件处理器。
必记闭环逻辑(核心考点)
Setup 是"组装工厂":它接收命令行选项,依次经过配置加载、校验、实例化,最终返回一个可运行的 Scheduler 实例。scheduler.New 是"组装车间":它接收各种 Option 函数,组装出包含 Cache、Queue、Profiles、EventHandlers 的完整调度器。
思考记忆提示 — Options.Config 是将配置转换为运行时对象的关键一步
cmd/kube-scheduler/app/options/options.go 中的 Options.Config 方法(第 300-345 行)是连接配置和运行时对象的桥梁:
// cmd/kube-scheduler/app/options/options.go
func (o *Options) Config(ctx context.Context) (*schedulerappconfig.Config, error) {
logger := klog.FromContext(ctx)
// 1. 自签名证书(如果启用了安全服务)
if o.SecureServing != nil {
if err := o.SecureServing.MaybeDefaultWithSelfSignedCerts(
"localhost", nil, []net.IP{netutils.ParseIPSloppy("127.0.0.1")}); err != nil {
return nil, fmt.Errorf("error creating self-signed certificates: %v", err)
}
}
c := &schedulerappconfig.Config{}
// 2. 应用配置到 Config 对象
if err := o.ApplyTo(logger, c); err != nil {
return nil, err
}
// 3. 创建两个 clientset
client, eventClient, err := createClients(c.KubeConfig)
if err != nil {
return nil, err
}
// 4. 创建 EventBroadcaster
c.EventBroadcaster = events.NewEventBroadcasterAdapterWithContext(ctx, eventClient)
// 5. Leader Election 配置(如果启用)
var leaderElectionConfig *leaderelection.LeaderElectionConfig
if c.ComponentConfig.LeaderElection.LeaderElect {
schedulerName := corev1.DefaultSchedulerName
if len(c.ComponentConfig.Profiles) != 0 {
schedulerName = c.ComponentConfig.Profiles[0].SchedulerName
}
coreRecorder := c.EventBroadcaster.DeprecatedNewLegacyRecorder(schedulerName)
leaderElectionConfig, err = makeLeaderElectionConfig(
c.ComponentConfig.LeaderElection, c.KubeConfig, coreRecorder)
if err != nil {
return nil, err
}
}
// 6. 创建 InformerFactory 和 DynInformerFactory
c.Client = client
c.InformerFactory = scheduler.NewInformerFactory(client, 0)
dynClient := dynamic.NewForConfigOrDie(c.KubeConfig)
c.DynInformerFactory = dynamicinformer.NewFilteredDynamicSharedInformerFactory(
dynClient, 0, corev1.NamespaceAll, nil)
c.LeaderElection = leaderElectionConfig
c.ComponentGlobalsRegistry = o.ComponentGlobalsRegistry
return c, nil
}
createClients 函数创建了两个 Kubernetes 客户端:
// cmd/kube-scheduler/app/options/options.go
func createClients(kubeConfig *restclient.Config) (clientset.Interface, clientset.Interface, error) {
client, err := clientset.NewForConfig(
restclient.AddUserAgent(kubeConfig, "scheduler"))
if err != nil {
return nil, nil, err
}
eventClient, err := clientset.NewForConfig(kubeConfig)
if err != nil {
return nil, nil, err
}
return client, eventClient, nil
}
两个 clientset 使用同一个 kubeConfig,但第一个加了 UserAgent 标识用于调度器核心操作,第二个不加用于事件记录。
makeLeaderElectionConfig 构建 leader election 配置:
// cmd/kube-scheduler/app/options/options.go
func makeLeaderElectionConfig(
config componentbaseconfig.LeaderElectionConfiguration,
kubeConfig *restclient.Config,
recorder record.EventRecorder) (*leaderelection.LeaderElectionConfig, error) {
hostname, err := os.Hostname()
if err != nil {
return nil, fmt.Errorf("unable to get hostname: %v", err)
}
id := hostname + "_" + string(uuid.NewUUID()) // 唯一标识符
rl, err := resourcelock.NewFromKubeconfig(
config.ResourceLock,
config.ResourceNamespace,
config.ResourceName,
resourcelock.ResourceLockConfig{
Identity: id,
EventRecorder: recorder,
},
kubeConfig,
config.RenewDeadline.Duration)
if err != nil {
return nil, fmt.Errorf("couldn't create resource lock: %v", err)
}
return &leaderelection.LeaderElectionConfig{
Lock: rl,
LeaseDuration: config.LeaseDuration.Duration,
RenewDeadline: config.RenewDeadline.Duration,
RetryPeriod: config.RetryPeriod.Duration,
WatchDog: leaderelection.NewLeaderHealthzAdaptor(time.Second * 20),
Name: "kube-scheduler",
ReleaseOnCancel: true,
}, nil
}
这里有几个关键点:每个调度器实例的 identity 由 hostname + UUID 组成,保证了在同一台机器上运行多个副本时 identity 也不会冲突。ResourceLock 指定了使用 leases 资源(kubernetes 1.14+ 推荐,比 configmaps 更高效)。WatchDog 是一个健康检查适配器,会在 leader election 期间持续检查 leader 是否还活着。
必记闭环逻辑(核心考点)
Options.Config 将命令行选项转换为运行时对象:KubeConfig、Client、EventBroadcaster、InformerFactory、LeaderElectionConfig。这些对象在 Setup 中被组装成 CompletedConfig,然后传给 scheduler.New 进一步实例化。
思考记忆提示 — EventBroadcaster 在 Run 函数中被启动,它负责将调度事件写入 API Server
EventBroadcaster 在 cmd/kube-scheduler/app/server.go 的 Run 函数中被启动:
// cmd/kube-scheduler/app/server.go Run 函数
cc.EventBroadcaster.StartRecordingToSink(ctx.Done())
defer cc.EventBroadcaster.Shutdown()
EventBroadcaster 是 Kubernetes 事件机制的核心组件。它接收来自调度器和各个插件的事件(Pod 调度成功、调度失败、节点变化等),然后异步写入 API Server 的 events 资源。这些事件可以通过 kubectl describe pod <name> 查看,是调试调度问题的重要工具。
StartRecordingToSink 启动事件记录的异步处理管道:从内部 channel 消费事件,批量写入 API Server。defer cc.EventBroadcaster.Shutdown() 确保在进程退出前所有事件都被写入。
recorderFactory 在 Setup 中创建:
func getRecorderFactory(cc *schedulerserverconfig.CompletedConfig,
) profile.RecorderFactory {
return func(name string) events.EventRecorderLogger {
return cc.EventBroadcaster.NewRecorder(name)
}
}
每次调用这个工厂函数创建一个新的事件记录器,传入的名称会出现在事件消息中。例如,"default-scheduler" 记录的事件看起来像:
Events:
Type Reason Age From Message
──── ───── ─── ─── ───────
Normal Scheduled 12s default-scheduler Successfully assigned
default/nginx-deployment-7fb96c846b-abcde to node-1
我的理解的意思是说
EventBroadcaster 就像一个电视台的导播台:调度器和各个插件是"摄像机",它们把"节目"(事件)发给导播台,导播台负责把内容"直播"到 API Server(相当于电视台的播出频道)。StartRecordingToSink 是"开始直播",Shutdown 是"关闭直播并确保最后的画面播出完成"。
必记闭环逻辑(核心考点)
EventBroadcaster 的启动是 Run 函数的第一步,它的关闭是 defer 的最后一步。这种"启动最早、关闭最晚"的模式确保了整个调度器生命周期中的所有事件都能被记录。
思考记忆提示 — kube-scheduler 不仅是一个调度器,还是一个 HTTP Server——它暴露了健康检查和指标端点
Run 函数中 HTTPS Server 的启动逻辑(第 257-276 行):
gracefulShutdownSecureServer := func() {}
if cc.SecureServing != nil {
handler := buildHandlerChain(
newEndpointsHandler(
&cc.ComponentConfig, cc.InformerFactory, isLeader,
checks, readyzChecks, cc.Flagz),
cc.Authentication.Authenticator,
cc.Authorization.Authorizer)
internalStopCh := make(chan struct{})
shutdownTimeout := 5 * time.Second
stoppedCh, listenerStoppedCh, err := cc.SecureServing.Serve(
handler, shutdownTimeout, internalStopCh)
if err != nil {
close(internalStopCh)
return fmt.Errorf("failed to start secure server: %v", err)
}
gracefulShutdownSecureServer = func() {
close(internalStopCh)
buildHandlerChain 是 kube-scheduler 中一个经典的设计模式——HTTP 过滤器链(也称为"洋葱模型"):
// cmd/kube-scheduler/app/server.go
func buildHandlerChain(handler http.Handler,
authn authenticator.Request,
authz authorizer.Authorizer) http.Handler {
requestInfoResolver := &apirequest.RequestInfoFactory{}
failedHandler := genericapifilters.Unauthorized(scheme.Codecs)
handler = genericapifilters.WithAuthorization(handler, authz, scheme.Codecs)
handler = genericapifilters.WithAuthentication(handler, authn, failedHandler, nil, nil)
handler = genericapifilters.WithRequestInfo(handler, requestInfoResolver)
handler = genericapifilters.WithCacheControl(handler)
handler = genericfilters.WithHTTPLogging(handler)
handler = genericfilters.WithPanicRecovery(handler, requestInfoResolver)
return handler
}
请求进入时首先经过 PanicRecovery(防止 panic 导致进程崩溃),然后是 HTTPLogging(日志记录)、CacheControl(缓存控制头)、RequestInfo(提取请求元数据)、Authentication(验证身份)、Authorization(检查权限)。每层都是一个装饰器,返回一个新的 Handler 包裹原来的 Handler。
newEndpointsHandler 在 staging/src/k8s.io/apiserver/pkg/server/routes.go 中定义,它注册了 /healthz、/readyz、/livez、/metrics 等核心端点,并传入 isLeader 函数让某些端点在非 leader 状态下返回错误。
健康检查(healthzChecks 和 readyzChecks)的注册在 Run 函数开头(第 199-228 行):
var checks, readyzChecks []healthz.HealthChecker
if cc.ComponentConfig.LeaderElection.LeaderElect {
checks = append(checks, cc.LeaderElection.WatchDog)
readyzChecks = append(readyzChecks, cc.LeaderElection.WatchDog)
}
readyzChecks = append(readyzChecks, healthz.NewShutdownHealthz(ctx.Done()))
handlerSyncCheck := healthz.NamedCheck("sched-handler-sync",
func(_ *http.Request) error {
select {
case
这三个健康检查分别是:WatchDog(leader election 的看门狗)、ShutdownHealthz(检查是否正在关闭)和 sched-handler-sync(检查 Informer handler 是否同步完成)。只有当所有健康检查都通过时,kube-scheduler 才会被认为已就绪,可以接受调度请求。
设计精髓
HTTPS Server 的设计体现了关注点分离:健康检查端点由 Run 函数组装后传入 newEndpointsHandler,而不是硬编码在某个函数里。isLeader 函数使得 /metrics/resources 等资源密集型端点只在 leader 实例上响应,非 leader 实例直接跳过处理,节省资源。
思考记忆提示 — Informer 是 kube-scheduler 的"眼睛"——它让调度器能实时感知集群状态变化
Run 函数中的 Informer 启动逻辑(第 278-300 行):
startInformersAndWaitForSync := func(ctx context.Context) {
// 阶段一:启动所有 Informer 的 Reflector goroutine
cc.InformerFactory.Start(ctx.Done())
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.Start(ctx.Done())
}
// 阶段二:等待初始 List 同步完成
cc.InformerFactory.WaitForCacheSync(ctx.Done())
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.WaitForCacheSync(ctx.Done())
}
// 阶段三:等待所有事件处理器收到初始 Add 事件
if err := sched.WaitForHandlersSync(ctx); err != nil {
logger.Error(err, "handlers are not fully synchronized")
}
close(handlerSyncReadyCh)
logger.V(3).Info("Handlers synced")
}
if !cc.ComponentConfig.DelayCacheUntilActive || cc.LeaderElection == nil {
startInformersAndWaitForSync(ctx)
}
阶段一:Start。每个 Informer 的 Start 方法启动其内部的 Reflector goroutine。Reflector 是 Kubernetes Informer 架构中的"协调员",它负责从 API Server 拉取(通过 List)和接收(通过 Watch)资源变更。
阶段二:WaitForCacheSync。Reflector 启动后第一步是执行 List 操作(从 API Server 获取全量数据)来初始化本地缓存(Indexer)。WaitForCacheSync 阻塞直到所有 Informer 的 List 操作完成。在这一步完成前,调度器如果开始调度可能会看到不完整的节点列表。
阶段三:WaitForHandlersSync。即使 Indexer 初始化完成,每个事件处理器的初始 Add 事件也可能还没被送达。WaitForHandlersSync 检查 registeredHandlers 中的每个 handler 是否都收到了初始列表数据。只有这一步完成后,/readyz/sched-handler-sync 检查才会通过。
scheduler.New 中通过 addAllEventHandlers 注册了所有事件处理器:
// pkg/scheduler/eventhandlers.go
if err = addAllEventHandlers(sched, informerFactory, dynInformerFactory,
resourceClaimCache, resourceSliceTracker, draManager,
unionedGVKs(queueingHintsPerProfile)); err != nil {
return nil, fmt.Errorf("adding event handlers: %w", err)
}
addAllEventHandlers 注册的事件处理器包括:
| 资源 | 事件类型 | 处理函数 | 作用 |
|---|---|---|---|
| Pod | Add/Update/Delete | addPod/updatePod/deletePod | 更新调度器本地缓存 |
| Node | Add/Update/Delete | addNodeToCache/updateNodeInCache/deleteNodeFromCache | 更新节点缓存 |
| StorageClass | Add/Update | MoveAllToActiveOrBackoffQueue | 触发调度重评估 |
| PV/PVC | Add/Update/Delete | MoveAllToActiveOrBackoffQueue | 触发调度重评估 |
| CSINode | Add/Update/Delete | MoveAllToActiveOrBackoffQueue | 触发调度重评估 |
| VolumeAttachment | Add/Update/Delete | MoveAllToActiveOrBackoffQueue | 触发调度重评估 |
DelayCacheUntilActive 是一个性能优化选项。当设置为 true 时,只有获得 leader 身份的实例才会启动 Informer。这意味着如果部署了 3 个 scheduler 副本,只有 1 个会启动 Informer(并给 API Server 带来初始 list 的压力),其他 2 个暂时不启动 Informer,节省 CPU 和内存。
注意
WaitForCacheSync 使用 ctx.Done() 作为超时机制。如果在这个等待期间收到关闭信号,WaitForCacheSync 会立即返回 false,导致调度器启动失败。这是正确的行为——在关闭期间不应该继续调度。
必记闭环逻辑(核心考点)
Informer 的三阶段同步是调度器就绪的前置条件:先让 Reflector 启动(能接收变更),再等 List 完成(看到完整初始数据),最后等 Handler 收到初始事件(能处理变更)。只有完成这三步,调度器才真正准备好工作了。
思考记忆提示 — Leader Election 是高可用调度的基石——理解它才能理解为什么多副本调度器不会重复调度
Leader Election 的配置和启动在 Run 函数中(第 305-341 行):
if cc.LeaderElection != nil {
if utilfeature.DefaultFeatureGate.Enabled(
kubefeatures.CoordinatedLeaderElection) {
cc.LeaderElection.Coordinated = true
}
cc.LeaderElection.Callbacks = leaderelection.LeaderCallbacks{
OnStartedLeading: func(ctx context.Context) {
close(waitingForLeader)
if cc.ComponentConfig.DelayCacheUntilActive {
logger.Info("Starting informers and waiting for sync...")
startInformersAndWaitForSync(ctx)
logger.Info("Sync completed")
}
sched.Run(ctx)
},
OnStoppedLeading: func() {
gracefulShutdownSecureServer()
select {
case <-ctx.Done():
logger.Info("Requested to terminate, exiting")
os.Exit(0)
default:
logger.Error(nil, "Leaderelection lost")
klog.FlushAndExit(klog.ExitFlushTimeout, 1)
}
},
}
leaderElector, err := leaderelection.NewLeaderElector(*cc.LeaderElection)
if err != nil {
return fmt.Errorf("couldn't create leader elector: %v", err)
}
leaderElector.Run(ctx)
return fmt.Errorf("lost lease")
}
Leader Election 的工作原理基于一个简单的观察:如果有多个调度器副本都尝试在同一时间运行调度循环,那么同一个 Pod 可能会被多个副本同时评估和绑定,导致重复调度和竞争条件。Leader Election 确保同一时间只有一个副本是 leader——只有 leader 才真正执行调度。
OnStartedLeading 是"当选为 leader"的回调。当前的实现逻辑:先关闭 waitingForLeader channel(让 isLeader() 返回 true),然后根据配置决定是否启动 Informer,最后调用 sched.Run(ctx) 启动调度循环。
OnStoppedLeading 是"失去 leader 身份"的回调。先调用 gracefulShutdownSecureServer() 关闭 HTTPS Server(因为非 leader 不应该接受任何请求),然后根据上下文是否已取消决定是正常退出还是异常退出。
CoordinatedLeaderElection 是 Kubernetes 1.36 中引入的增强机制。当启用时,每个候选者通过 lease.candidateleases.coordination.k8s.io 资源声明自己的身份,API Server 自动根据版本号(binaryVersion + emulationVersion)确定优先级,优先级最高的候选者获得 leader 身份。这种方式比传统单锁模式减少了竞争,提高了选举效率。
isLeader 函数和 waitingForLeader channel 的协作:
waitingForLeader := make(chan struct{})
isLeader := func() bool {
select {
case _, ok :=
初始时 waitingForLeader 处于打开状态,isLeader() 返回 false。OnStartedLeading 关闭这个 channel 后,isLeader() 返回 true。这个 isLeader 函数被传入 newEndpointsHandler,用于控制 /metrics/resources 等仅 leader 访问的端点。
设计精髓
Leader Election 的回调模式是一种"控制反转"设计——不是 leader election 代码主动调用调度器,而是调度器注册回调函数,告诉 leader election 代码"当选 leader 时请执行这些操作"。这种设计让 leader election 库与业务逻辑完全解耦。
思考记忆提示 — Scheduler.Run 是调度循环的起点——从这开始,调度器真正开始工作了
pkg/scheduler/scheduler.go 第 545-573 行的 Scheduler.Run:
// pkg/scheduler/scheduler.go
func (sched *Scheduler) Run(ctx context.Context) {
logger := klog.FromContext(ctx)
// 1. 启动调度队列的内部处理循环
sched.SchedulingQueue.Run(logger)
// 2. 启动异步 API 调度器(如果启用)
if sched.APIDispatcher != nil {
sched.APIDispatcher.Run(logger)
}
// 3. 在独立 goroutine 中启动 ScheduleOne 循环
go wait.UntilWithContext(ctx, sched.ScheduleOne, 0)
// 4. 阻塞等待上下文取消(收到关闭信号后从这里恢复)
<-ctx.Done()
// 5. 优雅关闭
if sched.APIDispatcher != nil {
sched.APIDispatcher.Close()
}
sched.SchedulingQueue.Close()
// 6. 关闭所有 framework 插件
err := sched.Profiles.Close()
if err != nil {
logger.Error(err, "Failed to close plugins")
}
}
这里有三个关键点:
为什么是独立 goroutine? 如果把 wait.UntilWithContext(ctx, sched.ScheduleOne, 0) 直接放在主 goroutine 中调用,当调度队列为空时,NextPod() 会一直阻塞在 channel 读取上,<-ctx.Done() 永远不会被执行。这意味着即使用户发送了 SIGTERM 信号,调度器也无法响应,优雅关闭变成了一句空话。
wait.UntilWithContext 是 Kubernetes 封装的工具函数,它反复调用传入的函数(这里是 ScheduleOne),每次调用间隔由第三个参数决定(0 表示无等待,立即开始下一次)。当上下文取消时,这个函数返回。
关闭顺序:先关闭 APIDispatcher(停止异步 API 调用),再关闭 SchedulingQueue(停止接收新 Pod),最后关闭所有 Framework 插件。这个顺序确保了在处理中的任务能安全地完成或取消。
必记闭环逻辑(核心考点)
Scheduler.Run 是调度循环的"发动机点火":SchedulingQueue 提供待调度的 Pod,APIDispatcher 处理异步 API 调用,ScheduleOne 反复从队列中取 Pod 并执行完整的调度算法,直到收到关闭信号。
思考记忆提示 — ScheduleOne 是调度器的心脏——每个 Pod 的调度决策都在这里产生
pkg/scheduler/schedule_one.go 的 ScheduleOne 是调度器的主循环函数:
// pkg/scheduler/schedule_one.go
func (sched *Scheduler) ScheduleOne(ctx context.Context) {
logger := klog.FromContext(ctx)
podInfo, err := sched.NextPod(logger)
if err != nil {
utilruntime.HandleErrorWithLogger(logger, err,
"Error while retrieving next pod from scheduling queue")
return
}
if podInfo == nil || podInfo.Pod == nil {
return // 队列已关闭
}
logger = klog.LoggerWithValues(logger, "pod", klog.KObj(podInfo.Pod))
ctx = klog.NewContext(ctx, logger)
logger.V(4).Info("About to try and schedule pod", "pod", klog.KObj(podInfo.Pod))
fwk, err := sched.frameworkForPod(podInfo.Pod)
if err != nil {
sched.SchedulingQueue.Done(podInfo.Pod.UID)
return
}
if sched.skipPodSchedule(ctx, fwk, podInfo) {
sched.SchedulingQueue.Done(podInfo.Pod.UID)
return
}
logger.V(3).Info("Attempting to schedule pod", "pod", klog.KObj(podInfo.Pod))
start := time.Now()
state := framework.NewCycleState()
scheduleResult, assumedPodInfo, status := sched.schedulingCycle(
schedulingCycleCtx, state, fwk, podInfo, start, podsToActivate)
if !status.IsSuccess() {
sched.FailureHandler(schedulingCycleCtx, fwk, assumedPodInfo,
status, scheduleResult.nominatingInfo, start)
return
}
go sched.runBindingCycle(ctx, state, fwk, scheduleResult,
assumedPodInfo, start, podsToActivate)
}
整个流程分为三个层次:
1. 取出 Pod(NextPod)。从 SchedulingQueue 的 activeQ 中阻塞获取下一个待调度的 Pod。如果队列为空,这里会一直等待直到有新 Pod 入队或收到关闭信号。
2. 调度周期(schedulingCycle)。这是同步执行的调度算法,包括 findNodesThatFitPod(过滤)和 prioritizeNodes(评分)。pkg/scheduler/schedule_one.go 第 567-624 行的 schedulePod 函数展示了核心逻辑:
// pkg/scheduler/schedule_one.go
func (sched *Scheduler) schedulePod(ctx context.Context,
fwk framework.Framework, state fwk.CycleState,
podInfo *framework.QueuedPodInfo) (result ScheduleResult, err error) {
if sched.nodeInfoSnapshot.NumNodesInPlacement() == 0 {
return result, ErrNoNodesAvailable
}
// 阶段一:找到所有适合的节点(过滤)
feasibleNodes, diagnosis, nodeHint, err := sched.findNodesThatFitPod(
ctx, fwk, state, podInfo)
if len(feasibleNodes) == 0 {
return result, &framework.FitError{
Pod: pod,
NumAllNodes: sched.nodeInfoSnapshot.NumNodesInPlacement(),
Diagnosis: diagnosis,
}
}
// 如果只有一个节点,直接使用(优化路径)
if len(feasibleNodes) == 1 {
return ScheduleResult{
SuggestedHost: feasibleNodes[0].Node().Name,
EvaluatedNodes: 1 + diagnosis.NodeToStatus.Len(),
FeasibleNodes: 1,
}, nil
}
// 阶段二:对节点进行评分(排序)
priorityList, err := prioritizeNodes(ctx, sched.Extenders, fwk,
state, pod, feasibleNodes)
sortedPrioritizedNodes := newSortedNodeScores(priorityList)
node := sortedPrioritizedNodes.Pop()
return ScheduleResult{
SuggestedHost: node,
EvaluatedNodes: len(feasibleNodes) + diagnosis.NodeToStatus.Len(),
FeasibleNodes: len(feasibleNodes),
}, err
}
过滤阶段调用 findNodesThatFitPod,它依次执行:
评分阶段调用 prioritizeNodes,它依次执行:
3. 绑定周期(runBindingCycle)。调度算法完成后,绑定操作在独立 goroutine 中异步执行。这包括:AssumPod(假设绑定到缓存)、执行 Bind 插件(Permit、PreBind)、向 API Server 写入 Binding 对象。
我的理解的意思是说
可以把 ScheduleOne 想象成一场考试:监考老师(调度器)从待考名单(SchedulingQueue)中叫一个学生(Pod)进场,然后对他进行两步考核:
必记闭环逻辑(核心考点)
ScheduleOne 的三阶段流水线(过滤 → 评分 → 绑定)是一个 Pod 从"待调度"到"已绑定"的核心路径。过滤和评分是同步的(保证决策一致性),绑定是异步的(提高并发吞吐量)。
思考记忆提示 — 本节将前面所有知识串成一张图,帮助你建立整体视野
让我们把整个启动流程串联成一张完整的图:
┌──────────────────────────────────────────────────────────────────────────────────┐
│ kube-scheduler 启动流程全景图 │
│ │
│ 阶段一:入口与命令创建(cmd/kube-scheduler/main.go + server.go) │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ main() → NewSchedulerCommand() → cli.Run() │ │
│ │ - 创建 cobra.Command(Use/Long/PreRunE/RunE/Args) │ │
│ │ - 注册所有 FlagSets(secure serving/authn/authz/leader election 等) │ │
│ └─────────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ 阶段二:启动准备(runCommand) │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ runCommand(): │ │
│ │ - verflag 检查(--version/-v) │ │
│ │ - logsapi.ValidateAndApply(初始化日志系统) │ │
│ │ - SetupSignalHandler → goroutine 监听 SIGTERM/SIGINT → cancel ctx │ │
│ └─────────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ 阶段三:配置与实例化(Setup) │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ Setup(): │ │
│ │ - latest.Default() → 加载默认 KubeSchedulerConfiguration │ │
│ │ - opts.Validate() → 校验所有配置项 │ │
│ │ - opts.Config(ctx) → 创建 KubeConfig/Client/InformerFactory │ │
│ │ - c.Complete() → 合并默认配置和命令行覆盖 │ │
│ │ - scheduler.New() → 创建 Scheduler 实例 │ │
│ │ ├─ internalcache.New() → 调度器本地缓存 │ │
│ │ ├─ profile.NewMap() → 创建 Framework(Plugins) │ │
│ │ ├─ NewSchedulingQueue() → 三队列(activeQ/backoffQ/unschedulable)│ │
│ │ └─ addAllEventHandlers() → 注册 Pod/Node/StorageClass 等 handler │ │
│ └─────────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ 阶段四:运行(Run) │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ Run(): │ │
│ │ ① EventBroadcaster.StartRecordingToSink() → 启动事件记录管道 │ │
│ │ ② 配置 healthz/readyz 健康检查(WatchDog/ShutdownHealthz/SyncCheck) │ │
│ │ ③ isLeader() / waitingForLeader channel 设置 │ │
│ │ ④ CoordinatedLeaderElection(如果启用) │ │
│ │ └─ NewCandidate → leaseCandidate.Run() → /readyz 增加检查 │ │
│ │ ⑤ HTTPS Server 启动(buildHandlerChain → Serve) │ │
│ │ └─ 过滤器链:PanicRecovery → HTTPLogging → CacheControl │ │
│ │ → RequestInfo → Authentication → Authorization │ │
│ │ ⑥ startInformersAndWaitForSync()(可能延迟到 leader election 后) │ │
│ │ ├─ InformerFactory.Start() → 启动 Reflector goroutine │ │
│ │ ├─ WaitForCacheSync() → 等待初始 List 完成 │ │
│ │ └─ WaitForHandlersSync() → 等待事件处理器收到初始事件 │ │
│ │ ⑦ Leader Election 配置(如果启用) │ │
│ │ └─ NewLeaderElector → leaderElector.Run(ctx) │ │
│ │ OnStartedLeading: │ │
│ │ - 关闭 waitingForLeader(isLeader→true) │ │
│ │ - startInformersAndWaitForSync(可选) │ │
│ │ - sched.Run(ctx) → 启动调度循环 │ │
│ │ OnStoppedLeading: │ │
│ │ - gracefulShutdownSecureServer() │ │
│ │ - ctx.Done() ? os.Exit(0) : os.Exit(1) │ │
│ └─────────────────────────────────────────────────────────────────────────┘ │
│ ↓ │
│ 阶段五:调度循环(Scheduler.Run → ScheduleOne) │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ Scheduler.Run(): │ │
│ │ - SchedulingQueue.Run() → 启动队列处理 goroutine │ │
│ │ - APIDispatcher.Run() → 启动异步 API 调度器(可选) │ │
│ │ - go wait.UntilWithContext(ScheduleOne) → 调度主循环 │ │
│ │ - <-ctx.Done() → 优雅关闭 │ │
│ │ │ │
│ │ ScheduleOne(): │ │
│ │ - NextPod() → 从 activeQ 阻塞获取下一个 Pod │ │
│ │ - schedulingCycle (同步): │ │
│ │ ├─ findNodesThatFitPod → 过滤(PreFilter + Filter + Extender) │ │
│ │ └─ prioritizeNodes → 评分(PreScore + Score + Extender) │ │
│ │ - go runBindingCycle (异步) → 绑定(PreBind + Bind + 写入 API Server)│ │
│ └─────────────────────────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────────────────────┘
还有一些关键细节值得特别关注:
第一,ConfigZ 注册。Run 函数在启动早期就向 /configz 端点注册了调度器的配置对象(通过 configz.New 和 cz.Set)。这意味着运维人员可以在运行时通过 GET 请求查看调度器的完整配置。
第二,Feature Gate 指标上报。runCommand 在调用 Setup 之后会调用 fg.AddMetrics() 和 opts.ComponentGlobalsRegistry.AddMetrics(),将特性门控的状态作为 Prometheus 指标暴露出来。这对于监控哪些实验性特性被启用非常有用。
第三,WriteConfigTo 选项。如果用户传了 --write-config-to /path/to/config.yaml,Setup 会将合并后的完整配置写入该文件。这是一个实用的运维工具,可以快速生成配置模板。
第四,SchedulingQueue 的事件驱动重调度。addAllEventHandlers 中为 StorageClass、PV、PVC 等资源注册的 MoveAllToActiveOrBackoffQueue 处理器,实际上实现了事件驱动的"动态重调度"。当一个 PV 被绑定时,所有等待该 PV 的 Pod 会从 unschedulablePods 被移回 activeQ,重新参与调度。这种设计避免了定时轮询的资源浪费。
第五,NominatedNodeName 的优化。当一个 Pod 因为资源不足暂时无法调度时,调度器会设置它的 nominatedNodeName 字段,暗示下次调度时应该优先尝试这个节点(因为其他 Pod 可能已经释放了资源)。findNodesThatFitPod 在开始遍历所有节点前,会先尝试 nominatedNodeName,这避免了每次都从第一个节点开始搜索的低效。
注意
如果你需要从头调试 kube-scheduler 启动过程,建议按以下顺序阅读源码:main.go → NewSchedulerCommand → runCommand → Setup → scheduler.New → Run → ScheduleOne。每个函数的输入输出都要清楚——前一个函数的返回值就是后一个函数的输入参数。
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。