containerManager
# 1.qos
# 1.1.简介
quality of service意为服务质量,作为一种控制机制限制不同用户或数据流的优先级,以影响pod的调度和驱逐策略,保证节点资源紧张或不足时可以进行资源释放,维持节点运行的稳定性。--- Guaranteed pod中所有container的resource申请的limit和request相同且不为0 --- Burstable 至少一个container设置resource的limit和request --- BestEffort pod container未设置resource的limit和request1
2
3
4
5
6
7
8
不可压缩资源和节点稳定性
1.资源出现压力时,
kubelet gc会优先清理停止的pod和未使用镜像,其次才会根据pod优先级进行驱逐2.容器创建时,
kubelet会设置oom_score_adj,进而影响oom killer杀死进程
# 1.2.oomScore
qos会根据不同的分类进行打分,分值越大越先被oom,范围处于-1000~1000之间,kubelet和kube-proxy进程的打分为-999,也就是最后被oomKill,guaranteed类型设置为-997,besteffort类型会设置为1000,也就是最先被回收。const ( KubeletOOMScoreAdj int = -999 // kubelet score KubeProxyOOMScoreAdj int = -999 // kube-proxy score guaranteedOOMScoreAdj int = -997 besteffortOOMScoreAdj int = 1000 )1
2
3
4
5
6此外,核心容器的
score就是guaranteedOOMScoreAdj,qos为burstable的pod计算稍微复杂一些,边界为3~999。func GetContainerOOMScoreAdjust(pod *v1.Pod, container *v1.Container, memoryCapacity int64) int { // 核心pod // 1.pod.Spec.PriorityClassName=system-node-critical // 2.static pod || mirror pod // 3.pod.Spec.Priority >= 2000000000 if types.IsNodeCriticalPod(pod) { return guaranteedOOMScoreAdj } switch v1qos.GetPodQOS(pod) { // guaranteed pod case v1.PodQOSGuaranteed: return guaranteedOOMScoreAdj // besteffort pod case v1.PodQOSBestEffort: return besteffortOOMScoreAdj } // 可突增型容器,介于保证型和尽力而为型之间 memoryRequest := container.Resources.Requests.Memory().Value() // 基础分数计算 oomScoreAdjust := 1000 - (1000*memoryRequest)/memoryCapacity // 下边界不能小于3 if int(oomScoreAdjust) < (1000 + guaranteedOOMScoreAdj) { return (1000 + guaranteedOOMScoreAdj) } // 上边界不能大于999 if int(oomScoreAdjust) == besteffortOOMScoreAdj { return int(oomScoreAdjust - 1) } return int(oomScoreAdjust) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
# 1.3.生成score
kubelet启动容器前会生成linuxContainerConfig配置,根据不同级别设置容器的oomScore信息,资源不足时根据score进行驱逐及回收策略调整。func (m *kubeGenericRuntimeManager) startContainer(...) (string, error) { ... containerConfig, cleanupAction, err := m.generateContainerConfig(container, pod, restartCount, podIP, imageRef, podIPs, target) if cleanupAction != nil { defer cleanupAction() } ... } func (m *kubeGenericRuntimeManager) generateContainerConfig(...) (*runtimeapi.ContainerConfig, func(), error) { ... // set platform specific configurations. m.applyPlatformSpecificContainerConfig(config, container, pod, uid, username, nsTarget) ... } func (m *kubeGenericRuntimeManager) applyPlatformSpecificContainerConfig(...) error { enforceMemoryQoS := false // MemoryQoS enabled with cgroups v2 if Enabled(kubefeatures.MemoryQoS) && libcontainercgroups.IsCgroup2UnifiedMode() { enforceMemoryQoS = true } // generate oomScore cl, err := m.generateLinuxContainerConfig(container, pod, uid, username, nsTarget, enforceMemoryQoS) ... config.Linux = cl return nil } func (m *kubeGenericRuntimeManager) generateLinuxContainerConfig(...) *runtimeapi.LinuxContainerConfig { ... lc.Resources.OomScoreAdj := int64(qos.GetContainerOOMScoreAdjust(pod, container, int64(m.machineInfo.MemoryCapacity))) ... return lc }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
# 2.cgroup
# 2.1.简介
cgroup(control group)意为控制组,是linux内核提供的一种机制,支持根据需求把一系列系统任务及其子任务整合或分隔到按资源划分等级的不同组内,用以限制、记录和隔离进程组所使用的物理资源,覆盖cpu、memory、network和磁盘IO。目前,cgroup已经成为容器技术实现资源限制和隔离的重要基础。
# 2.2.核心组件
task、cgroup、hierarchy和subsystem是控制资源的核心组件,通过协同工作实现进程资源控制。hierarchy会将cgroup组织为树状结构,subsystem作为资源控制模块依附于hierarchy,task需加入cgroup接受cgroup继承的subsystem管理。--- task(任务) 代表系统的一个进程 --- cgroup(控制组) 受相同资源限制的一组进程 --- hierarchy(层级) 类似树形结构,每个节点代表一个控制组,根节点表示系统所有任务的默认控制组 --- subsystem(子系统): 资源控制器,用以管理特定类型资源,subsystem加入hierarchy控制层级所有cgroup资源使用1
2
3
4
5
6
7
8
9
10
11
# 2.3.子系统
subsystem是linux内核实现具体资源控制功能的核心模块,每个系统专注于一种或一类资源(cpu/IO/mem)的管理,通过cgroup的层级结构和任务协作,实现对进程组的精细化资源限制、监控和隔离。subsystem本质是资源控制器,设定资源上限、统计资源使用情况、隔离不同进程组的资源访问及对进程组执行冻结和恢复操作。1. cpu子系统: 限制进程的cpu使用率 2. cpuacct子系统: 统计cgroups中进程的cpu使用报告 3. cpuset子系统: 向cgroups中的进程分配单独的cpu节点或内存节点 4. memory子系统: 限制进程的memory使用量 5. blkio子系统: 限制进程的块设备io 6. devices子系统: 控制进程可访问的设备 7. net_cls子系统: 标记cgroups中进程的网络数据包,可以配合tc(traffic control)对数据包进行控制 8. net_prio子系统: 限制进程网络流量的优先级 9. huge_tlb子系统: 限制hugeTLB的使用 10.freezer子系统: 挂起或恢复cgroups中的进程 11.ns子系统: 隔离不同cgroups进程到不同命名空间1
2
3
4
5
6
7
8
9
10
11
阻止形式
1.
cgroup用于对进程分组2.
hierarchy根据继承关系,将多个cgroup组成一棵树3.
subsystem负责资源限制操作,将subsystem和hierarchy绑定后,hierarchy上的所有cgroup下的进程都会被subsystem限制
# 2.4.文件系统
virtual filesystem(VFS)是Linux内核的一层抽象,屏蔽了不同物理文件系统(ext4/tmpfs/proc)的实现差异,为用户态和内核态提供统一的文件操作接口(open/read/write/mkdir)。cgroup的用户态接口完全基于VFS实现,内核通过特殊的文件系统类型(cgroup文件系统)将层级、控制组和子系统等抽象为用户可操作的文件和目录,基于文件操作驱动完成资源管理配置。
注意
基于
systemd系统的操作系统中,sys/fs/cgroup目录都是由systemd在系统启动的过程中挂载的,挂载类型为只读。这意味着系统不再建议直接向/sys/fs/cgroup目录下创建新的目录及挂载其它子系统。
# 2.5.cgroup v1
cgroup v1在每种资源对应的层级都有一个根节点,代表该类资源的全部分配,子节点则继承父节点的基础上对资源进行限制和划分。用户进程需要对资源使用进行限制时,进程需要和指定资源下的某个cgroup关联,建立起对资源类型的限制。
注意
cgroup v1为了提供灵活性,允许进程属于多个hierarchy的不同group,以限制多类资源,这种多hierarchy结构导致内核实现较为混乱
# 2.6.cgroup v2
cgroup v2将层级进行简化,整个cgroup系统只有一个层级,每个节点都可以拥有对cpu/mem/IO等多种资源的限制,父节点开启的子系统控制器可以控制到子节点。
改进
1.
cgroup v2的controller都会被挂载到unified hierarchy,不再沿用多hierarchy挂载controller情况2.
process只能绑定到root cgroup和leaf cgroup3.
cgroup.controllers和cgroup.subtree_control限制可用的controller4.
v1版本的task文件和cpuset controller中的cgroup.clone_children文件移除5.
cgroup为空时的通知机制得到改进,通过cgroup.events文件通知6.
cgroups.max.depth和cgroup.max.descendants文件新增,用于查看及设置hierarchy下后代cgroups数量
# 3.cgroupManager
# 3.1.简介
cgroupManager是容器资源隔离的核心组件,用于管理节点上的cgroup层次结构,确保容器按照配额使用资源,负责根据不同驱动管理cgroup的更新,设置pod和容器对应的cgroup资源限制,监控节点cgroup状态,主动回收不再需要的cgroup释放资源。--- cgroupfs驱动 1.直接操作/sys/fs/cgroup目录下的文件,通过创建目录和读写配置文件管理cgroup 2.与运行时containerd/CRI-O的cgroup管理方式直接匹配 --- systemd驱动 1.借助systemd的slice机制管理cgroup,systemd将cgroup组织为.slice文件,每个slice对应一个cgroup 2.层次结构由systemd自动管理,与systemd的服务管理集成 3.路径格式为system.slice/<unit>.slice,如kubepods-burstable-pod123.slice1
2
3
4
5
6
7
8
# 3.2.接口
type CgroupManager interface { // 创建叶子cgroup(父cgroup需存在). Create(*CgroupConfig) error // 删除cgroup. Destroy(*CgroupConfig) error // 更新cgroup配置. Update(*CgroupConfig) error // 检查cgroup有效性(路径存在/配置合法). Validate(name CgroupName) error // 检查cgroup存在 Exists(name CgroupName) bool // 驱动程序转换后主机上的文本cgroupfs名称 // systemd: foo.slice/foo-bar.slice Name(name CgroupName) string // 实际路径转换为内部表示符 CgroupName(name string) CgroupName // cgroup内所有pid Pids(name CgroupName) []int // cpu cfs降到最小份额(用于资源回收) ReduceCPULimits(cgroupName CgroupName) error // 内存使用情况 MemoryUsage(name CgroupName) (int64, error) } // 依赖runc/libcontainer的无状态实现 type cgroupManagerImpl struct { // 节点已挂载的cgroup子系统信息 subsystems *CgroupSubsystems // 驱动是否为systemd useSystemd bool } type CgroupSubsystems struct { // Cgroup subsystem mounts. /sys/fs/cgroup/cpu -> [cpu, cpuacct] Mounts []libcontainercgroups.Mount // Cgroup subsystem to their mount location. cpu -> /sys/fs/cgroup/cpu MountPoints map[string]string }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39注意
cgroupManagerImpl封装子系统信息和驱动类型,提供对Linux cgroup的无状态、多驱动支持管理能力,高层接口调用会转换为底层的文件系统操作或systemd调用,向容器化平台提供统一的资源管理抽象。
# 3.3.初始化
NewCgroupManager()时创建cgroup管理器实例的工厂函数,根据传入的subsyetem和驱动类型,初始化一个实现CgroupManager接口的对象,用于后续的cgroup管理和资源配置更新。func NewCgroupManager(cs *CgroupSubsystems, cgroupDriver string) CgroupManager { return &cgroupManagerImpl{ subsystems: cs, useSystemd: cgroupDriver == "systemd", } } // 传入的subsystem func GetCgroupSubsystems() (*CgroupSubsystems, error) { if libcontainercgroups.IsCgroup2UnifiedMode() { return getCgroupSubsystemsV2() } return getCgroupSubsystemsV1() }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15getCgroupSubsystemsV2()用于获取cgroupv2支持子系统,底层或调用libcontainercgroups.GetAllSubsystems()获取子系统信息,获取时兼容统一模式(完整迁移)和混合模式(v1和v2)共存情况。func GetAllSubsystems() ([]string, error) { // 统一模式 if IsCgroup2UnifiedMode() { // cgroupv2中devices和freezer不会出现在cgroup.controllers,但实际存在且可用 pseudo := []string{"devices", "freezer"} // 其余子系统从cgroup.controllers读取 data, err := ReadFile("/sys/fs/cgroup", "cgroup.controllers") ... subsystems := append(pseudo, strings.Fields(data)...) return subsystems, nil } // 混合模式(/proc/cgroups包含完整子系统) f, err := os.Open("/proc/cgroups") ... defer f.Close() subsystems := []string{} s := bufio.NewScanner(f) for s.Scan() { text := s.Text() if text[0] != '#' { parts := strings.Fields(text) if len(parts) >= 4 && parts[3] != "0" { subsystems = append(subsystems, parts[0]) } } } ... return subsystems, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29getCgroupSubsystemsV1()用于获取cgroupv1支持子系统,底层调用libcontainercgroups.GetCgroupMounts()获取cgroup mount以解析子系统信息。func GetCgroupMounts(all bool) ([]Mount, error) { ... return getCgroupMountsV1(all) } func getCgroupMountsV1(all bool) ([]Mount, error) { // 读取/proc/self/mountinfo下所有cgroup挂载点 mi, err := readCgroupMountinfo() ... // 读取及解析/proc/self/cgroup allSubsystems, err := ParseCgroupFile("/proc/self/cgroup") ... // 初始化子系统映射 allMap := make(map[string]bool) for s := range allSubsystems { allMap[s] = false } // 合并 // []Mount{ // { // Mountpoint: "/sys/fs/cgroup/cpu", // Subsystems: []string{"cpu", "cpuacct"}, // }, // } return getCgroupMountsHelper(allMap, mi, all) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
# 4.qosContainerManager
# 4.1.简介
qos containerManager是实现资源质量保障(QOS)的核心组件,负责根据pod的资源请求和限制创建分层的cgroup结构,确保关键工作负载获得所需资源,同时限制低优先级pod的资源使用上限。type QOSContainerManager interface { // 启动后台监控和管理 Start(func() v1.ResourceList, ActivePodsFunc) error // 获取qos cgroup层次结构的信息(路径/资源限制) GetQOSContainersInfo() QOSContainersInfo // 根据资源使用情况和节点状态动态更新cgroup配置 UpdateCgroups() error } type qosContainerManagerImpl struct { ... // qos cgroup路径信息 qosContainersInfo QOSContainersInfo // cgroup子系统 subsystems *CgroupSubsystems // cgroupManager cgroupManager CgroupManager // 获取活跃pod的回调(podManager.Get) activePods ActivePodsFunc // 获取节点可分配资源回调(containerManager.Get) getNodeAllocatable func() v1.ResourceList // cgroup根路径 cgroupRoot CgroupName // qos资源保留配置 qosReserved map[v1.ResourceName]int64 } func NewQOSContainerManager(...) (QOSContainerManager, error) { // 未开启qos,返回空实现 if !nodeConfig.CgroupsPerQOS { return &qosContainerManagerNoop{ cgroupRoot: cgroupRoot, }, nil } return &qosContainerManagerImpl{ subsystems: subsystems, cgroupManager: cgroupManager, cgroupRoot: cgroupRoot, qosReserved: nodeConfig.QOSReserved, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
注意
qos containerManager本质上是对cgroupManager的复用,在此基础上加入qos cgroup和定时管理机制。
# 4.2.节点预留
createNodeAllocatableCgroups()负责创建节点可分配资源cgroup的核心逻辑,通过设置linux cgroup限制pod总体资源使用,确保节点保留足够资源用于系统组件运行,提升集群稳定性。// node allocatable cgroup维护 func (cm *containerManagerImpl) createNodeAllocatableCgroups() error { // 节点可分配资源总量 nodeAllocatable := cm.internalCapacity // Use Node Allocatable limits instead of capacity if the user requested enforcing node allocatable. nc := cm.NodeConfig.NodeAllocatableConfig // qos和节点可分配资源限制(--enforce-node-allocatable=pods)开启 if cm.CgroupsPerQOS && nc.EnforceNodeAllocatable.Has(kubetypes.NodeAllocatableEnforcementKey) { // 节点资源总容量调整为扣减系统预留的分配量 nodeAllocatable = cm.getNodeAllocatableInternalAbsolute() } cgroupConfig := &CgroupConfig{ Name: cm.cgroupRoot, // The default limits for cpu shares can be very low which can lead to CPU starvation for pods. ResourceParameters: getCgroupConfig(nodeAllocatable), } // qos根cgroup存在 if cm.cgroupManager.Exists(cgroupConfig.Name) { return nil } // qos根cgroup创建 cm.cgroupManager.Create(cgroupConfig) ... return nil } // 可分配资源扣减 func (cm *containerManagerImpl) getNodeAllocatableAbsoluteImpl(capacity v1.ResourceList) v1.ResourceList { ... // 遍历资源类型(cpu/mem/ephemeralStorage) for k, v := range capacity { // 复制原始值,避免修改原始数据 value := v.DeepCopy() // 扣除系统预留资源 if cm.NodeConfig.SystemReserved != nil { value.Sub(cm.NodeConfig.SystemReserved[k]) } // 扣除kube组件预留资源 if cm.NodeConfig.KubeReserved != nil { value.Sub(cm.NodeConfig.KubeReserved[k]) } // 负值修正 if value.Sign() < 0 { // Negative Allocatable resources don't make sense. value.Set(0) } // 保留计算结果 result[k] = value } return result }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
注意
节点预留由
kubelet启动时指定,本质上还是扣除预留资源后设置pod可用总资源上限
# 4.3.定时更新
containerManager启动时,会创建qos root cgroup,调用qosContainerManager.Start()开启定期更新机制,创建及配置不同qos级别的顶级cgroup容器,确保系统资源按优先级分配。func (m *qosContainerManagerImpl) Start(...) error { cm := m.cgroupManager // pod qos cgroup根目录(启动时指定的cgroup-root+kubepods) rootContainer := m.cgroupRoot // 检查根目录是否存在(/sys/fs/cgroup/{subsystem}/kubepods) if !cm.Exists(rootContainer) { return fmt.Errorf("root container %v doesn't exist", rootContainer) } // burstable和bestEffort创建子cgroup // guaranteed直接使用root cgroup qosClasses := map[v1.PodQOSClass]CgroupName{ v1.PodQOSBurstable: NewCgroupName(rootContainer, strings.ToLower(string(v1.PodQOSBurstable))), v1.PodQOSBestEffort: NewCgroupName(rootContainer, strings.ToLower(string(v1.PodQOSBestEffort))), } for qosClass, containerName := range qosClasses { resourceParameters := &ResourceConfig{} // bestEffort qos的cpu份额最低为2 if qosClass == v1.PodQOSBestEffort { minShares := uint64(MinShares) resourceParameters.CpuShares = &minShares } containerConfig := &CgroupConfig{ // /sys/fs/croup/{subsystem/kubepods/{qos} Name: containerName, ResourceParameters: resourceParameters, } // 设置hugePages无限制 m.setHugePagesUnbounded(containerConfig) // qos pod cgroup不存在则创建 if !cm.Exists(containerName) { ... cm.Create(containerConfig) // 存在则更新配置 } else { // ... cm.Update(containerConfig) } } // 记录qos cgroup根目录和qos pod cgroup目录 m.qosContainersInfo = QOSContainersInfo{ Guaranteed: rootContainer, Burstable: qosClasses[v1.PodQOSBurstable], BestEffort: qosClasses[v1.PodQOSBestEffort], } // 获取减去系统预留和组件预留后pod可分配资源(/cpu/mem/hugePage) m.getNodeAllocatable = getNodeAllocatable // 获取podManager维护的活跃pod m.activePods = activePods // 间隔1min刷新cgroup go wait.Until(func() { err := m.UpdateCgroups() ... }, periodicQOSCgroupUpdateInterval, wait.NeverStop) return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55m.UpdateCgroups()实现了qos cgroup配置的周期性更新机制,根据当前最新// 更新qos cgroup配置 func (m *qosContainerManagerImpl) UpdateCgroups() error { ... // 三种级别的qos cgroup qosConfigs := map[v1.PodQOSClass]*CgroupConfig{ v1.PodQOSGuaranteed: { Name: m.qosContainersInfo.Guaranteed, ResourceParameters: &ResourceConfig{}, }, v1.PodQOSBurstable: { Name: m.qosContainersInfo.Burstable, ResourceParameters: &ResourceConfig{}, }, v1.PodQOSBestEffort: { Name: m.qosContainersInfo.BestEffort, ResourceParameters: &ResourceConfig{}, }, } // 获取最新activePods,设置burstable和bestEffort两个qos级别的req cpu配置(guaranteed cpu相对固定,无需更新) // limit cpu配置由运行时设置 m.setCPUCgroupConfig(qosConfigs) ... // 设置三个qos等级的hugePage无限制,避免影响大页内存分配 m.setHugePagesConfig(qosConfigs) ... // 设置内存(内存qos新功能仅在cgroup v2可用) if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.MemoryQoS) && libcontainercgroups.IsCgroup2UnifiedMode() { // 设置guaranteed和burstable两个qos等级的req mem配置(guaranteedMin= request[guaranteed]+burstableMin) // guaranteed(/kubepods): sum of all guaranteed and burstable pods // guaranteed和burstable的保障性内存,besteffort不会挤压上述内存 // besteffortMem = memory.limit_in_bytes - memory.min m.setMemoryQoS(qosConfigs) } // 设置qos预留资源(req mem) if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.QOSReserved) { // 设置burstable和bestEffort两个qos级别的可用内存 // burstableLimit = allocatable - guaranteedMem // bestEffortLimit = bestEffortLimit - burstableMem // guaranteed实际使用可能低于limit,此处是为了提高内存利用率,假设预留内存设置为20%,总内存15G,两者分别申请8G,3G // 最可观的情况下,guaranteed内存设置为[1.6G,8G]---优先级最高 // burstable内存设置为[0.6G,13.4G]---优先级其次 // bestEffort内存设置为[0,12.8G]---优先级最低 // 内存不足时,按照bestEffort-->burstable--->guaranteed顺序优先级回收,确保guaranteed可用,尽量保证burstable可用 for resource, percentReserve := range m.qosReserved { switch resource { case v1.ResourceMemory: m.setMemoryReserve(qosConfigs, percentReserve) } } ... // 尝试更新cgroup配置(runc cgroupfs实现) for _, config := range qosConfigs { m.cgroupManager.Update(config) ... } // 使用量超出限制,调整预留限制为当前使用量 for resource, percentReserve := range m.qosReserved { switch resource { case v1.ResourceMemory: // 获取已使用内存(memory.usage_in_bytes/memory.current) // 设置burstable和bestEffort预留内存为usage m.retrySetMemoryReserve(qosConfigs, percentReserve) } } } // 更新qos cgroup配置 for _, config := range qosConfigs { m.cgroupManager.Update(config) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
注意
这里的
root cgroup内存和CPU设置只用到request,但并非说明limit没有意义,严格来说,limit上限会设置到container层级的限制,同时kubelet触发驱逐时,优先级相同的pod会根据request超出limit的情况执行killPod。
# 5.podContainerManager
# 5.1.简介
podContainerManager负责pod级别cgroup管理的核心组件,主要处理pod cgroup层级管理,根据qos级别向pod分配资源进行硬性限制,处理pod创建、更新、删除时的cgroup管理,确保资源及时释放。// 管理接口 type PodContainerManager interface { // 获取pod的cgroup标识及主机实际路径 GetPodContainerName(*v1.Pod) (CgroupName, string) // 检查pod的cgroup路径存在及创建 EnsureExists(*v1.Pod) error // 检查pod的cgroup路径存在 Exists(*v1.Pod) bool // 销毁pod的cgroup及其资源限制 Destroy(name CgroupName) error // 将pod的cpu cfs限额降低到最小,用于资源紧张时临时回收CPU资源 ReduceCPULimits(name CgroupName) error // 根据cgroupfs枚举所有pod uid及对应cgroup GetAllPodsFromCgroups() (map[types.UID]CgroupName, error) // 检查cgroup路径属于pod及提取pod uid IsPodCgroup(cgroupfs string) (bool, types.UID) } // 结构定义 type podContainerManagerImpl struct { // qos根目录信息 qosContainersInfo QOSContainersInfo // cgroup支持子系统 subsystems *CgroupSubsystems // cgroup管理器 cgroupManager CgroupManager // pod允许最大进程数 podPidsLimit int64 // cpu cfs配额强制限制 enforceCPULimits bool // 容器的cpu cfs周期 cpuCFSQuotaPeriod uint64 } // 初始化podContainerManager func (cm *containerManagerImpl) NewPodContainerManager() PodContainerManager { // qos开启 if cm.NodeConfig.CgroupsPerQOS { return &podContainerManagerImpl{ qosContainersInfo: cm.GetQOSContainersInfo(), subsystems: cm.subsystems, cgroupManager: cm.cgroupManager, podPidsLimit: cm.ExperimentalPodPidsLimit, enforceCPULimits: cm.EnforceCPULimits, cpuCFSQuotaPeriod: uint64(cm.CPUCFSQuotaPeriod / time.Microsecond), } } // 空实现 return &podContainerManagerNoop{ cgroupRoot: cm.cgroupRoot, } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52注意
podContainerManager底层依赖cgroupManager,针对Guaranteed、Burstable和BestEffort三个qos等级实现差异化管理,维护pod维度cgroup的生命周期。
# 5.2.sync podCgroup
syncPod()调度应用期间,podContainerManager会根据qos策略为pod创建并配置cgroup,确保容器运行在正确的资源隔离环境,以及处理重启或配置变更场景时的资源维护问题。func (kl *Kubelet) syncPod(...) (isTerminal bool, err error) { ... // 初始化podContainerManager pcm := kl.containerManager.NewPodContainerManager() if !kl.podWorkers.IsPodTerminationRequested(pod.UID) { ... // 非首次调度,pod cgroup不存在,kill pod podKilled := false if !pcm.Exists(pod) && !firstSync { ... // pod cgroup不存在视为异常,kill pod kl.killPod(pod, p, nil) ... } // podKill或不重启的不执行pod cgroup创建或更新 if !(podKilled && pod.Spec.RestartPolicy == v1.RestartPolicyNever) { // pod cgroup不存在 if !pcm.Exists(pod) { // 触发qos cgroup资源更新 kl.containerManager.UpdateQOSCgroups() // 创建及更新pod cgroup pcm.EnsureExists(pod) ... } } } } // pod cgroup检查 func (m *podContainerManagerImpl) Exists(pod *v1.Pod) bool { // 获取pod cgroup路径 podContainerName, _ := m.GetPodContainerName(pod) // 检查是否存在 return m.cgroupManager.Exists(podContainerName) } // 获取pod cgroup name func (m *podContainerManagerImpl) GetPodContainerName(pod *v1.Pod) (CgroupName, string) { // 获取pod qos等级 podQOS := v1qos.GetPodQOS(pod) // 获取qos根目录 var parentContainer CgroupName switch podQOS { case v1.PodQOSGuaranteed: parentContainer = m.qosContainersInfo.Guaranteed case v1.PodQOSBurstable: parentContainer = m.qosContainersInfo.Burstable case v1.PodQOSBestEffort: parentContainer = m.qosContainersInfo.BestEffort } // 生成pod cgroup名称 podContainer := GetPodCgroupNameSuffix(pod.UID) // pod cgroup路径(../kubepods/pod-uid) cgroupName := NewCgroupName(parentContainer, podContainer) // cgroupName转换(cgroupfs/systemd格式) cgroupfsName := m.cgroupManager.Name(cgroupName) return cgroupName, cgroupfsName } // pod cgroup检查及创建 func (m *podContainerManagerImpl) EnsureExists(pod *v1.Pod) error { podContainerName, _ := m.GetPodContainerName(pod) // 检查pod cgroup存在 alreadyExists := m.Exists(pod) // 不存在 if !alreadyExists { ... // 生成pod cgroup配置 containerConfig := &CgroupConfig{ Name: podContainerName, ResourceParameters: ResourceConfigForPod(pod, m.enforceCPULimits, m.cpuCFSQuotaPeriod, enforceMemoryQoS), } // 追加pod进程限制 if m.podPidsLimit > 0 { containerConfig.ResourceParameters.PidsLimit = &m.podPidsLimit } m.cgroupManager.Create(containerConfig) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
注意
调度
pod期间初始化podContainerManager,调用cgroupManager创建cgroups,资源配额就是容器配置的资源request
# 5.3.remove podCgroup
syncPod()调度时,执行killPod停止的pod不会立即删除,而是由家务线程(<-housekeeping)周期性回收,负责终止调度协程podWorker、清理pod及cgroup资源回收。func (kl *Kubelet) syncLoopIteration(...) bool { select { ... case <-housekeepingCh: // http/fs/apiserver资源未就绪 if !kl.sourcesReady.AllReady() { ... } else { ... // pod资源清理 handler.HandlePodCleanups() ... } } } func (kl *Kubelet) HandlePodCleanups() error { ... if kl.cgroupsPerQOS { // 获取pod cgroup pcm := kl.containerManager.NewPodContainerManager() cgroupPods, err = pcm.GetAllPodsFromCgroups() ... } ... // 清理不需要的pod ... // 移除不需要的cgroup pod if kl.cgroupsPerQOS { pcm := kl.containerManager.NewPodContainerManager() kl.cleanupOrphanedPodCgroups(pcm, cgroupPods, possiblyRunningPods) } ... return nil } // 清理pod cgroup func (kl *Kubelet) cleanupOrphanedPodCgroups(...) { // 遍历所有cgrouo pod for uid, val := range cgroupPods { // pod可能运行,跳过 if _, ok := possiblyRunningPods[uid]; ok { continue } // pod volume存在(asw维护或fs有pod volume目录),未声明为终止保留 if podVolumesExist := kl.podVolumesExist(uid); podVolumesExist && !kl.keepTerminatedPodVolumes { ... // cpu配额缩减至最小(2) pcm.ReduceCPULimits(val) continue } // 杀死cgroup附加pid,清理pod cgroup go pcm.Destroy(val) } } // 清理cgroup附加pid及pod cgroup func (m *podContainerManagerImpl) Destroy(podCgroup CgroupName) error { // 杀死cgroup附加的pid m.tryKillingCgroupProcesses(podCgroup) ... // 清理pod cgroup m.cgroupManager.Destroy(containerConfig) ... return nil } // 遍历及杀死cgroup附加pid func (m *podContainerManagerImpl) tryKillingCgroupProcesses(podCgroup CgroupName) error { // 获取cgroup附加pid pidsToKill := m.cgroupManager.Pids(podCgroup) ... for i := 0; i < 5; i++ { ... for _, pid := range pidsToKill { // pid已处理 if _, ok := removed[pid]; ok { continue } ... // 杀死pid m.killOnePid(pid) ... removed[pid] = true } if len(errlist) == 0 { return nil } } return utilerrors.NewAggregate(errlist) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
注意
清理
pod期间,pod删除后杀死pod cgroup下所有pid,再调用底层删除pod cgroup
# 6.cpuManager
# 6.1.numa
现代计算机的
CPU结构多采用NUMA(Non-Uniform Memory Access,非统一内存)结构,将cpu资源分开,以node为单位进行分组,每个node拥有独立的cpu/memory资源。NUMA架构中,内存直接attach到CPU上,CPU访问直接attach内存对应的物理地址具有较短的访问时间,访问其它CPU附着的内存需要通过inter-connect通道,响应时间就会降低。因此,对于CPU密集型的程序,应该减少工作负载在节流时引起的上下文切换,避免频繁的调度进程到不同的核心导致的缓存失效开销。
内存策略
1.
localalloc: 默认策略,优先从当前节点分配内存,本节点内存不足再从其它节点分配2.
preferred: 倾向从特定节点分配内存,指定节点的内存不足再从其它节点分配3.
membind: 从传入的几个节点上分配内存,指定节点内存不足拒绝分配4.
interleave: 从传入的几个节点依次轮询,指定节点内存不足再从其它节点分配
# 6.2.cpuset
cpuset是Linux内核cgroup子系统的一部分,用于精确控制进程或线程的cpu亲和性和内存结点绑定,优化调度性能及减少资源竞争。cpuset整体上为层级树结构,有一个根节点包含系统所有的cpu和内存节点资源。根节点可以分支出一个或多个子节点,将根的资源划分为多个子集。
注意
cpuset是Linux系统实现高性能资源隔离和numa优化的强大工具,特别适合于cpu亲和性敏感的关键应用(实时系统/数据库)。相对的,其配置复杂性和静态特性更适合在硬件拓扑明确、工作负载稳定的环境。同时,精细的资源绑定可能导致独占的资源利用率降低,增大调度器的决策复杂度。
# 6.3.Manager
cpuManager提供了可选的cpu管理策略支持cpuset的资源控制能力,实现特定container绑定指定的cpu,提升cpu敏感型任务的性能。实现上,Guaanteed类型的pod container申请的cpu是超过1的整数时,cpuManager会将pid分配到指定的cpus,减少工作负载在节流时引起的上下文切换。type Manager interface { // kubelet初始化时调用 Start(activePods ActivePodsFunc, sourcesReady config.SourcesReady, podStatusProvider status.PodStatusProvider, containerRuntime runtimeService, initialContainers containermap.ContainerMap) error // 容器创建前CPU分配决策 Allocate(pod *v1.Pod, container *v1.Container) error // 添加容器到CPU管理器映射 AddContainer(p *v1.Pod, c *v1.Container, containerID string) // 容器删除后调用释放分配的CPU资源 RemoveContainer(containerID string) error // 返回CPU管理器状态 State() state.Reader // 获取容器的CPU分配拓扑提示,用于协作Topology Manager GetTopologyHints(*v1.Pod, *v1.Container) map[string][]topologymanager.TopologyHint // 获取容器独占分配的CPU集合 GetExclusiveCPUs(podUID, containerName string) cpuset.CPUSet // 获取pod级别CPU分配拓扑提示 GetPodTopologyHints(pod *v1.Pod) map[string][]topologymanager.TopologyHint // 获取可用于分配的CPU集合 GetAllocatableCPUs() cpuset.CPUSet // 获取容器完整的CPU亲和性集合(包括共享和独占) GetCPUAffinity(podUID, containerName string) cpuset.CPUSet } type manager struct { // 当前使用的CPU分配策略 policy Policy // 状态检查周期,定期检查修正分配的状态 reconcilePeriod time.Duration // 当前缓存的CPU分配状态 state state.State ... // 容器运行时,用于更新容器资源配置 containerRuntime runtimeService // podManager获取活跃pod的回调 activePods ActivePodsFunc // statusManager获取pod状态回调 podStatusProvider status.PodStatusProvider // 容器ID映射表 containerMap containermap.ContainerMap // CPU拓扑信息 topology *topology.CPUTopology // 节点预留资源 nodeAllocatableReservation v1.ResourceList ... // 可分配的CPU集合(去除预留) allocatableCPUs cpuset.CPUSet // 正在进行准入检查的pod pendingAdmissionPod *v1.Pod }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68补充
cpuManager通过策略抽象、状态持久化和拓扑感知,实现容器级的CPU资源精细化控制。
# 6.4.NewManager
NewManager()是初始化cpu管理器的工厂函数,负责根据指定的策略类型(static/none)创建及初始化cpu管理实例,解析系统的cpu拓扑结构提供硬件信息支持,检查系统预留cpu资源确保策略正常工作。func NewContainerManager(...) (ContainerManager, error) { ... // Initialize CPU manager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.CPUManager) { cm.cpuManager, err = cpumanager.NewManager( nodeConfig.ExperimentalCPUManagerPolicy, nodeConfig.ExperimentalCPUManagerPolicyOptions, nodeConfig.ExperimentalCPUManagerReconcilePeriod, machineInfo, nodeConfig.NodeAllocatableConfig.ReservedSystemCPUs, cm.GetNodeAllocatableReservation(), nodeConfig.KubeletRootDir, cm.topologyManager, ) ... // 注册到拓扑管理器 cm.topologyManager.AddHintProvider(cm.cpuManager) } ... } // NewManager creates new cpu manager based on provided policy func NewManager(...) (Manager, error) { ... switch policyName(cpuPolicyName) { // 未指定策略 case PolicyNone: // 空策略(CFS按照配额公平调度) policy, err = NewNonePolicy(cpuPolicyOptions) ... // 静态策略 case PolicyStatic: // 根据cadvisor提供的machineInfo初始化cpuTopology topo, err = topology.Discover(machineInfo) ... // 获取系统预留CPU(必须预留系统资源确保核心组件可用) reservedCPUs, ok := nodeAllocatableReservation[v1.ResourceCPU] ... // 计算系统预留核心数 reservedCPUsFloat := float64(reservedCPUs.MilliValue()) / 1000 numReservedCPUs := int(math.Ceil(reservedCPUsFloat)) // 初始化静态策略 policy, err = NewStaticPolicy(topo, numReservedCPUs, specificCPUs, affinity, cpuPolicyOptions) ... default: return nil, fmt.Errorf("unknown policy: \"%s\"", cpuPolicyName) } manager := &manager{ policy: policy, reconcilePeriod: reconcilePeriod, lastUpdateState: state.NewMemoryState(), topology: topo, nodeAllocatableReservation: nodeAllocatableReservation, stateFileDirectory: stateFileDirectory, } manager.sourcesReady = &sourcesReadyStub{} return manager, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
# 6.5.NewStaticPolicy
NewStaticPolicy()是静态策略管理的核心函数,用于为符合条件的guaranteed pod分配专用的cpu核心,控制容器生命周期内cpu分配不变,提供稳定的cpu资源和隔离。// 初始化静态策略 func NewStaticPolicy(topology *topology.CPUTopology, numReservedCPUs int, reservedCPUs cpuset.CPUSet, affinity topologymanager.Store, cpuPolicyOptions map[string]string) (Policy, error) { ... policy := &staticPolicy{ topology: topology, affinity: affinity, cpusToReuse: make(map[string]cpuset.CPUSet), options: opts, } // cpu总数 allCPUs := topology.CPUDetails.CPUs() var reserved cpuset.CPUSet // 指定专用核心 if reservedCPUs.Size() > 0 { reserved = reservedCPUs // 选择低编号物理核心的cpu预留 } else { reserved, _ = policy.takeByTopology(allCPUs, numReservedCPUs) } ... policy.reserved = reserved return policy, nil } // cpu核心分配决策函数 func (p *staticPolicy) takeByTopology(availableCPUs cpuset.CPUSet, numCPUs int) (cpuset.CPUSet, error) { // 跨NUMA均衡分配 if p.options.DistributeCPUsAcrossNUMA { // 计算CPU分配的最小单元 cpuGroupSize := 1 if p.options.FullPhysicalCPUsOnly { cpuGroupSize = p.topology.CPUsPerCore() } return takeByTopologyNUMADistributed(p.topology, availableCPUs, numCPUs, cpuGroupSize) } // 紧凑分配 return takeByTopologyNUMAPacked(p.topology, availableCPUs, numCPUs) } // 均衡型CPU分配 func takeByTopologyNUMADistributed(topo *topology.CPUTopology, availableCPUs cpuset.CPUSet, numCPUs int, cpuGroupSize int) (cpuset.CPUSet, error) { // 无法按照cgroupSize整分,走紧凑型分配 if (numCPUs % cpuGroupSize) != 0 { return takeByTopologyNUMAPacked(topo, availableCPUs, numCPUs) } // 初始化CPU分配器(NUMA维度优先/socket维度优先) acc := newCPUAccumulator(topo, availableCPUs, numCPUs) ... // 获取可用的排序后的NUMA节点(NUMA/Socket-->NUMA) numas := acc.sortAvailableNUMANodes() // 计算满足CPU分配的最小和最大NUMA节点 // 1.NUMA节点数量 // 2.可用NUMA节点数量 // 3.逻辑核心总数 // 4.物理核心数 // 5.单NUMA节点物理核心 // 6.需要预留的物理核心数 // 7.计算预留需要的最小和最大NUMA节点数 minNUMAs, maxNUMAs := acc.rangeNUMANodesNeededToSatisfy(cpuGroupSize) // 枚举组合,预计算的NUMA节点范围内寻找最佳的CPU分配组合 for k := minNUMAs; k <= maxNUMAs; k++ { ... // n选k组合 acc.iterateCombinations(numas, k, func(combo []int) LoopControl { ... // 组合的numa节点提供的cpu总数检查 cpus := acc.details.CPUsInNUMANodes(combo...) if cpus.Size() < numCPUs { return Continue } // 组合的numa节点提供的cpu物理核,用于检查是否可均分 numCPUGroups := 0 for _, numa := range combo { numCPUGroups += (acc.details.CPUsInNUMANodes(numa).Size() / cpuGroupSize) } if (numCPUGroups * cpuGroupSize) < numCPUs { return Continue } // 计算单numa节点基础分配量 distribution := (numCPUs / len(combo) / cpuGroupSize) * cpuGroupSize for _, numa := range combo { // 当前numa节点的cpu逻辑核心 cpus := acc.details.CPUsInNUMANodes(numa) // 当前numa节点提供的cpu不够基础分配量 if cpus.Size() < distribution { return Continue } } // 均匀分配,记录numa节点分配基础量后可用的cpu availableAfterAllocation := make(mapIntInt, len(numas)) for _, numa := range numas { availableAfterAllocation[numa] = acc.details.CPUsInNUMANodes(numa).Size() } for _, numa := range combo { availableAfterAllocation[numa] -= distribution } // 计算还需要申请多少cpu remainder := numCPUs - (distribution * len(combo)) // 剩余可分配cpu筛选 var remainderCombo []int for _, numa := range combo { if availableAfterAllocation[numa] >= cpuGroupSize { remainderCombo = append(remainderCombo, numa) } } // 均分后已满足cpu申请,计算标准差度量分布均衡性 if remainder == 0 { bestLocalBalance = standardDeviation(availableAfterAllocation.Values()) bestLocalRemainder = nil } // cpu申请还有剩余,继续分配 for k := len(remainderCombo); remainder > 0 && k >= 1; k-- { // 剩余可用numa节点均分 acc.iterateCombinations(remainderCombo, k, func(subset []int) LoopControl { ... // 当前numa组合提供的总cpu逻辑核无法满足剩余申请 if sum(availableAfterAllocation.Values(subset...)) < remainder { return Continue } // 当前组合中均分cpu for remainder > 0 { for _, numa := range subset { // cpu申请已覆盖 if remainder == 0 { break } // 当前numa节点无法按组分 if availableAfterAllocation[numa] < cpuGroupSize { continue } // 当前numa节点可以分,更新剩余申请量及numa节点可用cpu数量 availableAfterAllocation[numa] -= cpuGroupSize remainder -= cpuGroupSize } } // 计算分配后的标准差衡量均衡性 balance := standardDeviation(availableAfterAllocation.Values()) if balance < bestLocalBalance { bestLocalBalance = balance bestLocalRemainder = subset } return Continue }) } // 记录分配组合 if bestLocalBalance < bestBalance { // 标准差 bestBalance = bestLocalBalance // 参与剩余分配的numa bestRemainder = bestLocalRemainder // 参与均分的numa bestCombo = combo } return Continue }) // k节点划分无法满足分配,尝试更多节点 if bestCombo == nil { continue } // numa节点均分的基础量(整组) distribution := (numCPUs / len(bestCombo) / cpuGroupSize) * cpuGroupSize for _, numa := range bestCombo { // 向numa节点实际申请 cpus, _ := takeByTopologyNUMAPacked(acc.topo, acc.details.CPUsInNUMANodes(numa), distribution) // 分配结果记录到分配器 acc.take(cpus) } // 剩余申请cpu分配 remainder := numCPUs - (distribution * len(bestCombo)) for remainder > 0 { // 尝试参与剩余分配的numa节点申请 for _, numa := range bestRemainder { // cpu申请已覆盖 if remainder == 0 { break } // 当前numa节点无法按组分 if acc.details.CPUsInNUMANodes(numa).Size() < cpuGroupSize { continue } // 向当前numa节点实际申请 cpus, _ := takeByTopologyNUMAPacked(acc.topo, acc.details.CPUsInNUMANodes(numa), cpuGroupSize) // 分配结果记录到分配器 acc.take(cpus) // 更新剩余申请量 remainder -= cpuGroupSize } } ... // Otherwise, return the result return acc.result, nil } // 无法按照组合均分,走紧凑型分配 return takeByTopologyNUMAPacked(topo, availableCPUs, numCPUs) } // 紧凑型分配 func takeByTopologyNUMAPacked(topo *topology.CPUTopology, availableCPUs cpuset.CPUSet, numCPUs int) (cpuset.CPUSet, error) { // 初始化cpu分配器 acc := newCPUAccumulator(topo, availableCPUs, numCPUs) ... // 策略1: 可用cpu关联的numa节点空闲,取低序numa/socket全部分给当前申请(need>cpus) acc.numaOrSocketsFirst.takeFullFirstLevel() ... // 策略2: 第一策略分配后无法满足,取可用cpu关联的空闲socket/numa全部分配给当前剩余申请(need>cpus) acc.numaOrSocketsFirst.takeFullSecondLevel() ... // 策略3: 第一和第二策略分配后无法满足,取可用cpu关联的socket/numa节点的可用core全部分给当前剩余申请(need>cpus) // numaFirst: 可用cpu按照socket排序,从socket分配core(避免跨socket) // socketFirst: 可用cpu按照numa排序,从numa分配core(避免跨numa) acc.takeFullCores() ... // 策略4: 剩余申请分配到core thread acc.takeRemainingCPUs() ... return cpuset.NewCPUSet(), fmt.Errorf("failed to allocate cpus") }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
# 6.6.Start
cpuManager.Start()负责初始化状态管理器,启动策略引擎及周期同步,维护pod可用cpu资源,确保节点上运行的pod的cpu分配与预期申请一致。func (m *manager) Start(...) error { ... // sourceReady回调 m.sourcesReady = sourcesReady // podManager维护的活跃pod m.activePods = activePods // statusManager的状态回调 m.podStatusProvider = podStatusProvider // 容器运行时 m.containerRuntime = containerRuntime // pod--container内存映射 m.containerMap = initialContainers // 初始化cpu分配状态(用于重启恢复) stateImpl, err := state.NewCheckpointState(m.stateFileDirectory, cpuManagerStateFileName, m.policy.Name(), m.containerMap) ... m.state = stateImpl // 启动策略分配器 err = m.policy.Start(m.state) ... // 获取策略可用cpu(static策略去除系统预留) m.allocatableCPUs = m.policy.GetAllocatableCPUs(m.state) // None策略不执行周期同步 if m.policy.Name() == string(PolicyNone) { return nil } // 间隔10s同步状态 go wait.Until(func() { m.reconcileState() }, m.reconcilePeriod, wait.NeverStop) return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29reconcileState()是cpuManager的核心逻辑,用于周期性同步容器的cpuset,确保内存中记录的cpu分配信息被应用到容器运行环境。// 状态周期同步 func (m *manager) reconcileState() (success []reconciledContainer, failure []reconciledContainer) { ... // 清理不再活跃的pod状态,防止内存泄漏 m.removeStaleState() // 遍历活跃pod for _, pod := range m.activePods() { // 获取pod运行时状态 pstatus, ok := m.podStatusProvider.GetPodStatus(pod.UID) ... // 获取pod所有容器(init+业务) allContainers := pod.Spec.InitContainers allContainers = append(allContainers, pod.Spec.Containers...) // 遍历所有容器 for _, container := range allContainers { // 查找containerID containerID, err := findContainerIDByName(&pstatus, container.Name) ... // 查找container status cstatus, err := findContainerStatusByName(&pstatus, container.Name) ... m.Lock() // container终止(此处不会释放cpu资源,避免重启的container丢失cpu亲和) if cstatus.State.Terminated != nil { _, _, err := m.containerMap.GetContainerRef(containerID) ... m.Unlock() continue } // 记录container映射 m.containerMap.Add(string(pod.UID), container.Name, containerID) m.Unlock() // 获取容器应该分配的cpu(container cpuset/default cpuset) cset := m.state.GetCPUSetOrDefault(string(pod.UID), container.Name) ... // 获取上次设置的cpu状态 lcset := m.lastUpdateState.GetCPUSetOrDefault(string(pod.UID), container.Name) // 对比更新 if !cset.Equals(lcset) { // 调用CRI更新container cpu(无法获取container cpuset时回退为default cpuset) err = m.updateContainerCPUSet(containerID, cset) ... // 更新本地成功更新的container cpu m.lastUpdateState.SetCPUSet(string(pod.UID), container.Name, cset) } ... } } return success, failure } // 状态清理 func (m *manager) removeStaleState() { // 数据源就绪检查 if !m.sourcesReady.AllReady() { return } // 加锁防止清理期间新的container分配cpu m.Lock() defer m.Unlock() // 向podManager获取活跃pod activeAndAdmittedPods := m.activePods() // 正在申请cpu的pod if m.pendingAdmissionPod != nil { activeAndAdmittedPods = append(activeAndAdmittedPods, m.pendingAdmissionPod) } // 整理活跃容器 activeContainers := make(map[string]map[string]struct{}) for _, pod := range activeAndAdmittedPods { activeContainers[string(pod.UID)] = make(map[string]struct{}) for _, container := range append(pod.Spec.InitContainers, pod.Spec.Containers...) { activeContainers[string(pod.UID)][container.Name] = struct{}{} } } // 获取cpuManager state的cpu分配信息 assignments := m.state.GetCPUAssignments() // 对比清理非活跃pod的分配记录 for podUID := range assignments { for containerName := range assignments[podUID] { if _, ok := activeContainers[podUID][containerName]; !ok { ... m.policyRemoveContainerByRef(podUID, containerName) ... } } } // 二次确认,清理非活跃pod的映射缓存 m.containerMap.Visit(func(podUID, containerName, containerID string) { if _, ok := activeContainers[podUID][containerName]; !ok { ... m.policyRemoveContainerByRef(podUID, containerName) ... } }) } // 清理分配信息 func (m *manager) policyRemoveContainerByRef(podUID string, containerName string) error { // 释放state中container分配的cpu资源,更新checkpoint err := m.policy.RemoveContainer(m.state, podUID, containerName) if err == nil { // 清理lastUpdateState中container分配的cpu资源 m.lastUpdateState.Delete(podUID, containerName) // 清理containerMap内存映射 m.containerMap.RemoveByContainerRef(podUID, containerName) } return err } // 释放container的cpu资源 func (p *staticPolicy) RemoveContainer(s state.State, podUID string, containerName string) error { // 当前container外已分配cpu cpusInUse := getAssignedCPUsOfSiblings(s, podUID, containerName) // 当前container分配到cpu if toRelease, ok := s.GetCPUSet(podUID, containerName); ok { // 清理state中container分配的cpu,更新checkpoint s.Delete(podUID, containerName) // 过滤释放资源 toRelease = toRelease.Difference(cpusInUse) // 释放cpu合并到default,更新checkpoint s.SetDefaultCPUSet(s.GetDefaultCPUSet().Union(toRelease)) } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
# 6.7.Allocate
cpuManager.Allocate()负责向Guaranteed qos级别的容器分配专属的独占cpu,分配时融合cpu拓扑、numa节点亲和和资源复用机制。// 准入阶段 func (m *resourceAllocator) Admit(attrs *lifecycle.PodAdmitAttributes) lifecycle.PodAdmitResult { ... for _, container := range append(pod.Spec.InitContainers, pod.Spec.Containers...) { ... if m.cpuManager != nil { m.cpuManager.Allocate(pod, &container) ... } ... } return admission.GetPodAdmitResult(nil) } func (m *manager) Allocate(p *v1.Pod, c *v1.Container) error { // 准入检查阶段的pod标记为待处理 m.setPodPendingAdmission(p) // 释放过期资源 m.removeStaleState() m.Lock() defer m.Unlock() // 执行策略分配 m.policy.Allocate(m.state, p, c) ... return nil } // 策略分配 func (p *staticPolicy) Allocate(s state.State, pod *v1.Pod, container *v1.Container) error { // 申请整数资源的guaranteed qos container if numCPUs := p.guaranteedCPUs(pod, container); numCPUs != 0 { // 物理核心对齐检查 if p.options.FullPhysicalCPUsOnly && ((numCPUs % p.topology.CPUsPerCore()) != 0) { return SMTAlignmentError{ RequestedCPUs: numCPUs, CpusPerCore: p.topology.CPUsPerCore(), } } // 尝试获取缓存的分配信息 if cpuset, ok := s.GetCPUSet(string(pod.UID), container.Name); ok { // 资源复用 p.updateCPUsToReuse(pod, container, cpuset) return nil } // 调用topologyManager获取资源亲和(mem/gpu) hint := p.affinity.GetAffinity(string(pod.UID), container.Name) // 申请cpu cpuset, err := p.allocateCPUs(s, numCPUs, hint.NUMANodeAffinity, p.cpusToReuse[string(pod.UID)]) ... // 持久化分配结果 s.SetCPUSet(string(pod.UID), container.Name, cpuset) // 更新复用池信息 p.updateCPUsToReuse(pod, container, cpuset) } // 共享池container不参与cpu分配 return nil } // 维护cpu复用池 func (p *staticPolicy) updateCPUsToReuse(pod *v1.Pod, container *v1.Container, cset cpuset.CPUSet) { // 清理其它pod复用信息 for podUID := range p.cpusToReuse { if podUID != string(pod.UID) { delete(p.cpusToReuse, podUID) } } // 初始化当前pod复用池 if _, ok := p.cpusToReuse[string(pod.UID)]; !ok { p.cpusToReuse[string(pod.UID)] = cpuset.NewCPUSet() } // 记录init容器待释放的cpu for _, initContainer := range pod.Spec.InitContainers { if container.Name == initContainer.Name { p.cpusToReuse[string(pod.UID)] = p.cpusToReuse[string(pod.UID)].Union(cset) return } } // 更新其它container已申请的cpu p.cpusToReuse[string(pod.UID)] = p.cpusToReuse[string(pod.UID)].Difference(cset) } // cpu申请 func (p *staticPolicy) allocateCPUs(s state.State, numCPUs int, numaAffinity bitmask.BitMask, reusableCPUs cpuset.CPUSet) (cpuset.CPUSet, error) { // 获取可用cpu(去除系统预留可用的+可复用的) allocatableCPUs := p.GetAvailableCPUs(s).Union(reusableCPUs) // 优先从numa亲和节点分配 result := cpuset.NewCPUSet() if numaAffinity != nil { alignedCPUs := cpuset.NewCPUSet() // 初始化位于numa亲和节点的可用cpu for _, numaNodeID := range numaAffinity.GetBits() { alignedCPUs = alignedCPUs.Union(allocatableCPUs.Intersection(p.topology.CPUDetails.CPUsInNUMANodes(numaNodeID))) } // 计算申请的cpu数量 numAlignedToAlloc := alignedCPUs.Size() if numCPUs < numAlignedToAlloc { numAlignedToAlloc = numCPUs } // 优选申请 alignedCPUs, err := p.takeByTopology(alignedCPUs, numAlignedToAlloc) ... result = result.Union(alignedCPUs) } // 剩余cpu申请 remainingCPUs, err := p.takeByTopology(allocatableCPUs.Difference(result), numCPUs-result.Size()) ... result = result.Union(remainingCPUs) // 更新共享池cpu,持久化checkpoint s.SetDefaultCPUSet(s.GetDefaultCPUSet().Difference(result)) return result, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
注意
1.
Allocate()由kubelet调度pod准入时触发2.
cpu申请时会先释放过期资源,复用init容器释放的cpu3.
cpu申请时会根据numa节点亲和优先
# 7.memManager
# 7.1.简介
memoryManager是内存管理组件,旨在向guaranteed qos pod提供保证内存及大页内存分配能力。类似于cpuManager,内存管理器支持跨NUMA和紧凑NUMA分配策略,基于topologyManager获取拓扑提示实现NUMA节点亲和,保证pod内容器的内存和大页内存与相同的NUMA节点关联。// 接口 type Manager interface { Start(...) error // 内存分配 AddContainer(p *v1.Pod, c *v1.Container, containerID string) // 内存申请 Allocate(pod *v1.Pod, container *v1.Container) error // 内存释放 RemoveContainer(containerID string) error // 状态查询 State() state.Reader // 拓扑亲和提示 GetTopologyHints(*v1.Pod, *v1.Container) map[string][]topologymanager.TopologyHint GetPodTopologyHints(*v1.Pod) map[string][]topologymanager.TopologyHint // 状态监控 GetMemoryNUMANodes(pod *v1.Pod, container *v1.Container) sets.Int GetAllocatableMemory() []state.Block GetMemory(podUID, containerName string) []state.Block } type manager struct { ... // 策略(none/static) policy Policy // 状态缓存 state state.State // 运行时 containerRuntime runtimeService // podManager activePods ActivePodsFunc // statusManager podStatusProvider status.PodStatusProvider // 容器和pod内存映射 containerMap containermap.ContainerMap // 数据源 sourcesReady config.SourcesReady // checkpoint目录 stateFileDirectory string // NUMA节点可分配内存 allocatableMemory []state.Block // 处于准入阶段的pod pendingAdmissionPod *v1.Pod }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
注意
1.
memoryManager维护内存映射跟踪已分配和可用内存情况2.
memoryManager内存分配遵循粒度由大到小原则,优先用粗粒度完整分配3.
memoryManager内存分配更适合数据库及其它内存性能要求较高的的应用4.
memoryManager分割numa节点内存势必带来内存碎片,暂时没有合适的机制平衡pod和碎片整理5.
memoryManager会维护缓存的内存分配状态,同时定时同步、持久化分配状态
# 7.2.NewManager
NewManager()是初始化mem管理器的工厂函数,负责根据指定的策略类型(static/none)创建及初始化内存管理实例,解析系统的内存拓扑结构提供硬件信息支持,检查系统预留内存资源确保策略正常工作。func NewContainerManager(...) (ContainerManager, error) { ... // 开启memoryManager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.MemoryManager) { // 初始化 cm.memoryManager, err = memorymanager.NewManager( nodeConfig.ExperimentalMemoryManagerPolicy, machineInfo, cm.GetNodeAllocatableReservation(), nodeConfig.ExperimentalMemoryManagerReservedMemory, nodeConfig.KubeletRootDir, cm.topologyManager, ) ... // 注册到拓扑管理器 cm.topologyManager.AddHintProvider(cm.memoryManager) } ... } // NewManager returns new instance of the memory manager func NewManager(...) (Manager, error) { ... // 初始化分配策略 switch policyType(policyName) { case policyTypeNone: // 空策略 policy = NewPolicyNone() case policyTypeStatic: // 计算系统预留内存 systemReserved, err := getSystemReservedMemory(machineInfo, nodeAllocatableReservation, reservedMemory) ... // 静态策略 policy, err = NewPolicyStatic(machineInfo, systemReserved, affinity) ... default: return nil, fmt.Errorf("unknown policy: \"%s\"", policyName) } manager := &manager{ policy: policy, stateFileDirectory: stateFileDirectory, } manager.sourcesReady = &sourcesReadyStub{} return manager, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
# 7.3.Start
memoryManager.Start()负责初始化内存管理策略,恢复持久化状态,检查分配策略的预留内存合法性及初始化内存分配状态,用于后续内存分配检查。// Start starts the memory manager under the kubelet and calls policy start func (m *manager) Start(...) error { ... // 状态重建(machineMem/assignMem) stateImpl, err := state.NewCheckpointState(m.stateFileDirectory, memoryManagerStateFileName, m.policy.Name()) ... m.state = stateImpl // 内存分配状态校验 err = m.policy.Start(m.state) ... // 获取可用内存 m.allocatableMemory = m.policy.GetAllocatableMemory(m.state) return nil } func (p *staticPolicy) validateState(s state.State) error { // 持久化的机器内存状态 machineState := s.GetMachineState() // 持久化的分配状态 memoryAssignments := s.GetMemoryAssignments() // 检查及初始化机器内存状态 if len(machineState) == 0 { // Machine state cannot be empty when assignments exist if len(memoryAssignments) != 0 { return fmt.Errorf("[memorymanager] machine state can not be empty when it has memory assignments") } defaultMachineState := p.getDefaultMachineState() s.SetMachineState(defaultMachineState) return nil } // 扣减分配给容器内存 expectedMachineState := p.getDefaultMachineState() for pod, container := range memoryAssignments { for containerName, blocks := range container { for _, b := range blocks { // 容器申请的内存 requestedSize := b.Size // 容器内存归属的NUMA节点 for _, nodeID := range b.NUMAAffinity { // NUMA节点合法 nodeState, ok := expectedMachineState[nodeID] ... // 资源类型检查 memoryState, ok := nodeState.MemoryMap[b.Type] ... // 内存扣减 if memoryState.Free >= requestedSize { memoryState.Reserved += requestedSize memoryState.Free -= requestedSize requestedSize = 0 continue } requestedSize -= memoryState.Free memoryState.Reserved += memoryState.Free memoryState.Free = 0 } } } } // 预期的内存状态和持久化的内存状态一致 if !areMachineStatesEqual(machineState, expectedMachineState) { return fmt.Errorf("[memorymanager] the expected machine state is different from the real one") } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
# 7.4.Allocate
memoryManager.Allocate()负责pod准入阶段为容器预分配内存资源,确保容器启动前以满足资源分配策略要求。// 准入阶段 func (m *resourceAllocator) Admit(attrs *lifecycle.PodAdmitAttributes) lifecycle.PodAdmitResult { ... for _, container := range append(pod.Spec.InitContainers, pod.Spec.Containers...) { ... if m.memoryManager != nil { m.memoryManager.Allocate(pod, &container) ... } } return admission.GetPodAdmitResult(nil) } // Allocate is called to pre-allocate memory resources during Pod admission. func (m *manager) Allocate(pod *v1.Pod, container *v1.Container) error { // 标记pod处于准入阶段 m.setPodPendingAdmission(pod) // 先清理过期状态 m.removeStaleState() m.Lock() defer m.Unlock() // 进行本次内存分配 m.policy.Allocate(m.state, pod, container) ... return nil } // 过期状态清理 func (m *manager) removeStaleState() { // 数据源就绪 if !m.sourcesReady.AllReady() { return } m.Lock() defer m.Unlock() // 活跃pod整理 activeAndAdmittedPods := m.activePods() if m.pendingAdmissionPod != nil { activeAndAdmittedPods = append(activeAndAdmittedPods, m.pendingAdmissionPod) } // 活跃container整理 activeContainers := make(map[string]map[string]struct{}) for _, pod := range activeAndAdmittedPods { activeContainers[string(pod.UID)] = make(map[string]struct{}) for _, container := range append(pod.Spec.InitContainers, pod.Spec.Containers...) { activeContainers[string(pod.UID)][container.Name] = struct{}{} } } // 分配内存的容器整理,不在活跃容器列表进行资源释放 assignments := m.state.GetMemoryAssignments() for podUID := range assignments { for containerName := range assignments[podUID] { if _, ok := activeContainers[podUID][containerName]; !ok { m.policyRemoveContainerByRef(podUID, containerName) } } } // 二次遍历containerMap映射,不在活跃container列表进行资源释放 m.containerMap.Visit(func(podUID, containerName, containerID string) { if _, ok := activeContainers[podUID][containerName]; !ok { m.policyRemoveContainerByRef(podUID, containerName) } }) } func (m *manager) policyRemoveContainerByRef(podUID string, containerName string) { // 容器清理+资源释放 m.policy.RemoveContainer(m.state, podUID, containerName) // 缓存清理 m.containerMap.RemoveByContainerRef(podUID, containerName) } // 静态策略执行容器资源释放清理 func (p *staticPolicy) RemoveContainer(s state.State, podUID string, containerName string) { // 容器分配的内存信息 blocks := s.GetMemoryBlocks(podUID, containerName) ... // 容器分配资源释放,checkpoint持久化 s.Delete(podUID, containerName) // 更新空闲+分配预留内存 machineState := s.GetMachineState() for _, b := range blocks { releasedSize := b.Size // 按照容器分配内存NUMA节点更新 for _, nodeID := range b.NUMAAffinity { machineState[nodeID].NumberOfAssignments-- ... // 更新free和reserved if nodeResourceMemoryState.Reserved < releasedSize { releasedSize -= nodeResourceMemoryState.Reserved nodeResourceMemoryState.Free += nodeResourceMemoryState.Reserved nodeResourceMemoryState.Reserved = 0 continue } // the reserved memory big enough to satisfy the released memory nodeResourceMemoryState.Free += releasedSize nodeResourceMemoryState.Reserved -= releasedSize releasedSize = 0 } } // 更新machineState,持久化checkpoint s.SetMachineState(machineState) } // 内存分配 func (p *staticPolicy) Allocate(s state.State, pod *v1.Pod, container *v1.Container) error { // 专用内存管理只用于guaranteed pod if v1qos.GetPodQOS(pod) != v1.PodQOSGuaranteed { return nil } podUID := string(pod.UID) // 容器分配过内存 if blocks := s.GetMemoryBlocks(podUID, container.Name); blocks != nil { // 更新内存复用池信息 p.updatePodReusableMemory(pod, container, blocks) return nil } // 调用topologyManager获取亲和提示 hint := p.affinity.GetAffinity(podUID, container.Name) // 容器申请内存大小 requestedResources, err := getRequestedResources(container) ... // 获取机器内存分配信息 machineState := s.GetMachineState() bestHint := &hint // topologyManager无亲和提示,计算默认亲和性 if hint.NUMANodeAffinity == nil { // 计算numa组合 defaultHint, err := p.getDefaultHint(machineState, pod, requestedResources) ... // 计算的hint不是推荐的,topologyManager的推荐更优 if !defaultHint.Preferred && bestHint.Preferred { return fmt.Errorf("[memorymanager] failed to find the default preferred hint") } // 更新bestHint bestHint = defaultHint } // 计算亲和节点是否满足需求,不满足扩展hint if !isAffinitySatisfyRequest(machineState, bestHint.NUMANodeAffinity, requestedResources) { extendedHint, err := p.extendTopologyManagerHint(machineState, pod, requestedResources, bestHint.NUMANodeAffinity) ... // 扩展的不是推荐的,bestHint更优 if !extendedHint.Preferred && bestHint.Preferred { return fmt.Errorf("[memorymanager] failed to find the extended preferred hint") } // 更新bestHint bestHint = extendedHint } // 遍历资源请求,向每个请求分配block var containerBlocks []state.Block maskBits := bestHint.NUMANodeAffinity.GetBits() for resourceName, requestedSize := range requestedResources { // update memory blocks containerBlocks = append(containerBlocks, state.Block{ NUMAAffinity: maskBits, Size: requestedSize, Type: resourceName, }) // 分配内存 podReusableMemory := p.getPodReusableMemory(pod, bestHint.NUMANodeAffinity, resourceName) if podReusableMemory >= requestedSize { requestedSize = 0 } else { requestedSize -= podReusableMemory } // 更新节点内存分配状态,持久化checkpoint p.updateMachineState(machineState, maskBits, resourceName, requestedSize) } // 更新内存复用池 p.updatePodReusableMemory(pod, container, containerBlocks) // 更新机器内存分配状态 s.SetMachineState(machineState) // 更新分配的内存块 s.SetMemoryBlocks(podUID, container.Name, containerBlocks) // 同步init container的复用信息 p.updateInitContainersMemoryBlocks(s, pod, container, containerBlocks) return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
注意
1.
Allocate()申请的内存暂存于state,容器创建前生成配置会从state获取内存分配信息2.
memoryManager分配内存后,cgroup只会配置上限及可访问的NUMA节点,无法细分为NUMA节点限制多少内存3.
state维护的复用内存会在业务容器启动调用AddContainer()释放
# 8.devicePlugin
# 8.1.简介
k8s提供了一个设备插件框架,可以用来将系统硬件资源发布到kubelet,实现的设备插件可以手动部署或daemonset部署,无需定制集群组件本身,支持的设备涉及GPU、高性能NIC、FPGA和infiniBand适配器等。
分工
1.
kubelet:运行于节点。负责节点信息上报、启动和销毁Pod,向Pod分配网络、存储及资源2.
devicePlugin:设备插件,负责向kubelet注册设备,负载分配真实的资源,维护节点资源信息3.
deviceManager:设备插件管理器,负责管理注册的devicePlugin及设备资源调度
# 8.2.交互流程
devicePlugin的实现可以分为插件注册和插件调用两部分,设备插件启动会向kubelet发起注册以感知设备,pod准入申请资源会调用设备插件API分配资源,两者通信都基于unix.socket。// 插件注册(deviceManager实现) service Registration { rpc Register(RegisterRequest) returns (Empty) {} }1
2
3
4设备插件成功发布到
kubelet,Pod准入检查会根据命中的类型调用设备插件分配资源,相应的设备插件需实现规定接口以供上游调用。// DevicePlugin is the service advertised by Device Plugins service DevicePlugin { // 获取设备插件选项 rpc GetDevicePluginOptions(Empty) returns (DevicePluginOptions) {} // 返回Device列表构成的数据流,设备状态变化或消失会返回新的列表 rpc ListAndWatch(Empty) returns (stream ListAndWatchResponse) {} // 向一组可用的设备返回优选设备用来分配 // 返回的优选设备结果不一定是最终分配方案 // 此处只是为了让设备管理器在可能的情况下做出更有意义的决定 rpc GetPreferredAllocation(PreferredAllocationRequest) returns (PreferredAllocationResponse) {} // Pod准入期间向设备插件申请资源,通知kubelet令Device可在容器中访问的具体步骤 rpc Allocate(AllocateRequest) returns (AllocateResponse) {} // 设备插件注册阶段根据需要被调用,调用发生在容器启动前 // 将设备提供给容器使用前,设备插件可以运行一些诸如重置设备的动作 rpc PreStartContainer(PreStartContainerRequest) returns (PreStartContainerResponse) {} }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
注意
1.
kubelet已注册的插件信息不会持久化,重启会删除插件的socket文件,触发插件再次注册2.
devicePlugin注册后会向kubelet发送管理的设备列表,kubelet将资源发布到API服务器作为节点状态更新的一部分3.
devicePlugin以Pod形式部署需挂载kubelet.socket文件,确保交互正常
# 8.3.AMD算力插件
AMD GPU插件基于github.com/kubevirt/device-plugin-manager/pkg/dpm框架实现,基于框架注册grpc handle处理器及封装grpc server,设备插件只需实现对应接口及部分主动行为就可以接入deviceManager。func main() { ... // 传递心跳及资源更新 l := Lister{ ResUpdateChan: make(chan dpm.PluginNameList), Heartbeat: make(chan bool), } // 插件实现 manager := dpm.NewManager(&l) ... go func() { /sys/class/kfd only exists if ROCm kernel/driver is installed // 启动检测 if _, err := os.Stat(path); err == nil { l.ResUpdateChan <- []string{"gpu"} } }() // 启动管理器 manager.Run() }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20Run()是devicePlugin Manager的主循环实现,用于管理插件生命周期,负责启动gRPC服务,探测资源变更,监听socket变化以触发重新注册以及退出信号捕获。// Run starts the Manager. It sets up the infrastructure and handles system signals, Kubelet socket // watch and monitoring of available resources as well as starting and stoping of plugins. func (dpm *Manager) Run() { // 退出信号监听 signalCh := make(chan os.Signal, 1) signal.Notify(signalCh, syscall.SIGTERM, syscall.SIGQUIT, syscall.SIGINT) ... // 文件系统事件监听 fsWatcher, _ := fsnotify.NewWatcher() defer fsWatcher.Close() // 监听/var/lib/kubelet/device-plugins目录,主要是kubelet.sock的创建/删除 fsWatcher.Add(pluginapi.DevicePluginPath) // 维护运行的插件名和插件结构体 var pluginMap = make(map[string]devicePlugin) pluginsCh := make(chan PluginNameList) defer close(pluginsCh) // 启动协程动态探测设备插件列表(gpu) go dpm.lister.Discover(pluginsCh) // 主循环 HandleSignals: for { select { // 新插件 case newPluginsList := <-pluginsCh: // 初始化及启动插件 dpm.handleNewPlugins(pluginMap, newPluginsList) // 文件系统事件 case event := <-fsWatcher.Events: // kubelet.sock相关事件 if event.Name == pluginapi.KubeletSocket { // kubelet启动,重启插件服务 if event.Op&fsnotify.Create == fsnotify.Create { dpm.startPluginServers(pluginMap) } // kubelet重启,停止插件服务 if event.Op&fsnotify.Remove == fsnotify.Remove { dpm.stopPluginServers(pluginMap) } } // 退出信号 case s := <-signalCh: switch s { // 关心的退出信号 case syscall.SIGTERM, syscall.SIGQUIT, syscall.SIGINT: // 停止插件服务 dpm.stopPlugins(pluginMap) break HandleSignals } } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51ListAndWatch()用户获取当前设备列表,持续向kubelet报告GPU设备列表及其健康状态,作为插件和kubelet间的设备生命周期同步桥梁。// 设备列表上报 func (p *Plugin) ListAndWatch(...) error { // 探测本机可用的AMD GPU p.AMDGPUs = amdgpu.GetAMDGPUs() ... // 基于hwloc查询NUMA拓扑 func() { ... // 初始化硬件拓扑探测器,确认GPU和NUMA节点关系 hw.Init() defer hw.Destroy() ... for id := range p.AMDGPUs { // 每张GPU初始化为Healthy dev := &pluginapi.Device{ ID: id, Health: pluginapi.Healthy, } devs[i] = dev ... // 查询GPU绑定的NUMA节点列表 numas, err := hw.GetNUMANodes(id) ... // 记录GPU关联NUMA节点,用于kubelet亲和调度 for j, v := range numas { numaNodes[j] = &pluginapi.NUMANode{ ID: int64(v), } } dev.Topology = &pluginapi.TopologyInfo{ Nodes: numaNodes, } } }() // 首次发送设备列表给kubelet s.Send(&pluginapi.ListAndWatchResponse{Devices: devs}) for { select { // 心跳检测 case <-p.Heartbeat: ... // 健康探测 if simpleHealthCheck() { health = pluginapi.Healthy } // 更新设备健康状态 for i := 0; i < len(p.AMDGPUs); i++ { devs[i].Health = health } // 再次发送设备列表给kubelet s.Send(&pluginapi.ListAndWatchResponse{Devices: devs}) } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56Allocate()方法负责容器准入检查基于指定GPU设备向kubelet声明挂载的设备列表、环境变量和Mount目录,这些声明最终会关联到containerSpec,用于容器访问GPU设备。// 执行与设备相关的特定操作,向kubelet声明设备如何挂载到容器 func (p *Plugin) Allocate(ctx context.Context, r *pluginapi.AllocateRequest) (*pluginapi.AllocateResponse, error) { ... // 遍历容器的GPU请求 for _, req := range r.ContainerRequests { ... // ROCM的GPU容器必须访问/dev/kfd // 每个节点只有一个/dev/kfd(HSA计算框架),不属于特定GPU,统一挂载 dev = new(pluginapi.DeviceSpec) dev.HostPath = "/dev/kfd" dev.ContainerPath = "/dev/kfd" dev.Permissions = "rw" car.Devices = append(car.Devices, dev) // 遍历容器请求的GPUID for _, id := range req.DevicesIDs { // 插件维护的p.AMDGPUs映射获取GPU对应的设备信息 for k, v := range p.AMDGPUs[id] { // 构造设备路径 // /dev/dri/card0、/dev/dri/renderD128 devpath := fmt.Sprintf("/dev/dri/%s%d", k, v) dev = new(pluginapi.DeviceSpec) dev.HostPath = devpath dev.ContainerPath = devpath dev.Permissions = "rw" car.Devices = append(car.Devices, dev) } } // 容器的资源配置加入AllocateResponse response.ContainerResponses = append(response.ContainerResponses, &car) } return &response, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
注意
1.
devicePlugin启动上报设备后,kubelet反向拨号设备插件实现列表获取与资源绑定2.
kubelet和devicePlugin交互基于unix.sock+gRPC3.
Allocate阶段设备插件向kubelet声明需要绑定的驱动和设备(nvidia声明环境遍历,由nvidia-runtime-hook绑定)4.
devicePlugin实现上支持NUMA亲和拓扑调度
# 9.deviceManager
# 9.1.简介
deviceManager是负责交互设备插件的设备生命周期管理模块,用户管理外部硬件的发现、健康检查及资源分配,该模块管理的所有插件设备会写入node.status.Capacity和node.status.Allocatable资源字段,供调度器作为节点绑定基准。// Manager manages all the Device Plugins running on a node. type Manager interface { // 启动设备插件注册服务 Start(activePods ActivePodsFunc, sourcesReady config.SourcesReady) error // 设备分配申请 Allocate(pod *v1.Pod, container *v1.Container) error //根据已经分配设备更新节点的资源信息 UpdatePluginResources(node *schedulerframework.NodeInfo, attrs *lifecycle.PodAdmitAttributes) error // 停止设备管理器 Stop() error // 获取缓存的设备配置,用于容器启动时的runtime配置 GetDeviceRunContainerOptions(pod *v1.Pod, container *v1.Container) (*DeviceRunContainerOptions, error) // 获取节点设备资源总量/可分配量/已注册+不活跃设备插件列表 GetCapacity() (v1.ResourceList, v1.ResourceList, []string) // 插件监听机制,获取用于处理设备注册事件的handler GetWatcherHandler() cache.PluginHandler // 获取pod/container已分配的设备列表,用于运行时设备追踪和日志分析 GetDevices(podUID, containerName string) ResourceDeviceInstances // 获取当前manager已知的可用设备 GetAllocatableDevices() ResourceDeviceInstances // 检查设置相关资源是否重置(checkpoint文件不存在,节点重建) ShouldResetExtendedResourceCapacity() bool // 获取给定容器的NUMA拓扑结构的设备分配建议,支持拓扑感知调度 GetTopologyHints(pod *v1.Pod, container *v1.Container) map[string][]topologymanager.TopologyHint // 用于支持调度决策,影响Pod粒度的拓扑感知 GetPodTopologyHints(pod *v1.Pod) map[string][]topologymanager.TopologyHint // 清理终止Pod绑定的设备资源,防止设备泄漏和冲突 UpdateAllocatedDevices() } // ManagerImpl is the structure in charge of managing Device Plugins. type ManagerImpl struct { // 插件注册监听的unix socket socketname string socketdir string // 注册的设备插件维护的通信端点 endpoints map[string]endpointInfo // Key is ResourceName ... // 交互设备插件的gRPC服务 server *grpc.Server ... // podManager维护的活跃Pod activePods ActivePodsFunc sourcesReady config.SourcesReady // 设备状态更新的统一回调函数(ListAndWatch结果处理) callback monitorCallback // 缓存插件注册过的设备(resourceName -> DeviceInstances) allDevices ResourceDeviceInstances // 正常运行的设备列表(resourceName-->plugin-->deviceID) healthyDevices map[string]sets.String // 离线/故障的设备列表(resourceName-->plugin-->deviceID) unhealthyDevices map[string]sets.String // 已分配出去的设备ID列表(resourceName-->plugin-->deviceID) allocatedDevices map[string]sets.String // Pod已分配的设备列表(Pod-->Container--Resource-->Device ID) podDevices *podDevices // checkpoint检查点管理器 checkpointManager checkpointmanager.CheckpointManager // 本机的NUMA拓扑信息,决定设备绑定的NUMA节点(cpu/mem/device对齐) numaNodes []int // 提供设备与NUMA拓扑的亲和信息,协助topologyManager作出资源对齐决定 topologyAffinityStore topologymanager.Store // 资源复用维护(init container释放资源的二次分配) devicesToReuse PodReusableDevices // 准入阶段标识 pendingAdmissionPod *v1.Pod }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
# 9.2.NewManagerImpl
作为
containerManager的子模块,deviceManager会在containerManager初始化时创建,此时gRPC服务、插件监听和健康检查还未启动。// TODO(vmarmol): Add limits to the system containers. // Takes the absolute name of the specified containers. // Empty container name disables use of the specified container. func NewContainerManager(...) (ContainerManager, error) { ... // 开启设备插件 if devicePluginEnabled { // 初始化deviceManager cm.deviceManager, err = devicemanager.NewManagerImpl(machineInfo.Topology, cm.topologyManager) // 注册到拓扑管理器 cm.topologyManager.AddHintProvider(cm.deviceManager) } else { // 未开启走空实现 cm.deviceManager, err = devicemanager.NewManagerStub() } ... } // NewManagerImpl creates a new manager. func NewManagerImpl(topology []cadvisorapi.Node, topologyAffinityStore topologymanager.Store) (*ManagerImpl, error) { // kubelet.sock地址 socketPath := pluginapi.KubeletSocket ... // 初始化deviceManager return newManagerImpl(socketPath, topology, topologyAffinityStore) } // create a new deviceManagerImpl func newManagerImpl(...) (*ManagerImpl, error) { ... // 整理机器的NUMA节点 for _, node := range topology { numaNodes = append(numaNodes, node.Id) } ... manager := &ManagerImpl{ // 设备插件通信端点 endpoints: make(map[string]endpointInfo), // socket名称 socketname: file, // socket目录 socketdir: dir, // 设备相关 allDevices: NewResourceDeviceInstances(), healthyDevices: make(map[string]sets.String), unhealthyDevices: make(map[string]sets.String), allocatedDevices: make(map[string]sets.String), podDevices: newPodDevices(), // 当前机器NUMA节点 numaNodes: numaNodes, // topologyManager拓扑管理器 topologyAffinityStore: topologyAffinityStore, // 设备复用池(init container释放的二次分配) devicesToReuse: make(PodReusableDevices), } // 设备更新的状态回调 manager.callback = manager.genericDeviceUpdateCallback ... // checkpoint管理器,持久化分配状态(/var/lib/kubelet/device-plugins) manager.checkpointManager = checkpointmanager.NewCheckpointManager(dir) return manager, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63注意
1.
deviceManager和devicePlugindevicePlugin和kubelet交互基于/var/lib/kubelet/device-plugins/xxx.sock2.
callback用于将调整或注册的设备更新到维护的设备列表3.
checkpointManager用于持久化设备分配信息及重启恢复
# 9.3.Start
deviceManager作为containerManager启动的一部分,会初始化activePods和sourcesReady回调,读取checkpoint重建设备的分配状态,启动gRPC Server提供设备插件注册服务。func (cm *containerManagerImpl) Start(...) error { ... // Starts device manager. cm.deviceManager.Start(devicemanager.ActivePodsFunc(activePods), sourcesReady) ... } // 启动deviceManager,读取checkpoint重建设备分配状态 func (m *ManagerImpl) Start(activePods ActivePodsFunc, sourcesReady config.SourcesReady) error { // podManager维护的活跃Pod回调 m.activePods = activePods m.sourcesReady = sourcesReady // 加载checkpoint的设备分配状态(kubelet_internal_checkpoint) m.readCheckpoint() ... // 启动selinux时,设备适当的selinux label,避免设备插件无法访问device-plugins目录 if selinux.GetEnabled() { selinux.SetFileLabel(m.socketdir, config.KubeletPluginsDirSELinuxLabel) ... } // 清理device-plugins目录所有socket文件(设备插件会监听kubelet.sock) m.removeContents(m.socketdir) // 监听unix socket,等待插件连接及注册 s, err := net.Listen("unix", socketPath) ... m.wg.Add(1) m.server = grpc.NewServer([]grpc.ServerOption{}...) // 向gRPC注册插件注册服务 pluginapi.RegisterRegistrationServer(m.server, m) go func() { defer m.wg.Done() // 启动gRPC服务 m.server.Serve(s) }() return nil } // 基于checkpoint重建设备分配状态 func (m *ManagerImpl) readCheckpoint() error { // 优先读取v2格式checkpoint cp, err := m.getCheckpointV2() if err != nil { ... // 读取v1格式checkpoint cp, errv1 = m.getCheckpointV1() ... } m.mutex.Lock() defer m.mutex.Unlock() // 提取Pod/container分配资源及上次注册成功的设备资源 podDevices, registeredDevs := cp.GetDataInLatestFormat() // 重建Pod/container已分配设备状态 m.podDevices.fromCheckpointData(podDevices) // 重建已分配的设备状态 m.allocatedDevices = m.podDevices.devices() // 设备重新注册前,设备状态初始化为空 for resource := range registeredDevs { m.healthyDevices[resource] = sets.NewString() m.unhealthyDevices[resource] = sets.NewString() m.endpoints[resource] = endpointInfo{e: newStoppedEndpointImpl(resource), opts: nil} } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
# 9.4.Register
deviceManager对外提供Register()接口,供设备插件注册发布资源,生成Endpoint对象作为通信端点,基于长连接实时获取设备状态。// Register registers a device plugin. func (m *ManagerImpl) Register(...) (*pluginapi.Empty, error) { var versionCompatible bool // 检查插件版本是否支持(v1beta1) for _, v := range pluginapi.SupportedVersions { if r.Version == v { versionCompatible = true break } } ... // 检查resourceName是否合法 // 1.不能属于kubernetes.io,必须有自己的domain // 2.resourceName不能包含requests.前缀 // 3.resourceValue只能是整数 if !v1helper.IsExtendedResourceName(v1.ResourceName(r.ResourceName)) { return &pluginapi.Empty{}, err } // 注册插件endpoint go m.addEndpoint(r) return &pluginapi.Empty{}, nil } // 端点注册 func (m *ManagerImpl) addEndpoint(r *pluginapi.RegisterRequest) { // 初始化plugin endpoint(检查连接状态) new, err := newEndpointImpl(filepath.Join(m.socketdir, r.Endpoint), r.ResourceName, m.callback) ... // 注册 m.registerEndpoint(r.ResourceName, r.Options, new) go func() { // 启动plugin endpoint m.runEndpoint(r.ResourceName, new) }() } func (m *ManagerImpl) runEndpoint(resourceName string, e endpoint) { // 启动endpoint同步设备主循环 e.run() // run结束释放连接 e.stop() m.mutex.Lock() defer m.mutex.Unlock() // 标记当前设备不健康 if old, ok := m.endpoints[resourceName]; ok && old.e == e { m.markResourceUnhealthy(resourceName) } } // 基于gRPC与插件建立ListAndWatch流连接,持续接收设备列表更新,调用回调将设备状态传递给deviceManager func (e *endpointImpl) run() { // 获取streamClient stream, err := e.client.ListAndWatch(context.Background(), &pluginapi.Empty{}) ... // 主循环 for { // 接收设备状态 response, err := stream.Recv() ... devs := response.Devices ... // 回调传递设备状态 e.callback(e.resourceName, newDevs) } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
# 9.5.Allocate
PodAdmit准入阶段,kubelet会调用deviceManager.Allocate()申请分配设备及挂载信息,分配的设备会缓存及持久化到checkpoint,供后续容器创建使用。// 准入阶段 func (m *resourceAllocator) Admit(attrs *lifecycle.PodAdmitAttributes) lifecycle.PodAdmitResult { ... for _, container := range append(pod.Spec.InitContainers, pod.Spec.Containers...) { m.deviceManager.Allocate(pod, &container) ... } return admission.GetPodAdmitResult(nil) } // 调用插件分配设备 func (m *ManagerImpl) Allocate(pod *v1.Pod, container *v1.Container) error { // 标记Pod准入 m.setPodPendingAdmission(pod) // 初始化设备复用池 if _, ok := m.devicesToReuse[string(pod.UID)]; !ok { m.devicesToReuse[string(pod.UID)] = make(map[string]sets.String) } // 清理复用池其它Pod缓存(资源串行分配的) for podUID := range m.devicesToReuse { if podUID != string(pod.UID) { delete(m.devicesToReuse, podUID) } } // 当前是init容器,发起设备申请及记录分配设备 for _, initContainer := range pod.Spec.InitContainers { if container.Name == initContainer.Name { // 分配设备 m.allocateContainerResources(pod, container, m.devicesToReuse[string(pod.UID)]) ... // init容器分配的设备全部加入复用池 m.podDevices.addContainerAllocatedResources(string(pod.UID), container.Name, m.devicesToReuse[string(pod.UID)]) return nil } } // 普通容器申请分配 m.allocateContainerResources(pod, container, m.devicesToReuse[string(pod.UID)]) ... // 清理复用池分配出去的设备 m.podDevices.removeContainerAllocatedResources(string(pod.UID), container.Name, m.devicesToReuse[string(pod.UID)]) return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44allocateContainerResources会遍历容器的所有设备类型请求,调用设备插件的Allocate()接口完成资源申请,缓存结果供容器运行时注入。// 向container分配申请的设备(复用部分init container资源),分配状态持久化到checkpoint func (m *ManagerImpl) allocateContainerResources(pod *v1.Pod, container *v1.Container, devicesToReuse map[string]sets.String) error { ... // 只遍历limit(devicePlugin类型资源必须request==limit) for k, v := range container.Resources.Limits { resource := string(k) needed := int(v.Value()) // 不是可识别的设备资源请求(resourceName) if !m.isDevicePluginResource(resource) { continue } // 分配前释放终止Pod设备资源 if !allocatedDevicesUpdated { m.UpdateAllocatedDevices() allocatedDevicesUpdated = true } // 获取需要分配的设备(优先从设备复用池扣减) allocDevices, err := m.devicesToAllocate(podUID, contName, resource, needed, devicesToReuse[resource]) ... // 没有要分配的设备 if allocDevices == nil || len(allocDevices) <= 0 { continue } needsUpdateCheckpoint = true // 找出设备插件 m.mutex.Lock() eI, ok := m.endpoints[resource] m.mutex.Unlock() // 未找到设备插件,未注册 if !ok { m.mutex.Lock() // 重新同步allocatedDevices m.allocatedDevices = m.podDevices.devices() m.mutex.Unlock() return fmt.Errorf("unknown Device Plugin %s", resource) } // ... devs := allocDevices.UnsortedList() // 发起gRPC调用执行分配 resp, err := eI.e.allocate(devs) if err != nil { m.mutex.Lock() // 无法分配,重新同步allocatedDevices m.allocatedDevices = m.podDevices.devices() m.mutex.Unlock() return err } // 无法分配 if len(resp.ContainerResponses) == 0 { return fmt.Errorf("no containers return in allocation response %v", resp) } // 初始化NUMA拓扑 allocDevicesWithNUMA := checkpoint.NewDevicesPerNUMA() // Update internal cached podDevices state. m.mutex.Lock() // 解析NUMA拓扑构造分配结构 for dev := range allocDevices { // 申请的设备没有拓扑信息 if m.allDevices[resource][dev].Topology == nil || len(m.allDevices[resource][dev].Topology.Nodes) == 0 { // 归为无拓扑资源 allocDevicesWithNUMA[nodeWithoutTopology] = append(allocDevicesWithNUMA[nodeWithoutTopology], dev) continue } // 获取设备关联所有NUMA节点,记录NUMA节点与设备映射 for idx := range m.allDevices[resource][dev].Topology.Nodes { node := m.allDevices[resource][dev].Topology.Nodes[idx] allocDevicesWithNUMA[node.ID] = append(allocDevicesWithNUMA[node.ID], dev) } } m.mutex.Unlock() // 更新podDevices缓存 m.podDevices.insert(podUID, contName, resource, allocDevicesWithNUMA, resp.ContainerResponses[0]) } // 分配出设备,持久化checkpoint记录分配状态 if needsUpdateCheckpoint { return m.writeCheckpoint() } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
注意
1.
Allocate()分配的设备会缓存在podDevices,供container创建时生成配置使用2.
memoryManager和cpuManager分配的资源不同,供PreCreateContainer阶段影响容器配置
# 9.topologyManager
# 9.1.简介
topologyManager是5大准入控制器之一,用于全局资源角度给出拓扑建议。cpu/mem/device管理器实现HintProvider及注册到拓扑管理器,可以在资源分配时占用相同或贴近的NUMA节点,促使资源交互时提高性能,支持的作用域有Pod和container。// TopologyHint is a struct containing the NUMANodeAffinity for a Container type TopologyHint struct { // bitMask表示的NUMA节点组合 NUMANodeAffinity bitmask.BitMask // 用于标识当前组合是否优先考虑(资源足够前提下NUMA节点最少) Preferred bool } // 拓扑接口 type HintProvider interface { // 获取容器层级拓扑建议 GetTopologyHints(pod *v1.Pod, container *v1.Container) map[string][]TopologyHint // 获取Pod层级拓扑建议 GetPodTopologyHints(pod *v1.Pod) map[string][]TopologyHint // 资源分配接口 Allocate(pod *v1.Pod, container *v1.Container) error } type manager struct { // Topology Manager Scope scope Scope } // Manager interface provides methods for Kubelet to manage pod topology hints type Manager interface { // 准入控制 lifecycle.PodAdmitHandler // 注册HintProvider AddHintProvider(HintProvider) // 跟踪容器NUMA绑定 AddContainer(pod *v1.Pod, container *v1.Container, containerID string) // 清理容器NUMA绑定 RemoveContainer(containerID string) error // 存储pod/container的拓扑状态 Store }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42topologyManager提供4种拓扑策略,每种策略有各自的实现,代表不同的准入行为,拓扑管理器根据所选策略初始化及作用域初始化。// create containerManager for resource managing func NewContainerManager(...) (ContainerManager, error) { ... if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.TopologyManager) { // 初始化拓扑管理器 cm.topologyManager, err = topologymanager.NewManager( machineInfo.Topology, nodeConfig.ExperimentalTopologyManagerPolicy, nodeConfig.ExperimentalTopologyManagerScope, ) ... } else { cm.topologyManager = topologymanager.NewFakeManager() } ... } // NewManager creates a new TopologyManager based on provided policy and scope func NewManager(topology []cadvisorapi.Node, topologyPolicyName string, topologyScopeName string) (Manager, error) { ... // 根据cadvisor提供的节点NUMA拓扑获取NUMAID for _, node := range topology { numaNodes = append(numaNodes, node.Id) } // NUMA节点数量检查,非None策略最大支持8个(避免拓扑组合状态爆炸) if topologyPolicyName != PolicyNone && len(numaNodes) > maxAllowableNUMANodes { return nil, fmt.Errorf("unsupported on machines with more than %v NUMA Nodes", maxAllowableNUMANodes) } // 策略初始化 var policy Policy switch topologyPolicyName { case PolicyNone: // None策略,空实现 policy = NewNonePolicy() case PolicyBestEffort: // BestEffort策略,尽量优选 policy = NewBestEffortPolicy(numaNodes) case PolicyRestricted: // Restricted策略,必须优选 policy = NewRestrictedPolicy(numaNodes) case PolicySingleNumaNode: // SingleNumaNode策略,必须单节点优选实现 policy = NewSingleNumaNodePolicy(numaNodes) default: return nil, fmt.Errorf("unknown policy: \"%s\"", topologyPolicyName) } var scope Scope switch topologyScopeName { case containerTopologyScope: // 容器作用域 scope = NewContainerScope(policy) case podTopologyScope: // Pod作用域 scope = NewPodScope(policy) default: return nil, fmt.Errorf("unknown scope: \"%s\"", topologyScopeName) } manager := &manager{ scope: scope, } return manager, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
# 9.2.调用
topologyManager遍历Pod container,依次调用HintProvider.GetTopologyHints()方法产生TopologyHint及合并,寻求一个最优的TopologyHint,HintProvider基于最优的TopologyHint向容器分配资源。目前,mem/cpu/device均实现了HintProvider接口,初始化时会调用topologyManager.AddHintProvider()注册。// Takes the absolute name of the specified containers. // Empty container name disables use of the specified container. func NewContainerManager(...) (ContainerManager, error) { ... // memoryManager cm.memoryManager, err = memorymanager.NewManager(...,cm.topologyManager) cm.topologyManager.AddHintProvider(cm.memoryManager) ... // cpuManager cm.cpuManager, err = cpumanager.NewManager(...,cm.topologyManager) cm.topologyManager.AddHintProvider(cm.cpuManager) ... // deviceManager cm.deviceManager, err = devicemanager.NewManagerImpl(machineInfo.Topology, cm.topologyManager) cm.topologyManager.AddHintProvider(cm.deviceManager) ... } // 注册HintProvider func (m *manager) AddHintProvider(h HintProvider) { m.scope.AddHintProvider(h) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22kubelet调度Pod时,会先调用kl.canAdmitPod()进行准入检查,其实就是调用注册的准入控制器进行资源申请及校验,全部通过才会放行Pod。// canAdmitPod determines if a pod can be admitted, and gives a reason if it cannot. func (kl *Kubelet) canAdmitPod(pods []*v1.Pod, pod *v1.Pod) (bool, string, string) { // 1.evictionAdmitHandler(资源压力检查,内存、磁盘...) // 2.sysctlsAllowlist(Pod申请sysctls检查) // 3.containerManager.GetAllocateResourcesPodAdmitHandler(资源分配检查) // 4.lifecycle.NewPredicateAdmitHandler(常规条件检查,volume挂载、节点亲和、节点健康状态...) // 5.shutdownAdmitHandler(节点关机保护) attrs := &lifecycle.PodAdmitAttributes{Pod: pod, OtherPods: pods} for _, podAdmitHandler := range kl.admitHandlers { if result := podAdmitHandler.Admit(attrs); !result.Admit { return false, result.Reason, result.Message } } return true, "", "" }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15containerManager是资源管理的准入控制器,注册会调用containerManager.GetAllocateResourcesPodAdmitHandler()根据配置启用topologyManager或组合资源控制器。// registry AdmitHandler func (cm *containerManagerImpl) GetAllocateResourcesPodAdmitHandler() lifecycle.PodAdmitHandler { // 启用拓扑管理器 if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.TopologyManager) { // 总和调度cpuManager、memoryManager和deviceManager的拓扑亲和性 return cm.topologyManager } // 临时兼容方案(不具备拓扑亲和及可扩展性) return &resourceAllocator{cm.cpuManager, cm.memoryManager, cm.deviceManager} }1
2
3
4
5
6
7
8
9
10
# 9.3.cpuManager实现
cpuManager实现了HintProvider接口,初始化会保存topologyManager获取计算的bestHint,也会作为HintProvider注册到topologyManager。func (m *manager) GetTopologyHints(pod *v1.Pod, container *v1.Container) map[string][]topologymanager.TopologyHint { // 标记Pod准入 m.setPodPendingAdmission(pod) // 拓扑亲和计算前释放过期CPU资源 m.removeStaleState() // N选K生成Hint(container粒度) return m.policy.GetTopologyHints(m.state, pod, container) } func (m *manager) GetPodTopologyHints(pod *v1.Pod) map[string][]topologymanager.TopologyHint { // 标记Pod准入 m.setPodPendingAdmission(pod) // 释放过期资源 m.removeStaleState() // N选K生成Hint(Pod粒度,init容器和业务容器请求和取最大值) return m.policy.GetPodTopologyHints(m.state, pod) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# 9.4.memManager实现
memManager实现了HintProvider接口,初始化会保存topologyManager获取计算的bestHint,也会作为HintProvider注册到topologyManager。// GetTopologyHints returns the topology hints for the topology manager func (m *manager) GetTopologyHints(pod *v1.Pod, container *v1.Container) map[string][]topologymanager.TopologyHint { // 标记Pod准入 m.setPodPendingAdmission(pod) // 释放过期资源 m.removeStaleState() // N选K组合 return m.policy.GetTopologyHints(m.state, pod, container) } // GetPodTopologyHints returns the topology hints for the topology manager func (m *manager) GetPodTopologyHints(pod *v1.Pod) map[string][]topologymanager.TopologyHint { // The pod is during the admission phase. We need to save the pod to avoid it // being cleaned before the admission ended m.setPodPendingAdmission(pod) // Garbage collect any stranded resources before providing TopologyHints m.removeStaleState() // Delegate to active policy return m.policy.GetPodTopologyHints(m.state, pod) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 9.5.deviceManager
deviceManager实现了HintProvider接口,初始化会保存topologyManager获取计算的bestHint,也会作为HintProvider注册到topologyManager。// ensures the Device Manager is consulted when Topology Aware Hints for each container are created. func (m *ManagerImpl) GetTopologyHints(pod *v1.Pod, container *v1.Container) map[string][]topologymanager.TopologyHint { // 标记Pod准入 m.setPodPendingAdmission(pod) // 释放过期分配设备 m.UpdateAllocatedDevices() // 遍历可用设备完成Hint组合 deviceHints := make(map[string][]topologymanager.TopologyHint) for resourceObj, requestedObj := range container.Resources.Limits { ... // 插件注册过 if m.isDevicePluginResource(resource) { // 设备无亲和属性 if aligned := m.deviceHasTopologyAlignment(resource); !aligned { deviceHints[resource] = nil continue } // 容器已分配设备 allocated := m.podDevices.containerDevices(string(pod.UID), container.Name, resource) if allocated.Len() > 0 { // 已分配与请求不一致 if allocated.Len() != requested { deviceHints[resource] = []topologymanager.TopologyHint{} continue } // 重建容器设备Hint deviceHints[resource] = m.generateDeviceTopologyHints(resource, allocated, sets.String{}, requested) continue } // 获取可用设备(Healthy) available := m.getAvailableDevices(resource) // 获取Pod可复用设备 reusable := m.devicesToReuse[string(pod.UID)][resource] // 可用+复用无法覆盖申请 if available.Union(reusable).Len() < requested { deviceHints[resource] = []topologymanager.TopologyHint{} continue } // 基于可用+复用组合Hint deviceHints[resource] = m.generateDeviceTopologyHints(resource, available, reusable, requested) } } return deviceHints } // ensures the Device Manager is consulted when Topology Aware Hints for Pod are created. func (m *ManagerImpl) GetPodTopologyHints(pod *v1.Pod) map[string][]topologymanager.TopologyHint { // 标记Pod准入 m.setPodPendingAdmission(pod) // 释放过期资源 m.UpdateAllocatedDevices() ... // 获取Pod申请设备(init和业务容器和取最大) accumulatedResourceRequests := m.getPodDeviceRequest(pod) for resource, requested := range accumulatedResourceRequests { // 插件设备无亲和属性 if aligned := m.deviceHasTopologyAlignment(resource); !aligned { deviceHints[resource] = nil continue } // 获取Pod已分配设备 allocated := m.podDevices.podDevices(string(pod.UID), resource) if allocated.Len() > 0 { // 已分配与请求不一致 if allocated.Len() != requested { deviceHints[resource] = []topologymanager.TopologyHint{} continue } // 基于已分配重建 deviceHints[resource] = m.generateDeviceTopologyHints(resource, allocated, sets.String{}, requested) continue } // 获取可用设备 available := m.getAvailableDevices(resource) // 可用设备无法满足分配 if available.Len() < requested { deviceHints[resource] = []topologymanager.TopologyHint{} continue } // 基于可用计算Hint组合 deviceHints[resource] = m.generateDeviceTopologyHints(resource, available, sets.String{}, requested) } return deviceHints }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
# 9.6.全局调度
# 10.containerManager
# 10.1.简介
containerManager是kubelet资源调度管理控制器,负责节点上的容器运行环境资源配置,基于全局统一调度cpuManager、memoryManager、cgroupManager、qosContainerManager、deviceManager和topologyManager,其实现的AdmitHandler接口会作为准入控制器注册到kubelet进行资源调度及分配。// Manages the containers running on a machine. type ContainerManager interface { // 初始化system cgroup,启动各资源管理器 Start(*v1.Node, ActivePodsFunc, config.SourcesReady, status.PodStatusProvider, internalapi.RuntimeService) error // 获取系统预留资源,不参与pod调度 SystemCgroupsLimit() v1.ResourceList // 获取节点配置(cgroup驱动/qos开启/保留策略) GetNodeConfig() NodeConfig // containerManager健康状态(健康检查) Status() Status // 初始化podContainerManager(pod-level cgroup管理) NewPodContainerManager() PodContainerManager // 获取系统挂载cgroup subsystem GetMountedSubsystems() *CgroupSubsystems // 获取顶级qos cgroup信息 GetQOSContainersInfo() QOSContainersInfo // 获取预留资源(system/kube),不参与调度 GetNodeAllocatableReservation() v1.ResourceList // 获取节点资源容量 GetCapacity() v1.ResourceList // 获取设备插件资源 GetDevicePluginResourceCapacity() (v1.ResourceList, v1.ResourceList, []string) // 更新qos cgroup为期望状态 UpdateQOSCgroups() error // 获取container的附加资源(devicePlugin提供) GetResources(pod *v1.Pod, container *v1.Container) (*kubecontainer.RunContainerOptions, error) // 刷新设备插件资源信息 UpdatePluginResources(*schedulerframework.NodeInfo, *lifecycle.PodAdmitAttributes) error // 生命周期控制器,处理资源分配前后的Hook InternalContainerLifecycle() InternalContainerLifecycle // pod所属cgroup根目录 GetPodCgroupRoot() string // 获取设备插件注册Handle(nodeLeaseController调用) GetPluginRegistrationHandler() cache.PluginHandler // 扩展资源清理 ShouldResetExtendedResourceCapacity() bool // 获取资源准入控制器 GetAllocateResourcesPodAdmitHandler() lifecycle.PodAdmitHandler // 获取节点可分配资源绝对值 GetNodeAllocatableAbsolute() v1.ResourceList // 资源监控接口 podresources.CPUsProvider podresources.DevicesProvider podresources.MemoryProvider } type containerManagerImpl struct { ... // cadvisor通信模块 cadvisorInterface cadvisor.Interface // 挂载点管理 mountUtil mount.Interface ... // system container systemContainers []*systemContainer // 周期任务 periodicTasks []func() // 节点挂载的cgroup subsystem subsystems *CgroupSubsystems // 节点node对象 nodeInfo *v1.Node // cgroup管理器 cgroupManager CgroupManager // 节点资源容量 capacity v1.ResourceList // 加上保留资源的节点容量 internalCapacity v1.ResourceList // pod顶级cgroup cgroupRoot CgroupName ... // qos cgroup管理器 qosContainerManager QOSContainerManager // 设备管理器 deviceManager devicemanager.Manager // cpu管理器 cpuManager cpumanager.Manager // 内存管理器 memoryManager memorymanager.Manager // 拓扑管理器 topologyManager topologymanager.Manager }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
# 10.2.初始化
NewContainerManager()用于初始化完整的containerManager实例,检查节点启动条件及构造各资源管理器。// Takes the absolute name of the specified containers. // Empty container name disables use of the specified container. func NewContainerManager(...) (ContainerManager, error) { // 获取cgroup subsystem subsystems, err := GetCgroupSubsystems() ... // 禁止运行至swap环境 if failSwapOn { // 检查swap开启 swapFile := "/proc/swaps" swapData, err := ioutil.ReadFile(swapFile) ... swapData = bytes.TrimSpace(swapData) // extra trailing \n swapLines := strings.Split(string(swapData), "\n") // 存在非header内容,swap开启 if len(swapLines) > 1 { return nil, fmt.Errorf("running with swap on is not supported, please disable swap! or set --fail-swap-on flag to false. /proc/swaps contained: %v", swapLines) } } ... // 获取cadvisor采集的机器资源容量 machineInfo, err := cadvisorInterface.MachineInfo() ... capacity := cadvisor.CapacityFromMachineInfo(machineInfo) for k, v := range capacity { internalCapacity[k] = v } // 获取 /proc/sys/kernel/pid_max 限制最大进程 pidlimits, err := pidlimit.Stats() if err == nil && pidlimits != nil && pidlimits.MaxPID != nil { internalCapacity[pidlimit.PIDs] = *resource.NewQuantity( int64(*pidlimits.MaxPID), resource.DecimalSI) } cgroupRoot := ParseCgroupfsToCgroupName(nodeConfig.CgroupRoot) // 构造cgroupManager cgroupManager := NewCgroupManager(subsystems, nodeConfig.CgroupDriver) // qos启用 if nodeConfig.CgroupsPerQOS { ... // 校验cgroup root存在 cgroupManager.Validate(cgroupRoot) ... klog.InfoS("Container manager verified user specified cgroup-root exists", "cgroupRoot", cgroupRoot) // 构造qos pod根目录(kubepods) cgroupRoot = NewCgroupName(cgroupRoot, defaultNodeAllocatableCgroupName) } // 初始化qos管理器 qosContainerManager, err := NewQOSContainerManager(subsystems, cgroupRoot, nodeConfig, cgroupManager) // 初始化containerManager cm := &containerManagerImpl{ cadvisorInterface: cadvisorInterface, mountUtil: mountUtil, NodeConfig: nodeConfig, subsystems: subsystems, cgroupManager: cgroupManager, capacity: capacity, internalCapacity: internalCapacity, cgroupRoot: cgroupRoot, recorder: recorder, qosContainerManager: qosContainerManager, } // 启动拓扑 if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.TopologyManager) { // 初始化拓扑管理器 cm.topologyManager, err = topologymanager.NewManager( machineInfo.Topology, nodeConfig.ExperimentalTopologyManagerPolicy, nodeConfig.ExperimentalTopologyManagerScope, ) ... } else { cm.topologyManager = topologymanager.NewFakeManager() } // 启用设备插件 if devicePluginEnabled { // 初始化设备插件管理器 cm.deviceManager, err = devicemanager.NewManagerImpl(machineInfo.Topology, cm.topologyManager) ... // 注册HintProvider cm.topologyManager.AddHintProvider(cm.deviceManager) } else { cm.deviceManager, err = devicemanager.NewManagerStub() } ... // 启用cpuManager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.CPUManager) { cm.cpuManager, err = cpumanager.NewManager( nodeConfig.ExperimentalCPUManagerPolicy, nodeConfig.ExperimentalCPUManagerPolicyOptions, nodeConfig.ExperimentalCPUManagerReconcilePeriod, machineInfo, nodeConfig.NodeAllocatableConfig.ReservedSystemCPUs, cm.GetNodeAllocatableReservation(), nodeConfig.KubeletRootDir, cm.topologyManager, ) ... // 注册HintProvider cm.topologyManager.AddHintProvider(cm.cpuManager) } // 启用memoryManager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.MemoryManager) { cm.memoryManager, err = memorymanager.NewManager( nodeConfig.ExperimentalMemoryManagerPolicy, machineInfo, cm.GetNodeAllocatableReservation(), nodeConfig.ExperimentalMemoryManagerReservedMemory, nodeConfig.KubeletRootDir, cm.topologyManager, ) ... // 注册HintProvider cm.topologyManager.AddHintProvider(cm.memoryManager) } return cm, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
# 10.3.启用
由于依赖
cadvisor提供的节点文件系统容量信息,containerManager要在cadvisor后启动,会依次启动各管理器进行资源监听及分配。// start containerManager for resource scheduler func (cm *containerManagerImpl) Start(...) error { // 启动cpuManager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.CPUManager) { // 调用运行时构建containerMap缓存 containerMap := buildContainerMapFromRuntime(runtimeService) // 启动containerManager cm.cpuManager.Start(cpumanager.ActivePodsFunc(activePods), sourcesReady, podStatusProvider, runtimeService, containerMap) ... } // 启用memoryManager if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.MemoryManager) { // 调用运行时构建containerMap缓存 containerMap := buildContainerMapFromRuntime(runtimeService) // 启动memoryManager cm.memoryManager.Start(memorymanager.ActivePodsFunc(activePods), sourcesReady, podStatusProvider, runtimeService, containerMap) ... } // 缓存节点 cm.nodeInfo = node // 开启临时存储 if utilfeature.DefaultFeatureGate.Enabled(kubefeatures.LocalStorageCapacityIsolation) { // 获取rootfsInfo(/run/containerd/io.containerd.runtime.v2.task/k8s.io) rootfs, err := cm.cadvisorInterface.RootFsInfo() ... // 限制本地临时存储(避免根分区打满造成节点不稳定) for rName, rCap := range cadvisor.EphemeralStorageCapacityFromFsInfo(rootfs) { cm.capacity[rName] = rCap } } // 校验NodeAllocatable配置合法性 cm.validateNodeAllocatable() ... // 设置节点资源容器结构 cm.setupNode(activePods) ... // systemContainer开启ensureStateFunc if hasEnsureStateFuncs { // 间隔1min周期检查 go wait.Until(func() { for _, cont := range cm.systemContainers { if cont.ensureStateFunc != nil { // 检查系统容器cgroup/oom cont.ensureStateFunc(cont.manager) ... } } }, time.Minute, wait.NeverStop) } // 间隔5min执行周期任务 if len(cm.periodicTasks) > 0 { go wait.Until(func() { for _, task := range cm.periodicTasks { if task != nil { task() } } }, 5*time.Minute, wait.NeverStop) } // 启用设备管理器(调度设备及接受插件注册) cm.deviceManager.Start(devicemanager.ActivePodsFunc(activePods), sourcesReady) ... return nil } // 设置节点资源容器结构 func (cm *containerManagerImpl) setupNode(activePods ActivePodsFunc) error { // 检查当前环境支持cgroup subsystem f, err := validateSystemRequirements(cm.mountUtil) ... // 不支持CPU限制,状态中记录软错误 if !f.cpuHardcapping { cm.status.SoftRequirements = fmt.Errorf("CPU hardcapping unsupported") } // 设置内核参数 b := KernelTunableModify // 开启内核参数保护(--protect-kernel-defaults=true),禁止修改 if cm.GetNodeConfig().ProtectKernelDefaults { b = KernelTunableError } // 设置内核参数 setupKernelTunables(b) ... // 初始化qos cgroup层级 if cm.NodeConfig.CgroupsPerQOS { // 初始化qos根目录 cm.createNodeAllocatableCgroups() ... // 周期更新qos cgroup cm.qosContainerManager.Start(cm.GetNodeAllocatableAbsolute, activePods) ... } // 限制节点不同cgroup根可用资源 cm.enforceNodeAllocatableCgroups() ... systemContainers := []*systemContainer{} // system cgroup管理 if cm.SystemCgroupsName != "" { // 必须指定system cgroupName if cm.SystemCgroupsName == "/" { return fmt.Errorf("system container cannot be root (\"/\")") } // 初始化systemContainer,管理system cgroup相关 cont, err := newSystemCgroups(cm.SystemCgroupsName) ... // 注册状态回调 cont.ensureStateFunc = func(manager cgroups.Manager) error { return ensureSystemCgroups("/", manager) } // 记录为systemContainer systemContainers = append(systemContainers, cont) } // kubelet cgroup管理 if cm.KubeletCgroupsName != "" { // 初始化kubelet cgroup管理器 cont, err := newSystemCgroups(cm.KubeletCgroupsName) ... // 注册状态回调 cont.ensureStateFunc = func(_ cgroups.Manager) error { return ensureProcessInContainerWithOOMScore(os.Getpid(), int(cm.KubeletOOMScoreAdj), cont.manager) } // 记录为systemContainer systemContainers = append(systemContainers, cont) } else { // 注册周期任务 cm.periodicTasks = append(cm.periodicTasks, func() { // kubelet进程加入cgroup,设置oomScore ensureProcessInContainerWithOOMScore(os.Getpid(), int(cm.KubeletOOMScoreAdj), nil) ... // 获取kubelet cgroup路径 cont, err := getContainer(os.Getpid()) ... cm.Lock() defer cm.Unlock() // 记录kubelet cgroup路径 cm.KubeletCgroupsName = cont }) } cm.systemContainers = systemContainers return nil } // 强制cgroup资源限制 func (cm *containerManagerImpl) enforceNodeAllocatableCgroups() error { // 获取节点可分配资源 nc := cm.NodeConfig.NodeAllocatableConfig // 节点可分配资源默认为全部资源 nodeAllocatable := cm.internalCapacity // 开启--enforce-node-allocatable强制分配策略 if cm.CgroupsPerQOS && nc.EnforceNodeAllocatable.Has(kubetypes.NodeAllocatableEnforcementKey) { // 设置pods可调度资源为去除节点预留 nodeAllocatable = cm.getNodeAllocatableInternalAbsolute() } ... if len(cm.cgroupRoot) > 0 { go func() { for { // 更新kubepods cgroup资源限制 cm.cgroupManager.Update(cgroupConfig) ... } }() } // 设置system cgroup限制 if nc.EnforceNodeAllocatable.Has(kubetypes.SystemReservedEnforcementKey) { ... // system cgroup检查及资源限制 enforceExistingCgroup(cm.cgroupManager, cm.cgroupManager.CgroupName(nc.SystemReservedCgroupName), nc.SystemReserved) ... } // 设置kube cgroup限制 if nc.EnforceNodeAllocatable.Has(kubetypes.KubeReservedEnforcementKey) { // 检查kube cgroup及资源限制 enforceExistingCgroup(cm.cgroupManager, cm.cgroupManager.CgroupName(nc.KubeReservedCgroupName), nc.KubeReserved) ... } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
