基于 Kubernetes v1.37 源码(
pkg/scheduler/)
一、调度器的职责 @
调度器的核心任务只有一个:为待调度的 Pod 找到最合适的节点。
用户创建 Pod(未指定 nodeName)
↓
API Server 将 Pod 写入 etcd
↓
调度器监听到新 Pod
↓
为 Pod 选择最优节点
↓
将 Pod 绑定到该节点
二、整体架构 @
调度器由四大核心组件组成:
┌─────────────────────────────────────────────────────────────┐
│ SchedulingQueue │
│ ┌──────────┐ ┌──────────┐ ┌──────────────────────────┐ │
│ │ activeQ │ │ backoffQ │ │ unschedulableEntities │ │
│ │ 活跃队列 │ │ 退避队列 │ │ 不可调度的 Pod │ │
│ └──────────┘ └──────────┘ └──────────────────────────┘ │
└─────────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────────┐
│ Cache │
│ 维护集群节点状态的本地快照(NodeInfo Snapshot) │
│ 避免每次调度都查询 API Server │
└─────────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────────┐
│ Framework │
│ 插件框架,定义调度周期的各个扩展点 │
│ PreEnqueue → PreFilter → Filter → PostFilter → PreScore → │
│ Score → Reserve → Permit → PreBind → Bind → PostBind │
└─────────────────────────────────────────────────────────────┘
三、调度队列(SchedulingQueue) @
调度队列管理所有等待调度的 Pod,由三个子队列组成:
// scheduling_queue.go
type PriorityQueue struct {
activeQ activeQueuer // 活跃队列(即将被调度)
backoffQ backoffQueuer // 退避队列(调度失败后等待重试)
unschedulableEntities *unschedulableEntities // 不可调度的 Pod(暂时无法调度)
}
3.1 三个队列的职责 @
| 队列 | 作用 | 何时进入 | 何时移出 |
|---|---|---|---|
| activeQ | 存放即将被调度的 Pod | Pod 创建、从 backoffQ/unschedulableEntities 移入 | 调度成功或移入其他队列 |
| backoffQ | 存放调度失败需退避的 Pod | 调度失败后 | 退避时间到期,移入 activeQ |
| unschedulableEntities | 存放暂时无法调度的 Pod | 调度失败(无可用节点) | 集群事件触发时重新激活 |
3.2 出队方式 @
调度器的 Pop 操作只从 activeQ 取出 Pod 进行调度,三个队列之间并不存在优先级排序机制:
- backoffQ 中的 Pod 退避到期后被移入 activeQ,再从 activeQ 出队
- unschedulableEntities 中的 Pod 等待集群事件触发重新激活,从不被直接 pop
3.3 退避机制 @
调度失败后,Pod 会进入 backoffQ,按指数退避等待:
DefaultPodInitialBackoffDuration = 1 * time.Second // 初始退避 1s
DefaultPodMaxBackoffDuration = 10 * time.Second // 最大退避 10s
退避时间 = min(初始退避 × 2^重试次数, 最大退避)
3.4 激活机制 @
即使 Pod 在 unschedulableEntities 中,当集群发生变化时,Pod 会被重新激活:
// 当节点资源变化、Pod 被删除等事件发生时(参数已简化)
func (p *PriorityQueue) MoveAllToActiveOrBackoffQueue(logger, event, oldObj, newObj, preCheck) {
// 将 unschedulableEntities 中的 Pod 移入 activeQ 或 backoffQ
}
四、Cache 层 @
Cache 维护了集群中所有节点的本地快照,避免每次调度都查询 API Server。
4.1 NodeInfo @
// v1.37 中 NodeInfo 是接口,通过方法访问节点信息
type NodeInfo interface {
Node() *v1.Node // 节点信息
GetUsedPorts() HostPortInfo // 已使用端口
GetImageStates() map[string]*ImageStateSummary // 镜像状态
GetPods() []PodInfo // 节点上的 Pod 列表
// ...
}
4.2 Snapshot @
调度周期开始时,Cache 会生成一个快照(Snapshot),保证整个调度周期内看到一致的节点状态:
func (c *cacheImpl) UpdateSnapshot(logger, nodeInfoSnapshot) error {
// 更新节点快照
}
4.3 为什么需要 Cache? @
调度器每秒可能处理大量 Pod,如果每次都查询 API Server:
- API Server 和 etcd 压力巨大
- 延迟高,吞吐量低
使用本地 Cache 后:
- 调度决策基于本地内存,微秒级响应
- 通过 Informer 监听变化,保持最终一致性
五、Framework 插件框架 @
调度器的核心是一个可插拔的框架,定义了调度周期的多个扩展点。
5.1 扩展点一览 @
PreEnqueue → PreFilter → Filter → PostFilter → PreScore → Score → Reserve → Permit → PreBind → Bind → PostBind
| 扩展点 | 作用 | 失败后果 |
|---|---|---|
| PreEnqueue | Pod 入队前检查 | Pod 不入队 |
| PreFilter | 预处理(计算 Pod 需求) | 调度失败 |
| Filter | 过滤不满足条件的节点 | 节点被淘汰 |
| PostFilter | 过滤后处理(如抢占) | 触发抢占 |
| PreScore | 评分前准备 | 调度失败 |
| Score | 为节点打分 | 影响最终排名 |
| Reserve | 预留资源 | 调度失败 |
| Permit | 准入检查(等待/拒绝/允许) | Pod 等待或拒绝 |
| PreBind | 绑定前准备 | 调度失败 |
| Bind | 绑定 Pod 到节点 | 调度失败 |
| PostBind | 绑定后处理 | 记录日志 |
5.2 插件接口 @
每个插件需要实现对应的接口:
type Plugin interface {
Name() string
}
type PreFilterPlugin interface {
Plugin
PreFilter(ctx, state, pod) *Status
}
type FilterPlugin interface {
Plugin
Filter(ctx, state, pod, nodeInfo) *Status
}
type ScorePlugin interface {
Plugin
Score(ctx, state, pod, nodeInfo) (int64, *Status)
ScoreExtensions() ScoreExtensions
}
5.3 内置插件 @
k8s 内置了 20+ 插件,覆盖常见的调度需求:
| 插件 | 扩展点 | 作用 |
|---|---|---|
| NodeResourcesFit | Filter, Score | 检查节点资源是否充足 |
| NodeAffinity | Filter, Score | 节点亲和性 |
| PodTopologySpread | Filter, Score | Pod 拓扑分布 |
| InterPodAffinity | Filter, Score | Pod 间亲和/反亲和 |
| TaintToleration | Filter, Score | 污点与容忍 |
| NodePorts | Filter | 端口冲突检查 |
| NodeUnschedulable | Filter | 过滤不可调度节点 |
| NodeName | Filter | 指定节点名 |
| VolumeBinding | Filter, Score | 存储卷绑定 |
| GangScheduling | PreEnqueue, Permit | Gang 调度(all-or-nothing,全部或全无) |
| DefaultBinder | Bind | 默认绑定器 |
六、调度周期详解 @
一次完整的调度分为两个阶段:调度周期(Scheduling Cycle) 和 绑定周期(Binding Cycle)。
6.1 调度周期 @
func (sched *Scheduler) schedulingCycle(ctx, state, fwk, podInfo) {
// 1. 更新节点快照
sched.Cache.UpdateSnapshot(logger, sched.nodeInfoSnapshot)
// 2. 执行调度算法
scheduleResult, status := sched.schedulingAlgorithm(ctx, state, fwk, podInfo)
// 3. 准备绑定周期(在内存中假设 Pod 已调度)
assumedPodInfo, status := sched.prepareForBindingCycle(...)
}
6.2 调度算法流程 @
func (sched *Scheduler) schedulingAlgorithm(ctx, state, fwk, podInfo) {
// 1. PreFilter:预处理
fwk.RunPreFilterPlugins(ctx, state, pod)
// 2. Filter:过滤节点(findNodesThatPassFilters 汇总可行节点)
feasibleNodes := sched.findNodesThatPassFilters(ctx, fwk, state, pod, nodeInfos)
// 3. PostFilter:调度算法返回 *framework.FitError(无可行节点)
// 且注册了 PostFilter 插件时触发,尝试抢占
if len(feasibleNodes) == 0 && fwk.HasPostFilterPlugins() {
fwk.RunPostFilterPlugins(ctx, state, pod, fitError.Diagnosis.NodeToStatus)
}
// 4. PreScore:评分前准备
fwk.RunPreScorePlugins(ctx, state, pod, feasibleNodes)
// 5. Score:打分
priorityList := fwk.RunScorePlugins(ctx, state, pod, feasibleNodes)
// 6. 选择得分最高的节点(schedulePod 中)
selectedNode := framework.NewSortedScoredNodes(priorityList).Pop().Name
return ScheduleResult{selectedNode}
}
6.3 绑定周期 @
绑定周期是异步执行的,不阻塞调度周期处理下一个 Pod:
// 异步执行绑定
go sched.runBindingCycle(ctx, state, fwk, scheduleResult, assumedPodInfo, start)
注意 Reserve 和 Permit 扩展点在调度周期内同步执行(schedulingCycle → prepareForBindingCycle → assumeAndReserve),绑定周期只做剩余的绑定工作:
// 调度周期内(prepareForBindingCycle / assumeAndReserve)
fwk.RunReservePluginsReserve(ctx, state, assumedPod, scheduleResult.SuggestedHost)
fwk.RunPermitPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)
func (sched *Scheduler) bindingCycle(ctx, state, fwk, scheduleResult, assumedPodInfo) {
// 1. PreBindPreFlight:绑定前置检查
fwk.RunPreBindPreFlights(ctx, state, assumedPod, scheduleResult.SuggestedHost)
// 2. WaitOnPermit:等待 Permit 插件的准入结果
fwk.WaitOnPermit(ctx, assumedPod)
// 3. PreBind:绑定前准备
fwk.RunPreBindPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)
// 4. Bind:执行绑定(写入 API Server)
fwk.RunBindPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)
// 5. PostBind:绑定后处理
fwk.RunPostBindPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)
}
6.4 为什么绑定要异步? @
调度周期是串行的(同一时刻只调度一个 Pod),但绑定可以异步:
- 调度周期需要快速决策,不能等待 API Server 响应
- 绑定是 IO 操作(写入 API Server),可以异步执行
- 通过"假设(Assume)“机制,在内存中预占资源,不影响后续调度决策
七、CycleState @
CycleState 是调度周期内的状态传递机制,插件之间可以通过它共享数据:
state := framework.NewCycleState()
state.Write("preFilterState", preFilterResult)
state.Read("preFilterState")
好处:
- 避免重复计算(如 PreFilter 计算一次,Filter 复用)
- 插件间解耦,不直接依赖
八、节点过滤详解 @
8.1 Filter 阶段 @
Filter 阶段对每个节点调用所有 FilterPlugin,检查节点是否满足条件,并汇总出可行节点列表:
// schedule_one.go:聚合可行节点
func (sched *Scheduler) findNodesThatPassFilters(ctx, fwk, state, pod, nodes) []NodeInfo {
feasibleNodes := make([]NodeInfo, numNodesToFind)
checkNode := func(i int) {
status := fwk.RunFilterPluginsWithNominatedPods(ctx, state, pod, nodes[i])
if status.IsSuccess() {
feasibleNodes[...] = nodes[i] // 收集可行节点
}
}
fwk.Parallelizer().Until(ctx, numAllNodes, checkNode, ...) // 并行检查各节点
return feasibleNodes
}
8.2 并行过滤 @
Filter 阶段的并行化对象是节点,而不是插件:findNodesThatPassFilters 通过 Parallelizer().Until 并行检查各个节点;而在单个节点内部,所有 Filter 插件仍然串行执行:
// runtime/framework.go:针对单个节点,串行执行所有 Filter 插件
func (f *frameworkImpl) RunFilterPlugins(ctx, state, pod, nodeInfo) *Status {
for _, pl := range f.filterPlugins {
if status := f.runFilterPlugin(ctx, pl, state, pod, nodeInfo); !status.IsSuccess() {
return status
}
}
return nil
}
九、节点评分详解 @
9.1 Score 阶段 @
Score 阶段为每个可行节点打分,分数范围 0~100:
// plugins/noderesources/fit.go:插件结构体名为 Fit
func (f *Fit) Score(ctx, state, pod, nodeInfo) (int64, *Status) {
// Score 直接接收 NodeInfo,无需再按 nodeName 查询快照
// 计算资源利用率
requested := f.calculatePodResourceRequestList(pod, f.resources)
capacity := nodeInfo.GetAllocatable() // 节点容量
// 分数 = (1 - 利用率) × 100
return score, nil
}
9.2 评分策略 @
NodeResourcesFit 通过 scoringStrategy.type 选择评分策略:
| 策略 | 说明 | 适用场景 |
|---|---|---|
| LeastAllocated | 优先选择资源利用率低的节点 | 避免热点 |
| MostAllocated | 优先选择资源利用率高的节点 | 提高资源利用率 |
| RequestedToCapacityRatio | 自定义资源权重比 | 精细控制 |
另有独立插件 NodeResourcesBalancedAllocation:优先选择 CPU/内存利用率最均衡的节点,避免资源倾斜。
9.3 多插件加权 @
每个插件有独立的权重,最终分数为加权和:
score := int64(0)
for _, pluginScore := range pluginScores {
score += pluginScore.Score * pluginScore.Weight
}
十、抢占机制(Preemption) @
当 Pod 找不到可调度节点时,调度器会尝试抢占:
10.1 触发条件 @
// 没有可行节点(*framework.FitError),且注册了 PostFilter 插件
if len(feasibleNodes) == 0 && fwk.HasPostFilterPlugins() {
fwk.RunPostFilterPlugins(ctx, state, pod, fitError.Diagnosis.NodeToStatus)
}
10.2 抢占流程 @
1. 找到可以抢占的节点(牺牲低优先级 Pod 后能容纳新 Pod)
↓
2. 选择牺牲者(Victims):优先级最低的 Pod
↓
3. 检查 PDB(PodDisruptionBudget):确保不违反中断预算
↓
4. 通知被抢占的 Pod 删除
↓
5. 被抢占的 Pod 被删除(驱逐),不会回到调度队列;由其控制器重建的新 Pod 作为新对象从 activeQ 开始排队
↓
6. 新 Pod 调度到该节点
10.3 关键代码 @
// preemption.go
func (ev *Evaluator) evaluate(ctx, state, pod, m) Candidate {
// 1. 找出所有可能的受害者
// 2. 检查 PDB 约束
// 3. 选择最优的候选节点
}
十一、Profile 机制 @
调度器支持多 Profile,每个 Profile 是独立的插件组合:
type Map map[string]framework.Framework
func NewMap(ctx, cfgs []KubeSchedulerProfile, r Registry, recorderFact RecorderFactory) (Map, error) {
m := make(Map)
for _, cfg := range cfgs {
p, _ := newProfile(ctx, cfg, r, recorderFact, opts...)
m[cfg.SchedulerName] = p
}
return m, nil
}
11.1 使用场景 @
- 不同业务使用不同的调度策略
- 自定义调度器名称:
schedulerName: my-scheduler
apiVersion: v1
kind: Pod
metadata:
name: my-pod
spec:
schedulerName: my-scheduler # 指定调度器 Profile
containers:
- name: nginx
image: nginx
十二、调度器主循环 @
func (sched *Scheduler) ScheduleOne(ctx) {
// 1. 从队列获取下一个 Pod
entity, err := sched.NextEntity(logger)
// 2. 根据实体类型分发
switch entity.(type) {
case *framework.QueuedPodInfo:
sched.scheduleOnePod(ctx, podInfo)
case *framework.QueuedPodGroupInfo:
sched.scheduleOnePodGroup(ctx, podGroupInfo)
}
}
func (sched *Scheduler) scheduleOnePod(ctx, podInfo) {
// 1. 获取 Framework
schedFramework, _ := sched.frameworkForPod(pod)
// 2. 初始化 CycleState
state := framework.NewCycleState()
// 3. 执行调度周期
scheduleResult, assumedPodInfo, status := sched.schedulingCycle(ctx, state, schedFramework, podInfo)
// 4. 异步执行绑定周期
go sched.runBindingCycle(ctx, state, schedFramework, scheduleResult, assumedPodInfo, start)
}
十三、完整流程示例 @
假设用户创建一个需要 2 核 4G 内存的 Pod:
1. API Server 创建 Pod(无 nodeName)
↓
2. 调度器 PreEnqueue 检查通过
↓
3. Pod 加入 activeQ
↓
4. ScheduleOne 取出 Pod
↓
5. PreFilter:计算 Pod 资源需求(2C4G)
↓
6. Filter:过滤节点
- NodeUnschedulable:过滤掉 cordoned 节点
- NodeResourcesFit:过滤掉可用资源 < 2C4G 的节点
- NodeAffinity:过滤掉标签不匹配的节点
- TaintToleration:过滤掉有不容忍污点的节点
- 剩余 10 个可行节点
↓
7. PreScore:准备评分数据
↓
8. Score:为 10 个节点打分
- NodeResourcesFit:按资源利用率打分
- NodeAffinity:按标签匹配度打分
- 加权求和,排序
↓
9. 选择得分最高的节点:node-3
↓
10. Reserve:在 Cache 中预留 node-3 的资源
↓
11. Permit:准入检查(通过)
↓
12. PreBind:准备绑定(如等待 PV 绑定)
↓
13. Bind:调用 API Server 创建 Binding 对象
↓
14. API Server 更新 Pod 的 nodeName 字段
↓
15. kubelet 监听到 Pod 被分配到本节点,启动容器
十四、设计亮点总结 @
- 可插拔架构:插件框架设计使得调度器可以灵活扩展,无需修改核心代码。
- 三段式队列:activeQ + backoffQ + unschedulableEntities 的设计,平衡了调度效率和公平性。
- 快照机制:Cache 快照保证整个调度周期看到一致的节点状态,避免"读到旧数据”。
- 异步绑定:绑定异步执行,不阻塞调度周期,提高吞吐量。
- 抢占机制:通过抢占低优先级 Pod,保证高优先级 Pod 的调度成功。
- 多 Profile:支持多套调度策略,满足不同业务需求。
十五、总结 @
Kubernetes 调度器的设计体现了几个核心思想:
- 分离关注点:调度与绑定分离,调度周期与绑定周期解耦
- 可插拔:插件框架支持灵活扩展
- 高效:本地缓存 + 快照,避免频繁查询 API Server
- 公平:退避机制 + 激活机制,避免 Pod 饥饿
- 优先级:抢占机制保证高优先级 Pod
理解调度器,就掌握了 Pod 生命周期的关键环节。后续的 kubelet、容器运行时,都是调度结果的执行者。
下一篇:Kubernetes API Server——请求处理链路