基于 Kubernetes v1.37 源码(
pkg/kubelet/)
一、Kubelet 的职责 @
Kubelet 是运行在每个节点上的节点代理,其核心职责是管理 Pod 的生命周期:
调度器将 Pod 调度到节点
↓
kubelet 监听到 Pod 分配
↓
创建容器(通过 CRI)
↓
监控容器状态
↓
上报状态到 API Server
↓
处理删除/更新/驱逐
二、Kubelet 架构全景 @
┌────────────────────────────────────────────────────────────┐
│ Kubelet │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ Pod Sync Loop │ │
│ │ LIST/WATCH API Server → 对比期望与实际状态 → 执行 │ │
│ └──────────────────────────────────────────────────────┘ │
│ ↓ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ Container Manager (CM) │ │
│ │ ┌───────────┐ ┌───────────┐ ┌──────────────────┐ │ │
│ │ │ CRI Client │ │ cgroup │ │ Device Plugin │ │ │
│ │ │ 容器运行时 │ │ 资源隔离 │ │ 设备管理 │ │ │
│ │ └───────────┘ └───────────┘ └──────────────────┘ │ │
│ └──────────────────────────────────────────────────────┘ │
│ ↓ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ CRI (Container Runtime Interface) │ │
│ │ containerd / CRI-O / cri-dockerd │ │
│ └──────────────────────────────────────────────────────┘ │
└────────────────────────────────────────────────────────────┘
三、Pod Sync Loop(Pod 同步循环) @
Kubelet 的核心是一个持续运行的 Pod 同步循环,它监听多个通道来触发同步:
// pkg/kubelet/kubelet.go
func (kl *Kubelet) syncLoop(
ctx context.Context,
updates <-chan kubetypes.PodUpdate, // Pod 更新通道
handler SyncHandler, // 同步处理器
) {
for {
// 检查运行时错误,如有则指数退避
if err := kl.runtimeState.runtimeErrors(); err != nil {
time.Sleep(duration)
continue
}
// 监听多个通道,哪个有数据就处理哪个
kl.syncLoopIteration(ctx, updates, handler,
syncTicker.C, housekeepingTicker.C, plegCh)
}
}
3.1 syncLoopIteration @
// pkg/kubelet/kubelet.go
func (kl *Kubelet) syncLoopIteration(
ctx context.Context,
configCh <-chan kubetypes.PodUpdate, // 配置变更通道
handler SyncHandler,
syncCh <-chan time.Time, // 定期同步通道
housekeepingCh <-chan time.Time, // 清理通道
plegCh <-chan *pleg.PodLifecycleEvent, // Pod 生命周期事件通道
) bool {
select {
case u, open := <-configCh:
// 处理 Pod 配置变更(添加/更新/删除/调和)
// PodUpdate 只有三个字段:Pods []*v1.Pod、Op、Source
switch u.Op {
case kubetypes.ADD:
handler.HandlePodAdditions(ctx, u.Pods)
case kubetypes.UPDATE:
handler.HandlePodUpdates(ctx, u.Pods)
case kubetypes.REMOVE:
handler.HandlePodRemoves(ctx, u.Pods)
case kubetypes.RECONCILE:
handler.HandlePodReconcile(ctx, u.Pods)
}
case <-syncCh:
// 定期同步所有待同步的 Pod
podsToSync := kl.getPodsToSync()
handler.HandlePodSyncs(ctx, podsToSync)
case <-housekeepingCh:
// 清理已终止的 Pod、孤立的 Volume 和 Pod 目录等
handler.HandlePodCleanups(ctx)
case e := <-plegCh:
// Pod 生命周期事件(容器启动/停止),同步对应的 Pod
if pod, ok := kl.podManager.GetPodByUID(e.ID); ok {
handler.HandlePodSyncs(ctx, []*v1.Pod{pod})
}
}
}
3.2 Pod Workers(异步执行) @
Kubelet 没有 syncPod 方法,而是通过 podWorkers 异步管理每个 Pod:
// pkg/kubelet/pod_workers.go
type podWorkers struct {
podLock sync.Mutex
podsSynced bool // 所有工作中的 Pod 至少同步过一次后置为 true
// 每个 Pod 一个更新通道,向通道发信号通知对应 goroutine 消费更新
podUpdates map[types.UID]chan struct{}
// 按 UID 跟踪 Pod 的同步状态(syncing / terminating / terminated / evicted)
podSyncStatuses map[types.UID]*podSyncStatus
workQueue queue.WorkQueue
podSyncer podSyncer // 实际执行同步的接口,由 Kubelet 实现
}
// 更新 Pod(异步执行,不阻塞调用方)
func (p *podWorkers) UpdatePod(ctx context.Context, options UpdatePodOptions) {
p.podLock.Lock()
// 将更新写入该 Pod 的 podSyncStatus.pendingUpdate,
// 并向 podUpdates[uid] 通道发信号(goroutine 不存在则创建)
p.podLock.Unlock()
}
// 每个 Pod 的 goroutine 串行消费更新,按 kubetypes.SyncPodType 分类:
// SyncPodCreate / SyncPodUpdate / SyncPodSync / SyncPodKill
func (p *podWorkers) podWorkerLoop(uid types.UID) {
for range podUpdatesCh {
// 根据 podSyncStatus 状态机,调用 podSyncer 对应方法:
switch {
case 需要创建或更新:
p.podSyncer.SyncPod(ctx, updateType, pod, mirrorPod, podStatus)
case 需要终止:
p.podSyncer.SyncTerminatingPod(...) // 停止容器
case 已终止:
p.podSyncer.SyncTerminatedPod(...) // 清理资源
}
}
}
四、Pod 生命周期 @
4.1 Pod 状态机 @
Pending → Running → Succeeded/Failed(异常时为 Unknown)
| 状态 | 含义 |
|---|---|
| Pending | Pod 已被接受,但容器尚未创建 |
| Running | 所有容器已创建,至少一个在运行 |
| Succeeded | 所有容器成功终止 |
| Failed | 所有容器已终止,至少一个失败 |
| Unknown | 无法获取 Pod 状态 |
注意:Scheduled 不是 Pod Phase。Phase 只有上表五种;“已调度到节点"是通过 PodScheduled Condition 表示的(staging/src/k8s.io/api/core/v1/types.go)。
4.2 容器状态机 @
Waiting → Running → Terminated
kubelet 内部运行时层与 API 层对容器状态的表示不同:
// pkg/kubelet/container/runtime.go —— kubelet 内部运行时层
type State string
const (
ContainerStateCreated State = "created" // 已创建但未启动
ContainerStateRunning State = "running"
ContainerStateExited State = "exited"
ContainerStateUnknown State = "unknown" // restarting、paused、dead 等
)
// staging/src/k8s.io/api/core/v1/types.go —— API 层
// Waiting/Running/Terminated 是 v1.ContainerState 结构体的字段(互斥指针):
type ContainerState struct {
Waiting *ContainerStateWaiting
Running *ContainerStateRunning
Terminated *ContainerStateTerminated
}
五、CRI(容器运行时接口) @
Kubelet 通过 CRI 与容器运行时交互,解耦具体实现:
// k8s.io/cri-api/pkg/apis/runtime/v1/api.proto
service RuntimeService {
rpc RunPodSandbox(RunPodSandboxRequest) returns (RunPodSandboxResponse);
rpc StopPodSandbox(StopPodSandboxRequest) returns (StopPodSandboxResponse);
rpc RemovePodSandbox(RemovePodSandboxRequest) returns (RemovePodSandboxResponse);
rpc CreateContainer(CreateContainerRequest) returns (CreateContainerResponse);
rpc StartContainer(StartContainerRequest) returns (StartContainerResponse);
rpc StopContainer(StopContainerRequest) returns (StopContainerResponse);
rpc RemoveContainer(RemoveContainerRequest) returns (RemoveContainerResponse);
rpc ListContainers(ListContainersRequest) returns (ListContainersResponse);
rpc ContainerStatus(ContainerStatusRequest) returns (ContainerStatusResponse);
rpc ExecSync(ExecSyncRequest) returns (ExecSyncResponse);
rpc Attach(AttachRequest) returns (AttachResponse);
rpc PortForward(PortForwardRequest) returns (PortForwardResponse);
rpc UpdateContainerResources(UpdateContainerResourcesRequest) returns (UpdateContainerResourcesResponse);
}
5.1 容器创建流程 @
// pkg/kubelet/kuberuntime/kuberuntime_container.go
// Pod Sandbox(Pause 容器)已在 SyncPod 流程中先行创建,这里不再创建
func (m *kubeGenericRuntimeManager) startContainer(
ctx context.Context,
podSandboxID string,
podSandboxConfig *runtimeapi.PodSandboxConfig,
spec *startSpec,
pod *v1.Pod,
podStatus *kubecontainer.PodStatus,
pullSecrets []v1.Secret,
podIP string, podIPs []string,
imageVolumes kubecontainer.ImageVolumes,
) (string, error) {
// 1. 确保镜像存在(不存在则拉取)
imageRef, msg, err := m.imagePuller.EnsureImageExists(
ctx, ref, pod, container.Image, pullSecrets, podSandboxConfig, ...)
// 2. 生成容器配置并创建容器
containerConfig, cleanupAction, err := m.generateContainerConfig(...)
containerID, err := m.runtimeService.CreateContainer(
ctx, podSandboxID, containerConfig, podSandboxConfig)
// 3. 启动容器
err = m.runtimeService.StartContainer(ctx, containerID)
// 4. 执行 PostStart 钩子(失败则杀死容器)
if container.Lifecycle != nil && container.Lifecycle.PostStart != nil {
m.runner.Run(ctx, kubeContainerID, pod, container, container.Lifecycle.PostStart)
}
return "", nil
}
5.2 支持的运行时 @
| 运行时 | 说明 |
|---|---|
| containerd | 默认推荐,CNCF 项目 |
| CRI-O | 专为 Kubernetes 设计 |
| Docker | 通过 cri-dockerd 适配 |
六、Container Manager(CM) @
CM 负责节点上所有容器的资源管理和生命周期:
// pkg/kubelet/cm/container_manager.go
type ContainerManager interface {
Start(context.Context, *v1.Node, ActivePodsFunc, GetNodeFunc,
config.SourcesReady, status.PodStatusProvider,
internalapi.RuntimeService, bool) error
SystemCgroupsLimit() v1.ResourceList
GetNodeConfig() NodeConfig
GetResources(ctx context.Context, pod *v1.Pod, container *v1.Container) (*kubecontainer.RunContainerOptions, error)
UpdateQOSCgroups(logger klog.Logger) error
}
6.1 cgroup 资源隔离 @
// pkg/kubelet/cm/cgroup_manager_linux.go
// 公共部分抽为 cgroupCommon,cgroup v1/v2 各有具体实现
type cgroupCommon struct {
subsystems *CgroupSubsystems // 节点上已挂载的 cgroup 子系统
useSystemd bool // 是否使用 systemd cgroup 驱动
}
var _ CgroupManager = &cgroupV1impl{}
var _ CgroupManager = &cgroupV2impl{}
func (m *cgroupV1impl) Create(logger klog.Logger, config *CgroupConfig) error {
// 创建 cgroup 目录
// 设置资源限制:cpu.shares、memory.limit_in_bytes 等
// (v2 实现对应 cpu.weight、memory.max 等)
return nil
}
6.2 CPU 管理策略 @
| 策略 | 说明 |
|---|---|
| none | 默认,共享 CPU |
| static | 独占 CPU 核(Guaranteed QoS) |
// pkg/kubelet/cm/cpumanager/policy_static.go
func (p *staticPolicy) Allocate(logger klog.Logger, s state.State,
pod *v1.Pod, container *v1.Container, operation lifecycle.Operation) error {
// 为 Guaranteed Pod 的容器分配独占 CPU 核(目前仅支持 Add 操作)
// 另有 AllocatePod(logger, s, pod, operation) 处理 Pod 级资源分配
return p.allocateForAdd(logger, s, pod, container)
}
6.3 设备管理 @
// pkg/kubelet/cm/devicemanager/types.go
type Manager interface {
Start(logger klog.Logger, activePods ActivePodsFunc, sourcesReady config.SourcesReady, ...) error
Allocate(ctx context.Context, pod *v1.Pod, container *v1.Container, operation lifecycle.Operation) error
UpdatePluginResources(node *schedulerframework.NodeInfo, attrs *lifecycle.PodAdmitAttributes) error
GetDevices(podUID, containerName string) ResourceDeviceInstances
}
七、Pod 创建完整流程 @
1. 调度器将 Pod 调度到节点
↓
2. kubelet 监听到 Pod 分配(Watch API Server)
↓
3. Pod Sync Loop 检测到新的 Pod
↓
4. 创建 Pod 数据目录
/var/lib/kubelet/pods/{podUID}/
↓
5. 创建 Pod Sandbox(Pause 容器)
- 创建网络命名空间
- 配置网络(CNI)
↓
6. 拉取镜像
- 检查本地缓存
- 从 Registry 拉取
↓
7. 创建 Init 容器(如果有)
- 按顺序执行
- 必须全部成功
↓
8. 创建主容器
- 创建容器配置
- 设置资源限制(cgroup)
- 设置环境变量
- 挂载 Volume
↓
9. 启动容器
- 执行 PostStart 钩子
- 启动 Liveness/Readiness 探针
↓
10. 更新 Pod 状态
- 上报到 API Server
- 触发相关事件
八、健康检查(Probes) @
Kubelet 负责执行三种健康检查:
8.1 探针类型 @
| 探针 | 作用 | 失败后果 |
|---|---|---|
| LivenessProbe | 容器是否存活 | 重启容器 |
| ReadinessProbe | 容器是否就绪 | 从 Service 摘除 |
| StartupProbe | 容器是否启动完成 | 超过 failureThreshold 后杀死并重启容器;成功之前会禁用 Liveness/Readiness 探针 |
8.2 探针实现 @
// pkg/kubelet/prober/prober_manager.go(probeType 是小写未导出类型)
type probeType int
const (
liveness probeType = iota
readiness
startup
)
// pkg/kubelet/prober/prober.go
func (pb *prober) probe(ctx context.Context, probeType probeType, pod *v1.Pod,
status v1.PodStatus, container v1.Container, containerID kubecontainer.ContainerID,
) (results.Result, error) {
// v1.Probe 内嵌 ProbeHandler(Exec/HTTPGet/TCPSocket/GRPC 指针字段),
// 没有 Type 字段,通过判断哪个字段非 nil 来选择探测方式
switch probeSpec := getProbeSpec(container, probeType); {
case probeSpec.Exec != nil:
return pb.exec.Probe(...)
case probeSpec.HTTPGet != nil:
return pb.http.Probe(...)
case probeSpec.TCPSocket != nil:
return pb.tcp.Probe(...)
case probeSpec.GRPC != nil:
return pb.grpc.Probe(...)
}
}
8.3 探针执行循环 @
// pkg/kubelet/prober/worker.go
func (w *worker) run(ctx context.Context) {
// 用探针配置的 PeriodSeconds 构造 ticker(worker 没有 ProbeInterval 字段)
probeTicker := time.NewTicker(time.Duration(w.spec.PeriodSeconds) * time.Second)
for w.doProbe(ctx) {
select {
case <-w.stopCh:
return
case <-probeTicker.C:
// 到点继续下一次探测
}
}
}
// doProbe 执行一次探测,只把结果写入 resultsManager,从不直接杀死容器
func (w *worker) doProbe(ctx context.Context) (keepGoing bool) {
result, err := w.probeManager.prober.probe(
ctx, w.probeType, w.pod, status, w.container, w.containerID)
w.resultsManager.Set(w.containerID, result, w.pod)
// 结果由 syncLoopIteration 消费 kl.livenessManager.Updates() 等通道,
// 触发 handler.HandlePodSyncs,在后续 SyncPod 流程中间接杀死并重建失败的容器
return true
}
九、Volume 管理 @
9.1 Volume 挂载流程 @
// pkg/kubelet/volumemanager/volume_manager.go
func (vm *volumeManager) WaitForAttachAndMount(pod *v1.Pod) error {
for _, volume := range pod.Spec.Volumes {
// 1. 等待 Volume 由 Controller 完成 Attach(如 EBS)
err := volumePlugin.WaitForAttach(volume, pod)
// 2. 挂载 Volume 到节点
err = volumePlugin.MountDevice(volume, pod)
// 3. 挂载到 Pod 目录
err = volumePlugin.SetUp(volume, pod)
}
return nil
}
9.2 Volume 插件 @
| 类型 | 说明 |
|---|---|
| emptyDir | 临时目录,Pod 删除后消失 |
| hostPath | 主机目录映射 |
| configMap | ConfigMap 挂载为文件 |
| secret | Secret 挂载为文件 |
| persistentVolumeClaim | 持久化存储 |
| csi | CSI 驱动(外部存储) |
| nfs | NFS 网络存储 |
十、Pod 驱逐(Eviction) @
当节点资源不足时,kubelet 会主动驱逐部分 Pod:
10.1 驱逐触发条件 @
// pkg/kubelet/eviction/api/types.go
type Threshold struct {
Signal Signal // memory.available、nodefs.inodesFree 等
Operator ThresholdOperator // 目前只有 OpLessThan
Value ThresholdValue // 二选一:Quantity *resource.Quantity 或 Percentage float32
GracePeriod time.Duration
MinReclaim *ThresholdValue // 触发阈值后至少回收的资源量
}
type Signal string
const (
SignalMemoryAvailable Signal = "memory.available"
SignalNodeFsAvailable Signal = "nodefs.available"
SignalNodeFsInodesFree Signal = "nodefs.inodesFree"
SignalImageFsAvailable Signal = "imagefs.available"
SignalImageFsInodesFree Signal = "imagefs.inodesFree"
SignalPIDAvailable Signal = "pid.available"
)
10.2 驱逐流程 @
// pkg/kubelet/eviction/eviction_manager.go
func (m *managerImpl) Start(ctx context.Context, diskInfoProvider DiskInfoProvider,
podFunc ActivePodsFunc, podCleanedUpFunc PodCleanedUpFunc,
monitoringInterval time.Duration) {
go func() {
for {
// 1. 核心循环:采集节点摘要、判定阈值、必要时驱逐
evictedPods, err := m.synchronize(ctx, diskInfoProvider, podFunc)
if evictedPods != nil && err == nil {
// 2. 每个监控周期最多驱逐 1 个 Pod,
// 驱逐后阻塞等待其清理完成再进入下一周期
m.waitForPodsCleanup(logger, podCleanedUpFunc, evictedPods)
} else {
time.Sleep(monitoringInterval)
}
}
}()
}
// synchronize 内部:先尝试回收节点级资源(镜像/容器 GC),
// 压力仍未缓解时按信号对应的排序函数(rankMemoryPressure 等)
// 对活跃 Pod 排序,驱逐排名第一的 Pod 后立即返回
10.3 Pod 驱逐优先级 @
驱逐排序的首要依据是 Pod 资源用量是否超过其 requests,其次是 Pod Priority,最后是相对 requests 的超用量(见 pkg/kubelet/eviction/helpers.go 的 rankMemoryPressure 等排序函数)。QoS 等级本身不是排序键,但与被驱逐风险大致相关:
| QoS 等级 | 被驱逐风险 | 说明 |
|---|---|---|
| BestEffort | 最高(最易被驱逐) | 无 requests,用量必然"超过 requests” |
| Burstable | 中等 | 超出 requests 的部分越多越危险 |
| Guaranteed | 最低(最后被驱逐) | requests == limits,且通常 Priority 较高 |
十一、节点状态上报 @
Kubelet 定期向 API Server 上报节点状态:
// pkg/kubelet/kubelet_node_status.go
// 周期性循环在外层 syncNodeStatus 中,每次调用带重试的 updateNodeStatus
func (kl *Kubelet) updateNodeStatus(ctx context.Context) error {
for i := 0; i < nodeStatusUpdateRetry; i++ {
if err := kl.tryUpdateNodeStatus(ctx, i); err != nil {
// 失败则重试
} else {
return nil
}
}
return fmt.Errorf("update node status exceeds retry count")
}
// tryUpdateNodeStatus 内部经 setNodeStatus 填充节点条件:
// Ready / MemoryPressure / DiskPressure / PIDPressure / NetworkUnavailable
// 其中三种压力条件来自 eviction manager 的
// IsUnderMemoryPressure() / IsUnderDiskPressure() / IsUnderPIDPressure(),
// 最后调用 API Server 的 UpdateStatus 上报
11.1 节点条件 @
| 条件 | 含义 |
|---|---|
| Ready | 节点健康,可接收 Pod |
| MemoryPressure | 内存压力大 |
| DiskPressure | 磁盘压力大 |
| PIDPressure | PID 压力大 |
| NetworkUnavailable | 网络未就绪 |
十二、Pod 垃圾回收 @
Kubelet 负责清理已终止的 Pod。Pod 级清理由 syncLoopIteration 的 housekeeping 分支定期触发 HandlePodCleanups:
// pkg/kubelet/kubelet_pods.go
func (kl *Kubelet) HandlePodCleanups(ctx context.Context) error {
// 1. 对比期望与实际集合,停止已终止 Pod 的 worker
workingPods := kl.podWorkers.SyncKnownPods(logger, allPods)
// 2. 停止不再运行的 Pod 的探针
kl.probeManager.CleanupPods(possiblyRunningPods)
// 3. 清理孤立的 Pod 状态
kl.removeOrphanedPodStatuses(logger, allPods, mirrorPods)
// 4. 清理孤立的 Pod 数据目录和 Volume
kl.cleanupOrphanedPodDirs(logger, allPods, runningRuntimePods)
// 5. 删除孤立的 mirror Pod、终止残留的孤儿容器等
}
容器级 GC 由 containerGC 独立执行,按 GCPolicy 策略删除已退出的容器:
// pkg/kubelet/container/container_gc.go
type GCPolicy struct {
MinAge time.Duration // 容器可被回收的最小存在时间
MaxPerPodContainer int // 每个 Pod(UID+容器名)最多保留的已退出容器数
MaxContainers int // 全节点最多保留的已退出容器总数
}
十三、镜像管理 @
13.1 镜像拉取 @
// pkg/kubelet/images/image_manager.go
func (m *imageManager) EnsureImageExists(
ctx context.Context,
objRef *v1.ObjectReference,
pod *v1.Pod,
requestedImage string,
pullSecrets []v1.Secret,
podSandboxConfig *runtimeapi.PodSandboxConfig,
podRuntimeHandler string,
pullPolicy v1.PullPolicy,
) (imageRef, message string, err error) {
// 1. 通过 CRI ImageStatus 检查镜像是否已存在(不存在本地 imageCache)
// 2. 按 pullPolicy(Always/IfNotPresent/Never)决定是否需要拉取
// 3. 需要拉取时解析凭据并调用 CRI PullImage
return imageRef, message, nil
}
13.2 镜像垃圾回收 @
// pkg/kubelet/images/image_gc_manager.go
func (im *realImageGCManager) GarbageCollect(ctx context.Context, beganGC time.Time) error {
// 1. 按最近使用时间对镜像排序(最久未用的在前)
images, err := im.imagesInEvictionOrder(ctx, freeTime)
// 2. 获取镜像文件系统的磁盘使用率
fsStats, _, err := im.statsProvider.ImageFsStats(ctx)
// 3. 使用率高于 HighThresholdPercent 时,从旧到新删除未使用的镜像,
// 直到使用率降到 LowThresholdPercent
if usagePercent >= im.policy.HighThresholdPercent {
im.freeSpace(ctx, amountToFree, freeTime, images)
}
return nil
}
十四、设计亮点总结 @
- 异步处理:Pod Sync Loop 与容器启动异步执行,提高吞吐量。
- 分层管理:Pod 层管理 Pod 生命周期,容器层通过 CRI 与运行时交互,资源层通过 cgroup 实现资源隔离。
- 健康检查:三种探针(Liveness/Readiness/Startup)保证容器健康。
- 驱逐机制:基于资源压力信号的驱逐策略(按用量是否超过 requests 及 Pod Priority 排序),保证高优先级 Pod 的运行。
- 状态上报:定期上报节点状态,为调度器提供决策依据。
十五、总结 @
Kubelet 是 Kubernetes 节点上的"执行者":
- 监听:Watch API Server 获取 Pod 分配
- 创建:通过 CRI 创建容器
- 监控:执行健康检查,上报状态
- 清理:驱逐低优先级 Pod,回收资源
理解 Kubelet,就掌握了 Pod 从"调度结果"到"运行实体"的完整链路。