Kubernetes调度器源码分析与自定义调度

引言

凌晨两点,我被一通电话从睡梦中拽了起来——生产环境的Kubernetes集群出现了大量Pending状态的Pod,用户反馈服务已经不可用超过十分钟。

登录集群查看,节点资源明明充足,但Pod就像被施了定身咒一样卡在调度队列里。排查了半天,发现是某个节点被打上了node.kubernetes.io/unschedulable: NoSchedule污点,而我们的工作负载恰好没有容忍这个污点。

那一刻我意识到,虽然每天在用Kubernetes,但调度器对我而言始终是个黑盒。如果当时能深入理解调度器的内部机制,这个故障排查可能只需要几分钟。

这就是我想写这篇文章的初衷——带你深入Kubernetes调度器的核心源码,并手把手实现一个自定义调度器。读完这篇文章,你将能够回答以下问题:

  • Pod从创建到绑定节点,中间经历了什么?
  • 调度器的"预选"和"优选"策略到底是如何实现的?
  • 如何在不修改kube-scheduler源码的情况下,扩展调度能力?

核心概念

生活类比:高考录取系统

想象一下Kubernetes调度器是一个高考录取系统,而Pod就是考生,节点就是大学。

调度过程就像高考录取:

  1. 资格审查(Predicates/预选):考生分数是否达到最低投档线?是否有身体条件限制?这对应调度器检查Pod的请求资源是否满足节点可用资源、端口是否冲突、节点亲和性是否匹配等硬性条件。
  1. 综合评定(Priorities/优选):在过了投档线的考生中,按分数从高到低排序。对应调度器对每个满足条件的节点进行打分,选择得分最高的节点。
  1. 录取通知(Binding/绑定):将考生档案投递到目标大学。对应调度器创建Pod与节点的绑定关系。

这套设计的好处在于,预选保证可行性,优选保证最优性,两者解耦,各自可以独立扩展。

技术定义

Kubernetes调度器(kube-scheduler)是控制平面的核心组件之一,负责将Pending状态的Pod分配到最合适的节点上。它通过监听API Server获取待调度的Pod,经过调度算法(过滤+打分)选定目标节点,最后通过绑定操作将Pod与节点关联。

调度器架构

graph TD A[API Server] -->|Watch Pods| B[Informer] B -->|待调度Pod进入队列| C[Scheduling Queue] C -->|弹出Pod| D[调度主循环] subgraph 调度周期 D --> E[预选阶段
Filter/过滤] E --> F{是否通过?} F -->|否| G[标记不可调度
等待下次重试] F -->|是| H[优选阶段
Score/打分] H --> I[选择最高分节点] end I -->|假定绑定| J[Assume Pod] J -->|异步执行| K[绑定周期
Bind Pod到Node] subgraph 绑定周期 K --> L[Volume绑定] L --> M[Pod绑定API调用] M --> N[清理缓存] end D --> O{调度成功?} O -->|是| P[更新调度结果] O -->|否| Q[重新入队]

源码深度分析

调度器主流程

kube-scheduler的入口在cmd/kube-scheduler/app/server.go,但核心逻辑在pkg/scheduler/scheduler.go。我们来看关键的scheduleOne方法:

// pkg/scheduler/scheduler.go
func (sched *Scheduler) scheduleOne(ctx context.Context) {
    // 1. 从调度队列中获取下一个待调度的Pod
    podInfo := sched.NextPod()
    // podInfo 可能是通过优先级抢占来的,也可能是普通Pod
    
    // 2. 执行调度周期(Scheduling Cycle)
    scheduleResult, err := sched.schedulePod(ctx, fwk, podInfo.Pod)
    if err != nil {
        // 调度失败处理
        sched.handleSchedulingFailure(ctx, fwk, podInfo, err, ...)
        return
    }
    
    // 3. 假定绑定(Assume)
    // 这是关键优化:在真正持久化绑定之前,先更新内部缓存,
    // 避免其他Pod重复选择同一节点
    assumedPodInfo := podInfo.DeepCopy()
    assumedPod := assumedPodInfo.Pod
    err = sched.assume(assumedPod, scheduleResult.SuggestedHost)
    
    // 4. 执行绑定周期(Binding Cycle)
    // 绑定周期与调度周期是并发执行的
    go func() {
        err := sched.bindPod(ctx, fwk, assumedPod, scheduleResult)
    }()
}

这里最核心的设计是调度周期与绑定周期的分离。调度周期是串行的(同一时间只调度一个Pod),而绑定周期可以并行执行。这类似于餐厅后厨:点单(调度)必须一单一单来,但做菜(绑定)可以同时进行。

预选阶段——Filter插件

预选阶段的核心是NodeResourcesFit插件,它检查节点是否有足够的资源满足Pod的请求。我们深入看一下资源匹配的逻辑:

// pkg/scheduler/framework/plugins/noderesources/fit.go
func (f *Fit) Filter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeInfo *framework.NodeInfo) *framework.Status {
    // 获取节点上的已分配资源
    allocatedResources := nodeInfo.AllocatableResources()
    
    // 检查每个容器的资源请求
    for _, container := range pod.Spec.Containers {
        // CPU、内存、存储等资源逐一检查
        if !allocatedResources.CPU.MilliValue() >= container.Resources.Requests.Cpu().MilliValue() {
            return framework.NewStatus(framework.Unschedulable, "Insufficient cpu")
        }
    }
}

一个重要细节NodeInfo是调度器维护的节点资源缓存。当Pod被Assume后,调度器会乐观地更新这个缓存,即假设Pod已经在该节点上运行了。这样可以避免重复调度,但也引入了缓存与实际状态不一致的风险。

优选阶段——Score插件

优选阶段有多个打分插件,我们重点看NodeResourcesBalancedAllocation(资源均衡分配):

// pkg/scheduler/framework/plugins/noderesources/balanced_allocation.go
func (pl *BalancedAllocation) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
    nodeInfo, err := pl.handle.SnapshotSharedLister().NodeInfos().Get(nodeName)
    
    // 计算CPU和内存的使用率
    cpuFraction := fractionOfCapacity(requestedCPU, allocatableCPU)
    memoryFraction := fractionOfCapacity(requestedMemory, allocatableMemory)
    
    // 得分 = 10 - 资源使用率的方差 * 10
    // 方差越小,说明资源使用越均衡,得分越高
    variance := math.Abs(cpuFraction - memoryFraction)
    score := (1 - variance) * maxScore
    
    return int64(score), nil
}

这个插件的设计很巧妙——它鼓励将Pod调度到资源使用均衡的节点上,避免出现CPU打满而内存空闲的"偏科"节点。

调度队列的优先级管理

调度队列不是简单的FIFO,它实现了基于优先级的调度。Pod的优先级通过PriorityClass定义:

// pkg/scheduler/internal/queue/scheduling_queue.go
func (p *PriorityQueue) Add(pod *v1.Pod) error {
    // 根据Pod优先级插入到不同的子队列
    if p.podBackoffQ != nil && p.podBackoffQ.Has(pod) {
        // 如果Pod正在退避期(调度失败后),更新其退避时间
        return p.podBackoffQ.Update(pod)
    }
    
    // 高优先级Pod进入activeQ的前端
    // 低优先级Pod可能被抢占
    return p.activeQ.Add(pod)
}

实战:自定义调度器

示例一:基于Webhook的调度器

有时候,我们希望在外部系统中完成调度决策——比如根据业务指标、成本数据等。Webhook调度器允许我们将决策逻辑完全外置。

// custom-scheduler/main.go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"io/ioutil"
	"net/http"
	"time"

	v1 "k8s.io/api/core/v1"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/client-go/kubernetes"
	"k8s.io/client-go/tools/clientcmd"
)

// WebhookScheduler 是一个通过外部HTTP服务做调度决策的调度器
type WebhookScheduler struct {
	client       kubernetes.Interface
	webhookURL   string
	pollInterval time.Duration
}

// ScheduleRequest 发送给Webhook服务的请求体
type ScheduleRequest struct {
	Pod      *v1.Pod `json:"pod"`
	Nodes    []v1.Node `json:"nodes"`
	Metadata map[string]string `json:"metadata"`
}

// ScheduleResponse Webhook服务的响应体
type ScheduleResponse struct {
	NodeName string `json:"nodeName"`
	Reason   string `json:"reason,omitempty"`
}

// NewWebhookScheduler 创建基于Webhook的调度器
func NewWebhookScheduler(kubeconfig, webhookURL string) (*WebhookScheduler, error) {
	config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
	if err != nil {
		return nil, fmt.Errorf("加载kubeconfig失败: %v", err)
	}
	
	clientset, err := kubernetes.NewForConfig(config)
	if err != nil {
		return nil, fmt.Errorf("创建Kubernetes客户端失败: %v", err)
	}
	
	return &WebhookScheduler{
		client:       clientset,
		webhookURL:   webhookURL,
		pollInterval: 2 * time.Second,
	}, nil
}

// Run 启动调度器主循环
func (s *WebhookScheduler) Run(ctx context.Context) error {
	fmt.Println("Webhook调度器已启动,开始监听待调度的Pod...")
	
	// 持续监听未调度的Pod
	watcher, err := s.client.CoreV1().Pods("").Watch(ctx, metav1.ListOptions{
		FieldSelector: "spec.nodeName=", // 只关注未绑定节点的Pod
	})
	if err != nil {
		return err
	}
	
	for event := range watcher.ResultChan() {
		pod, ok := event.Object.(*v1.Pod)
		if !ok {
			continue
		}
		
		// 忽略已完成或删除的Pod
		if pod.DeletionTimestamp != nil {
			continue
		}
		
		// 调度决策
		nodeName, err := s.schedule(ctx, pod)
		if err != nil {
			fmt.Printf("调度Pod %s/%s 失败: %v\n", pod.Namespace, pod.Name, err)
			continue
		}
		
		// 绑定Pod到节点
		if err := s.bind(ctx, pod, nodeName); err != nil {
			fmt.Printf("绑定Pod %s/%s 到节点 %s 失败: %v\n", 
				pod.Namespace, pod.Name, nodeName, err)
		}
	}
	
	return nil
}

// schedule 调用外部Webhook服务进行调度决策
func (s *WebhookScheduler) schedule(ctx context.Context, pod *v1.Pod) (string, error) {
	// 获取所有可用节点
	nodes, err := s.client.CoreV1().Nodes().List(ctx, metav1.ListOptions{})
	if err != nil {
		return "", err
	}
	
	// 构造请求
	req := ScheduleRequest{
		Pod:      pod,
		Nodes:    nodes.Items,
		Metadata: map[string]string{
			"requestTime": time.Now().Format(time.RFC3339),
		},
	}
	
	// 序列化请求
	reqBody, err := json.Marshal(req)
	if err != nil {
		return "", err
	}
	
	// 调用Webhook服务
	httpReq, err := http.NewRequestWithContext(ctx, "POST", s.webhookURL, bytes.NewBuffer(reqBody))
	if err != nil {
		return "", err
	}
	httpReq.Header.Set("Content-Type", "application/json")
	
	resp, err := http.DefaultClient.Do(httpReq)
	if err != nil {
		return "", fmt.Errorf("调用Webhook服务失败: %v", err)
	}
	defer resp.Body.Close()
	
	if resp.StatusCode != http.StatusOK {
		return "", fmt.Errorf("Webhook服务返回非200状态码: %d", resp.StatusCode)
	}
	
	// 解析响应
	var result ScheduleResponse
	if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
		return "", err
	}
	
	if result.NodeName == "" {
		return "", fmt.Errorf("Webhook未返回有效的节点名称: %s", result.Reason)
	}
	
	return result.NodeName, nil
}

// bind 将Pod绑定到指定节点
func (s *WebhookScheduler) bind(ctx context.Context, pod *v1.Pod, nodeName string) error {
	binding := &v1.Binding{
		ObjectMeta: metav1.ObjectMeta{
			Name:      pod.Name,
			Namespace: pod.Namespace,
		},
		Target: v1.ObjectReference{
			Kind: "Node",
			Name: nodeName,
		},
	}
	
	// 通过API Server创建Binding对象
	err := s.client.CoreV1().Pods(pod.Namespace).Bind(ctx, binding, metav1.CreateOptions{})
	if err != nil {
		return fmt.Errorf("创建绑定失败: %v", err)
	}
	
	fmt.Printf("成功将Pod %s/%s 调度到节点 %s\n", pod.Namespace, pod.Name, nodeName)
	return nil
}

func main() {
	ctx := context.Background()
	
	scheduler, err := NewWebhookScheduler(
		"/root/.kube/config", 
		"http://localhost:8080/schedule",
	)
	if err != nil {
		panic(err)
	}
	
	if err := scheduler.Run(ctx); err != nil {
		panic(err)
	}
}

示例二:调度器扩展(Scheduler Extender)

如果不想完全替换调度器,而是希望扩展现有kube-scheduler的能力,可以使用Scheduler Extender机制。它允许在预选和优选阶段之后,追加自定义的过滤和打分逻辑。

// extender/main.go
package main

import (
	"encoding/json"
	"fmt"
	"net/http"
	"sort"
	"time"

	v1 "k8s.io/api/core/v1"
	"k8s.io/apimachinery/pkg/types"
)

// ExtenderArgs 是kube-scheduler传给extender的参数
type ExtenderArgs struct {
	Pod      *v1.Pod             `json:"pod"`
	Nodes    *v1.NodeList        `json:"nodes"`
	NodeNames *[]string          `json:"nodeNames"`
}

// ExtenderFilterResult 是预选阶段的返回结果
type ExtenderFilterResult struct {
	Nodes       *v1.NodeList `json:"nodes,omitempty"`
	NodeNames   *[]string    `json:"nodeNames,omitempty"`
	FailedNodes FailedNodesMap `json:"failedNodes,omitempty"`
	Error       string       `json:"error,omitempty"`
}

type FailedNodesMap map[string]string

// ExtenderPriorityResult 是优选阶段的返回结果
type ExtenderPriorityResult struct {
	// 每个节点的打分结果
	Priorities []HostPriority `json:"priorities"`
}

type HostPriority struct {
	Host string `json:"host"`
	Score int64  `json:"score"`
}

// SchedulerExtender 实现了一个自定义的调度器扩展
type SchedulerExtender struct {
	// 可配置的权重
	metadataWeight int64
}

// NewSchedulerExtender 创建扩展实例
func NewSchedulerExtender() *SchedulerExtender {
	return &SchedulerExtender{
		metadataWeight: 10,
	}
}

// Filter 预选阶段扩展
// 这里我们实现一个简单的逻辑:排除带有特定标签的节点
func (e *SchedulerExtender) Filter(w http.ResponseWriter, r *http.Request) {
	var args ExtenderArgs
	if err := json.NewDecoder(r.Body).Decode(&args); err != nil {
		http.Error(w, err.Error(), http.StatusBadRequest)
		return
	}
	
	result := ExtenderFilterResult{
		FailedNodes: make(FailedNodesMap),
	}
	
	if args.Nodes != nil {
		filteredNodes := &v1.NodeList{}
		for _, node := range args.Nodes.Items {
			// 排除带有"exclude-from-scheduling=true"标签的节点
			if val, ok := node.Labels["exclude-from-scheduling"]; ok && val == "true" {
				result.FailedNodes[node.Name] = "节点被标签排除"
				continue
			}
			
			// 排除处于NotReady状态的节点
			for _, condition := range node.Status.Conditions {
				if condition.Type == v1.NodeReady && condition.Status != v1.ConditionTrue {
					result.FailedNodes[node.Name] = fmt.Sprintf("节点状态: %s", condition.Status)
					continue
				}
			}
			
			filteredNodes.Items = append(filteredNodes.Items, node)
		}
		result.Nodes = filteredNodes
	}
	
	w.Header().Set("Content-Type", "application/json")
	json.NewEncoder(w).Encode(result)
}

// Prioritize 优选阶段扩展
// 这里我们实现一个基于节点CPU架构的打分策略
func (e *SchedulerExtender) Prioritize(w http.ResponseWriter, r *http.Request) {
	var args ExtenderArgs
	if err := json.NewDecoder(r.Body).Decode(&args); err != nil {
		http.Error(w, err.Error(), http.StatusBadRequest)
		return
	}
	
	// 假设Pod需要ARM架构
	needARM := false
	if args.Pod.Labels["architecture"] == "arm64" {
		needARM = true
	}
	
	var priorities []HostPriority
	if args.Nodes != nil {
		for _, node := range args.Nodes.Items {
			score := int64(50) // 基础分
			
			// ARM架构匹配则加分
			if needARM {
				if node.Status.NodeInfo.Architecture == "arm64" {
					score += e.metadataWeight * 5
				}
			} else if node.Status.NodeInfo.Architecture == "amd64" {
				score += e.metadataWeight * 3
			}
			
			// 节点上运行的Pod数量越少,得分越高
			podCount := len(node.Status.Images) // 简化的示例
			score -= int64(podCount) * 2
			
			priorities = append(priorities, HostPriority{
				Host:  node.Name,
				Score: score,
			})
		}
	}
	
	// 按分数排序
	sort.Slice(priorities, func(i, j int) bool {
		return priorities[i].Score > priorities[j].Score
	})
	
	result := ExtenderPriorityResult{
		Priorities: priorities,
	}
	
	w.Header().Set("Content-Type", "application/json")
	json.NewEncoder(w).Encode(result)
}

func main() {
	extender := NewSchedulerExtender()
	
	http.HandleFunc("/filter", extender.Filter)
	http.HandleFunc("/prioritize", extender.Prioritize)
	
	fmt.Println("Scheduler Extender已启动,监听端口 8888...")
	http.ListenAndServe(":8888", nil)
}

示例三:使用调度器框架(Framework)编写扩展插件

Kubernetes 1.15+提供了调度器框架(Scheduling Framework),这是最现代、最推荐的自定义调度方式。它允许你在不修改kube-scheduler源码的情况下,通过编译期注册或运行时配置来扩展调度器。

// plugins/customfilter/customfilter.go
package customfilter

import (
	"context"
	"fmt"
	"strings"

	v1 "k8s.io/api/core/v1"
	"k8s.io/apimachinery/pkg/runtime"
	"k8s.io/kubernetes/pkg/scheduler/framework"
)

// CustomFilter 自定义的过滤插件
// 这个插件会检查Pod是否声明了特定的资源需求,并验证节点是否支持
type CustomFilter struct {
	handle framework.Handle
	// 支持的资源类型列表
	supportedResources map[string]bool
}

// New 创建插件实例
func New(obj runtime.Object, handle framework.Handle) (framework.Plugin, error) {
	// 从配置中解析支持的资源类型
	// 这里使用硬编码作为示例
	supported := map[string]bool{
		"gpu":    true,
		"fpga":   true,
		"tpu":    true,
	}
	
	return &CustomFilter{
		handle:             handle,
		supportedResources: supported,
	}, nil
}

// Name 返回插件名称
func (cf *CustomFilter) Name() string {
	return "CustomFilter"
}

// Filter 实现预选阶段逻辑
func (cf *CustomFilter) Filter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeInfo *framework.NodeInfo) *framework.Status {
	node := nodeInfo.Node()
	
	// 检查Pod是否请求了特殊资源
	for _, container := range pod.Spec.Containers {
		for resourceName := range container.Resources.Requests {
			resourceStr := string(resourceName)
			
			// 检查是否为自定义的扩展资源
			if strings.Contains(resourceStr, "example.com/") {
				// 提取资源类型
				resourceType := strings.TrimPrefix(resourceStr, "example.com/")
				
				// 检查节点是否支持这种资源
				if !cf.supportedResources[resourceType] {
					return framework.NewStatus(
						framework.Unschedulable,
						fmt.Sprintf("节点 %s 不支持资源类型 %s", node.Name, resourceType),
					)
				}
				
				// 检查节点上的资源容量是否足够
				allocatable, ok := node.Status.Allocatable[resourceName]
				if !ok || allocatable.IsZero() {
					return framework.NewStatus(
						framework.Unschedulable,
						fmt.Sprintf("节点 %s 上没有可用的 %s 资源", node.Name, resourceType),
					)
				}
			}
		}
	}
	
	return framework.NewStatus(framework.Success, "")
}

// Score 实现优选阶段逻辑
func (cf *CustomFilter) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
	// 获取节点信息
	nodeInfo, err := cf.handle.SnapshotSharedLister().NodeInfos().Get(nodeName)
	if err != nil {
		return 0, framework.NewStatus(framework.Error, fmt.Sprintf("获取节点信息失败: %v", err))
	}
	
	node := nodeInfo.Node()
	score := int64(0)
	
	// 检查节点上的自定义资源利用率
	for resourceName, quantity := range node.Status.Allocatable {
		if strings.Contains(string(resourceName), "example.com/") {
			// 计算资源使用率
			requested := nodeInfo.Requested.Resource(resourceName)
			allocatable := quantity.Value()
			
			if allocatable > 0 {
				usageRatio := float64(requested.Value()) / float64(allocatable)
				// 资源使用率越低,得分越高(鼓励分散调度)
				score += int64((1 - usageRatio) * 100)
			}
		}
	}
	
	return score, framework.NewStatus(framework.Success, "")
}

// ScoreExtensions 返回打分扩展
func (cf *CustomFilter) ScoreExtensions() framework.ScoreExtensions {
	return nil
}

配置插件到调度器:

# scheduler-config.yaml
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
profiles:
- schedulerName: custom-scheduler
  plugins:
    filter:
      enabled:
      - name: "CustomFilter"
      disabled:
      - name: "NodeAffinity"  # 可以禁用内置插件
    score:
      enabled:
      - name: "CustomFilter"
        weight: 50
      - name: "NodeResourcesBalancedAllocation"
        weight: 30

然后通过命令行启动自定义调度器:

kube-scheduler --config=scheduler-config.yaml --v=5

方案对比

四种自定义调度方案对比

| 方案 | 复杂度 | 灵活性 | 维护成本 | 适用场景 |

|------|--------|--------|----------|----------|

| Webhook调度器 | 低 | 高 | 中 | 调度逻辑完全外置,需要集成外部系统 |

| Scheduler Extender | 中 | 中 | 中 | 在现有调度器基础上追加过滤/打分逻辑 |

| 调度器框架插件 | 中 | 高 | 低 | 需要深度定制调度行为,且希望与K8s版本兼容 |

| 修改kube-scheduler源码 | 高 | 最高 | 高 | 需要修改调度核心逻辑,且有足够人力维护 |

关键差异分析

调度器框架是官方推荐的方式,因为:

  1. 版本兼容性:框架API相对稳定,升级K8s版本时迁移成本低
  2. 可组合性:可以同时启用多个插件,每个插件关注一个关注点
  3. 可测试性:插件可以独立单元测试

Scheduler Extender的优势在于它通过HTTP接口工作,可以用任何语言实现,但性能开销较大——每次调度都要走一次HTTP调用。

完全替换调度器(Webhook方式)提供了最大的灵活性,但你需要自己实现调度队列管理失败重试抢占等复杂功能,工作量不容小觑。

最佳实践与避坑指南

实践一:调度性能调优

调度器的性能瓶颈往往不在算法本身,而在API Server的请求压力。以下是一些优化建议:

// 优化示例:批量获取节点信息
// 不要这样(逐个查询节点):
for _, nodeName := range nodeNames {
    node, _ := client.CoreV1().Nodes().Get(ctx, nodeName, metav1.GetOptions{})
}

// 应该这样(一次获取所有节点):
nodes, _ := client.CoreV1().Nodes().List(ctx, metav1.ListOptions{})

实践二:处理调度失败

调度失败的处理需要特别注意重试策略退避算法,避免出现"调度风暴":

// 合理设置重试退避时间
type RetryPolicy struct {
    MaxRetries   int           // 最大重试次数
    InitialDelay time.Duration // 初始延迟
    MaxDelay     time.Duration // 最大延迟
}

func (p *RetryPolicy) GetBackoff(retryCount int) time.Duration {
    // 指数退避 + 抖动
    delay := p.InitialDelay * time.Duration(1<<uint(retryCount))
    if delay > p.MaxDelay {
        delay = p.MaxDelay
    }
    
    // 添加±10%的随机抖动,避免惊群效应
    jitter := time.Duration(float64(delay) * 0.1 * rand.Float64())
    return delay + jitter
}

避坑指南

  1. 不要忽略Pod的优先级:自定义调度器需要正确处理PriorityClass,否则可能导致高优先级Pod饿死
  1. 注意缓存一致性:使用Assume机制时,务必在绑定失败后回滚缓存,否则会导致资源泄漏
  1. 处理节点状态变化:调度器必须监听节点状态变化(如NotReady、资源不足),及时从调度候选列表中移除
  1. 避免调度热点:对于大规模集群,要考虑调度器的并发能力。内置的调度器是单实例的,如果遇到性能瓶颈,可以考虑使用多调度器(每个调度器负责一部分资源类型)
  1. 测试要充分:自定义调度器上线前,务必进行混沌测试——模拟节点宕机、网络分区、API Server不可用等场景

总结

通过这篇文章,我们从源码层面剖析了Kubernetes调度器的核心机制:

  • 调度周期与绑定周期分离的设计,保证了调度的高吞吐
  • 预选与优选的两阶段策略,兼顾了可行性与最优性
  • 调度器框架的插件化设计,让我们可以灵活扩展调度能力

我们实现了三种自定义调度方案,从最简单的Webhook调度器到最优雅的框架插件,每种方案都有其适用场景。

延伸思考:Kubernetes社区正在讨论"多调度器"架构,即不同场景使用不同的调度器实例。例如,批处理任务使用一个调度器,在线服务使用另一个调度器。这种架构可以避免不同类型工作负载之间的调度干扰,但也带来了更复杂的集群管理挑战。

如果你正在使用Kubernetes,我强烈建议你花时间深入调度器源码。你会发现,那些看似神奇的功能,底层实现其实非常优雅且易于理解。正如开篇的故障排查经历所示,对调度器的深入理解,在关键时刻能让你从"排查半天"变成"几分钟搞定"。