【NPU】Ascend Device Plugin —Part 6 K8s客户端 + 重复检测模块 超深度源码分析之一
Part 6: K8s客户端 + 重复检测模块 超深度源码分析
分析范围:
pkg/kubeclient/(9个文件) +pkg/duplicatedetector/(5个文件)
总代码行数: ~2,473 行
模块定位: Kubernetes API 交互层 + 容器运行时 NPU 设备重复挂载检测层
一、模块定位
1.1 业务职责与功能定位
本模块由两个子模块构成,分别承担不同的业务职责:
| 子模块 | 业务职责 | 功能定位 |
|---|---|---|
kubeclient |
Kubernetes API 客户端封装 | 为整个 ascended-device-plugin 提供统一的 K8s 资源操作入口,包括 Node/Pod/ConfigMap/Event 的 CRUD 操作、Pod 缓存管理、Informer 监听、Kubelet 直连查询等 |
duplicatedetector |
NPU 设备重复挂载检测 | 运行时检测同一物理 NPU 设备被多个容器同时挂载的异常场景,支持 Docker 和 containerd 两种容器运行时 |
1.2 在系统中的位置
核心定位说明:
kubeclient是所有上层业务 Manager 与 K8s 交互的唯一通道,封装了 API Server 调用、Kubelet 直连、Informer 缓存三种数据获取方式duplicatedetector作为安全防护组件,在设备插件启动时初始化,持续监控容器级别的 NPU 设备挂载情况,防止设备重复分配
二、模块整体结构
2.1 kubeclient 子模块结构
2.2 duplicatedetector 子模块结构
2.3 核心方法清单
kubeclient 核心方法
| 文件 | 方法 | 作用 |
|---|---|---|
kubeclient.go |
NewClientK8s() |
创建 K8s 客户端实例,初始化 Clientset、NodeName、Kubelet HTTP 客户端 |
kubeclient.go |
GetNode() |
获取当前节点信息 |
kubeclient.go |
PatchNodeState() |
补丁更新节点状态 |
kubeclient.go |
AddAnnotation() |
为节点添加注解(带重试) |
kubeclient.go |
GetPod() |
获取指定 Pod |
kubeclient.go |
PatchPod() |
补丁更新 Pod 信息 |
kubeclient.go |
GetActivePodList() |
获取当前节点上活跃的 Pod 列表 |
kubeclient.go |
GetAllPodList() |
获取当前节点上所有 Pod 列表 |
kubeclient.go |
GetContainerRuntime() |
获取节点容器运行时类型 |
kubeclient.go |
CreateConfigMap() |
创建 ConfigMap |
kubeclient.go |
GetConfigMap() |
获取 ConfigMap |
kubeclient.go |
UpdateConfigMap() |
更新 ConfigMap |
kubeclient.go |
ResetDeviceInfo() |
重置设备信息 |
kubeclient.go |
ClearResetInfo() |
清除重置信息 |
kubeclient.go |
CreateEvent() |
创建 K8s Event |
kubeclient.go |
ResourceEventHandler() |
创建资源事件处理器 |
kubeclient.go |
RemoveOldResource() |
移除节点上旧的无用资源 |
kube_connect.go |
getPodsByKltPort() |
通过 Kubelet 端口直接获取 Pod 列表 |
kube_connect.go |
createKltPodsReqWithToken() |
创建带 Token 认证的 Kubelet HTTP 请求 |
kube_connect.go |
getKltPodsURL() |
构建 Kubelet /pods URL |
kube_cache.go |
PodInformerInspector() |
定期检查 Pod 缓存有效性 |
kube_cache.go |
UpdatePodList() |
通过 Informer 更新 Pod 缓存 |
kube_cache.go |
GetAllPodListCache() |
从缓存获取所有 Pod 列表 |
kube_cache.go |
GetActivePodListCache() |
从缓存获取活跃 Pod 列表 |
kube_cache.go |
GetNodeIpCache() |
带缓存的节点 IP 获取 |
kube_cache.go |
GetServerUsageLabelCache() |
带缓存的节点标签获取 |
kube_cache.go |
WriteDeviceInfoDataIntoCMCache() |
写设备信息到 CM 并缓存 |
cur_node_informer.go |
InitPodInformer() |
初始化 Pod Informer |
client_server.go |
TryUpdatePodAnnotation() |
更新 Pod 注解(带重试) |
client_server.go |
TryUpdatePodCacheAnnotation() |
同时更新 Pod 注解和缓存 |
client_server.go |
WriteDeviceInfoDataIntoCM() |
写设备信息到 ConfigMap |
client_server.go |
WriteResetInfoDataIntoCM() |
写重置信息到 ConfigMap |
client_server.go |
WriteFaultInfoDataIntoCM() |
写故障信息到 ConfigMap |
client_server.go |
AnnotationReset() |
重置节点注解 |
client_server.go |
GetPodsUsedNpuByCommon() |
从缓存获取已使用 NPU 列表 |
client_server.go |
GetPodsUsedNPUByKlt() |
从 Kubelet 获取已使用 NPU 列表 |
client_server.go |
GetNodeIp() |
获取节点 IP |
kubeclient_v2.go |
ParseFaultNetworkInfoCM() |
解析故障网络配置 ConfigMap |
kubeclient_v2.go |
IsConfigMapExists() |
判断 ConfigMap 是否存在 |
duplicatedetector 核心方法
| 文件 | 方法 | 作用 |
|---|---|---|
manager.go |
CheckDuplicateDevices() |
模块入口:单例初始化并启动检测 |
manager.go |
NewManager() |
创建重复检测管理器 |
manager.go |
Start() |
启动检测:全量扫描 + 事件监听 |
manager.go |
HandleNewContainer() |
处理新容器创建事件 |
manager.go |
HandleContainerRemoval() |
处理容器销毁事件 |
cache/container_cache.go |
StoreAllAndFindDuplicates() |
全量存储并查找重复 |
cache/container_cache.go |
StoreSingleAndFindDuplicates() |
增量存储并查找重复 |
cache/container_cache.go |
RemoveContainer() |
从缓存移除容器 |
containerruntime/interface.go |
NewClient() |
根据运行时类型创建客户端 |
containerruntime/interface.go |
autoDetectOciEndpoint() |
自动探测 OCI socket 路径 |
containerruntime/interface.go |
ParseSingleContainer() |
解析单个容器的 NPU 设备信息 |
docker_client.go |
NewDockerClient() |
创建 Docker 客户端 |
docker_client.go |
ParseAllContainers() |
列出并解析所有 Docker 容器 |
docker_client.go |
WatchContainerEvents() |
监听 Docker 容器事件 |
containerd_client.go |
NewContainerdClient() |
创建 containerd 客户端 |
containerd_client.go |
ParseAllContainers() |
列出并解析所有 containerd 容器 |
containerd_client.go |
WatchContainerEvents() |
监听 containerd 容器事件 |
2.4 内部调用关系
2.5 数据流入流出方式
三、核心业务逻辑深度解析
3.1 kubeclient 模块
3.1.1 客户端初始化流程 — NewClientK8s()
文件: kubeclient.go L72-L106
逐行解析:
// L72-73: 构建K8s配置,两个空字符串参数表示使用InClusterConfig
// 这是K8s推荐的方式:Pod内运行的服务自动使用ServiceAccount的凭据
clientCfg, err := clientcmd.BuildConfigFromFlags("", "")
// L77: 使用配置创建正式的K8s Clientset
// Clientset是K8s client-go提供的类型安全客户端,可操作所有K8s原生资源
client, err := kubernetes.NewForConfig(clientCfg)
// L82: 从环境变量 NODE_NAME 获取当前节点名
// 这是Device Plugin部署时通过K8s downward API注入的
nodeName, err := GetNodeNameFromEnv()
// L84-91: 构建Kubelet HTTP客户端的Transport
// 关键安全配置:
// - InsecureSkipVerify: 跳过证书验证(Kubelet通常使用自签证书)
// - CipherSuites: 限定使用ECDHE系列强加密套件,禁用弱算法
// - MinVersion: 强制TLS 1.3,拒绝旧版本协议
transport := &http.Transport{
TLSClientConfig: &tls.Config{
InsecureSkipVerify: true,
CipherSuites: defaultSafeCipherSuites,
MinVersion: tls.VersionTLS13,
},
}
// L93: 创建Kubelet专用的HTTP客户端,后续用于直连Kubelet /pods接口
kltClient := &http.Client{Transport: transport}
// L95-103: 组装ClientK8s结构体
// - DeviceInfoName: 设备信息ConfigMap的命名规则 = "mindx-dl-deviceinfo-" + nodeName
// - Queue: 速率限制工作队列,用于Informer事件去重和背压控制
return &ClientK8s{
Clientset: client,
NodeName: nodeName,
DeviceInfoName: common.DeviceInfoCMNamePrefix + nodeName,
Queue: workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter()),
IsApiErr: false,
KltClient: kltClient,
}, nil
安全加密套件设计意图:
// L40-47: defaultSafeCipherSuites 限定了6种ECDHE套件
// 设计目的: 在InsecureSkipVerify=true的前提下,仍保证传输层加密强度
var defaultSafeCipherSuites = []uint16{
tls.TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256, // ECDSA + AES-128-GCM
tls.TLS_ECDHE_ECDSA_WITH_AES_256_GCM_SHA384, // ECDSA + AES-256-GCM
tls.TLS_ECDHE_ECDSA_WITH_CHACHA20_POLY1305_SHA256, // ECDSA + ChaCha20
tls.TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256, // RSA + AES-128-GCM
tls.TLS_ECDHE_RSA_WITH_AES_256_GCM_SHA384, // RSA + AES-256-GCM
tls.TLS_ECDHE_RSA_WITH_CHACHA20_POLY1305_SHA256, // RSA + ChaCha20
}
3.1.2 节点名验证 — GetNodeNameFromEnv() + checkNodeName()
文件: kubeclient.go L263-L285
// 从环境变量 NODE_NAME 获取节点名
func GetNodeNameFromEnv() (string, error) {
nodeName := os.Getenv(api.NodeNameEnv) // 读取环境变量 NODE_NAME
if err := checkNodeName(nodeName); err != nil {
return "", fmt.Errorf("check node name failed: %v", err)
}
return nodeName, nil
}
// 节点名合法性校验:长度、正则
func checkNodeName(nodeName string) error {
// 检查1: 非空 — NODE_NAME 环境变量必须设置
if len(nodeName) == 0 {
return fmt.Errorf("the env variable whose key is NODE_NAME must be set")
}
// 检查2: 长度上限 230 字符 (KubeEnvMaxLength)
// 防止超长节点名导致K8s资源名截断或API拒绝
if len(nodeName) > common.KubeEnvMaxLength {
return fmt.Errorf("node name length %d is bigger than %d", len(nodeName), common.KubeEnvMaxLength)
}
// 检查3: 正则匹配 — 确保节点名只包含合法字符
// pattern 来自 common.GetPattern()["nodeName"],符合K8s节点名规范
pattern := common.GetPattern()["nodeName"]
if match := pattern.MatchString(nodeName); !match {
return fmt.Errorf("node name %s is illegal", nodeName)
}
return nil
}
3.1.3 Pod 列表获取 — 三级数据获取策略
本模块设计了三级 Pod 数据获取策略,体现了从"最新数据"到"缓存降级"的容错设计:
方式1详解 — Kubelet直连 getPodsByKltPort():
文件: kube_connect.go L82-L112
func (ki *ClientK8s) getPodsByKltPort() (*v1.PodList, error) {
// 步骤1: 创建带Token认证的HTTP请求
// 内部调用 getKltPodsURL() 构建URL:
// - 从 HOST_IP 环境变量获取主机IP
// - 从 KUBELET_PORT 环境变量获取端口(默认10250)
// - IPv6地址加方括号
// - 路径固定为 /pods
// 然后调用 createKltPodsReqWithToken():
// - 使用 rest.InClusterConfig() 获取ServiceAccount Token
// - 添加 Authorization: Bearer <token> 头
req, err := createKltPodsReqWithToken()
if err != nil {
return nil, err
}
// 步骤2: 发送HTTP请求(使用NewClientK8s中创建的TLS 1.3客户端)
resp, err := ki.KltClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
// 步骤3: 检查HTTP状态码
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("get kubelet http response failed, resp status code: %v is not %v",
resp.StatusCode, http.StatusOK)
}
// 步骤4: 读取响应体(使用LimitReader防止响应过大导致OOM)
// limiter.DefaultDataLimit 限制最大读取字节数
body, err := io.ReadAll(io.LimitReader(resp.Body, limiter.DefaultDataLimit))
if err != nil {
if err == io.EOF {
// 响应体超过限制,可能是异常的大响应
return nil, fmt.Errorf("response size exceeds limit")
}
return nil, err
}
// 步骤5: JSON反序列化为PodList
var podList v1.PodList
if err = json.Unmarshal(body, &podList); err != nil {
return nil, err
}
return &podList, nil
}
Kubelet URL构建详解 — getKltPodsURL():
文件: kube_connect.go L43-L67
func getKltPodsURL() (string, error) {
// 从环境变量获取主机IP
envHostIP := os.Getenv(HostIPEnv) // "HOST_IP"
parseIP := net.ParseIP(envHostIP)
if parseIP == nil {
return "", fmt.Errorf("host ip is invalid") // IP格式非法
}
hostIP := parseIP.String()
// 从环境变量获取Kubelet端口,默认10250
kubeletPort := os.Getenv(KubeletPortEnv) // "KUBELET_PORT"
if err := isValidPort(kubeletPort); err != nil {
// 端口无效时降级为默认端口
kubeletPort = DefaultKubeletPort // "10250"
}
// IPv6地址处理:需要加方括号
// 例如: [::1]:10250 vs 192.168.1.1:10250
if net2.IsIPv6(parseIP) {
hostIP = "[" + hostIP + "]"
}
// 构建完整URL: https://<hostIP>:<kubeletPort>/pods
host := fmt.Sprintf("%s:%s", hostIP, kubeletPort)
kltPodsURL := &url.URL{
Scheme: "https",
Host: host,
Path: "/pods",
}
return kltPodsURL.String(), nil
}
Token认证详解 — createKltPodsReqWithToken():
文件: kube_connect.go L70-L80
func createKltPodsReqWithToken() (*http.Request, error) {
// 1. 构建Kubelet URL
kltPodsURL, err := getKltPodsURL()
if err != nil {
return nil, err
}
// 2. 首次调用时记录URL(使用sync.Once确保只打印一次)
logUrlOnce.Do(func() {
hwlog.RunLog.Info("get kubelet pods url: ", kltPodsURL)
})
// 3. 创建HTTP GET请求
req, err := http.NewRequest("GET", kltPodsURL, nil)
if err != nil {
return nil, err
}
// 4. 首次调用时初始化InClusterConfig获取BearerToken
// 使用sync.Once确保只初始化一次,避免重复创建配置
kubeConfigOnce.Do(func() {
kubeConfig, err = rest.InClusterConfig()
// InClusterConfig自动读取:
// - Token: /var/run/secrets/kubernetes.io/serviceaccount/token
// - CA证书: /var/run/secrets/kubernetes.io/serviceaccount/ca.crt
})
if kubeConfig == nil {
return nil, fmt.Errorf("kubeConfig is nil")
}
// 5. 添加Bearer Token认证头
// Kubelet验证Token以确认调用者是集群内的合法Pod
req.Header.Add("Authorization", "Bearer "+kubeConfig.BearerToken)
return req, nil
}
方式2详解 — API Server缓存读取 GetActivePodListCache():
文件: kube_cache.go L193-L213
func (ki *ClientK8s) GetActivePodListCache() []v1.Pod {
// 关键: 检查IsApiErr标志
// 当API Server调用出错时,IsApiErr被置为true
// 下次读取缓存时会触发全量刷新,保证缓存数据不会过时太久
if ki.IsApiErr {
ki.refreshPodList() // 重新从API Server全量拉取
}
newPodList := make([]v1.Pod, 0, common.GeneralMapSize)
lock.Lock()
defer lock.Unlock()
// 遍历内存缓存 podCache (map[types.UID]*podInfo)
for _, pi := range podCache {
// 安全检查1: Pod名称合法性
if err := common.CheckPodNameAndSpace(pi.GetName(), common.PodNameMaxLength); err != nil {
continue // 跳过非法名称的Pod
}
// 安全检查2: Pod命名空间合法性
if err := common.CheckPodNameAndSpace(pi.GetNamespace(), common.PodNameSpaceMaxLength); err != nil {
continue
}
// 过滤条件: 排除已结束的Pod(Failed和Succeeded状态)
if pi.Status.Phase == v1.PodFailed || pi.Status.Phase == v1.PodSucceeded {
continue
}
newPodList = append(newPodList, *pi.Pod)
}
return newPodList
}
缓存刷新机制 — refreshPodList():
文件: kube_cache.go L155-L169
func (ki *ClientK8s) refreshPodList() {
// 从API Server全量拉取当前节点的Pod列表
newV1PodList, err := ki.GetAllPodList()
if err != nil {
// 拉取失败时保持旧缓存不变,避免错误数据覆盖
return
}
// 构建全新的缓存Map
newPodCache := map[types.UID]*podInfo{}
for _, pod := range newV1PodList.Items {
// 关键: Go的for range遍历slice时,循环变量地址固定
// 直接取 &pod 会导致所有缓存项指向同一个地址
// 因此使用闭包参数传递,确保每个pod有独立地址
func(pod v1.Pod) {
newPodCache[pod.UID] = &podInfo{
Pod: &pod,
updateTime: time.Now(),
}
}(pod)
}
lock.Lock()
podCache = newPodCache // 原子替换整个缓存
lock.Unlock()
ki.IsApiErr = false // 重置错误标志
}
3.1.4 Pod Informer 初始化与事件处理
文件: cur_node_informer.go L33-L56
func (ki *ClientK8s) InitPodInformer() {
// 创建SharedInformerFactory,关键配置:
// 1. 使用ki.Clientset作为K8s客户端
// 2. resyncPeriod=0 表示不做定期全量同步,依赖Watch事件驱动
// 3. WithTweakListOptions: 设置FieldSelector只关注当前节点的Pod
// 这大幅减少了API Server返回的数据量和网络开销
factory := informers.NewSharedInformerFactoryWithOptions(ki.Clientset, 0,
informers.WithTweakListOptions(func(options *v1.ListOptions) {
options.FieldSelector = "spec.nodeName=" + ki.NodeName
}))
// 创建Pod Informer
podInformer := factory.Core().V1().Pods().Informer()
// 注册第一个EventHandler: 维护podCache
// 所有Pod变更都会触发UpdatePodList,更新内存缓存
podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ki.UpdatePodList(obj, EventTypeAdd)
},
UpdateFunc: func(oldObj, newObj interface{}) {
// DeepEqual优化: 只有Pod对象实际变更时才处理
if !reflect.DeepEqual(oldObj, newObj) {
ki.UpdatePodList(newObj, EventTypeUpdate)
}
},
DeleteFunc: func(obj interface{}) {
ki.UpdatePodList(obj, EventTypeDelete)
},
})
// 注册第二个EventHandler: 过滤后入队
// 使用FilteringResourceEventHandler,过滤条件为checkPod
// checkPod: 只关注带有 HuaweiAscend910 注解的Pod
// 这些是使用了Ascend NPU的Pod,需要特殊处理
podInformer.AddEventHandler(ki.ResourceEventHandler(PodResource, checkPod))
// 设置Watch错误处理器
// 当Informer与API Server的Watch连接断开时触发
podInformer.SetWatchErrorHandler(func(r *cache.Reflector, err error) {
hwlog.RunLog.Errorf("pod informer watch error: %v", err)
// 如果配置了DealWatchHandler,则标记缓存需要刷新
// 下次读取缓存时会从API Server全量拉取
if common.ParamOption.DealWatchHandler {
ki.FlushPodCacheNextQuerying() // 设置IsApiErr=true
}
})
// 启动Informer Factory
factory.Start(make(chan struct{}))
// 保存PodInformer引用到ClientK8s
ki.PodInformer = podInformer
}
checkPod过滤函数:
// cur_node_informer.go L58-L64
// 过滤条件: Pod必须包含 HuaweiAscend910 注解
// 这意味着只有使用了Ascend 910 NPU的Pod才会被加入WorkQueue
// 其他普通Pod不会触发队列事件,减少不必要处理
func checkPod(obj interface{}) bool {
pod, ok := obj.(*corev1.Pod)
if !ok {
return false
}
_, exist := pod.Annotations[api.HuaweiAscend910]
return exist
}
3.1.5 Pod缓存更新与Predicate Time校正机制
文件: kube_cache.go L83-L130
updatePodPredicateTime() 详解 — 这是调度纠正机制的核心:
// kube_cache.go L132-L153
func (ki *ClientK8s) updatePodPredicateTime(newPod *v1.Pod) {
// 从缓存中查找该Pod的旧信息
cachedPodInfo, exists := podCache[newPod.UID]
if !exists {
return // 缓存中不存在,无需校正
}
// 检查缓存中Pod的注解是否存在
if cachedPodInfo.Pod.Annotations == nil {
return
}
// 读取缓存中的 predicate-time 值
// predicate-time 是调度阶段的标记,值为 MaxUint64 表示尚未被分配
cachedPredicateTime := cachedPodInfo.Pod.Annotations[common.PodPredicateTime]
if cachedPredicateTime != strconv.FormatUint(math.MaxUint64, common.BaseDec) {
return // 只有 MaxUint64 才需要校正
}
// 检查新Pod的注解是否与缓存一致
needsUpdate := false
if newPod.Annotations == nil {
// 新Pod没有注解,需要补充
newPod.Annotations = make(map[string]string)
needsUpdate = true
}
if newPod.Annotations[common.PodPredicateTime] != cachedPredicateTime {
// 新Pod的predicate-time与缓存不一致,需要校正
// 这种情况可能发生在: API Server返回的Pod数据丢失了annotation
needsUpdate = true
}
if needsUpdate {
// 将缓存中的正确值写回新Pod对象
newPod.Annotations[common.PodPredicateTime] = cachedPredicateTime
// 异步更新API Server上的Pod annotation
// 使用goroutine避免阻塞Informer事件处理循环
annotation := map[string]string{common.PodPredicateTime: strconv.FormatUint(math.MaxUint64, common.BaseDec)}
go ki.TryUpdatePodCacheAnnotation(newPod, annotation)
}
}
设计意图分析:predicate-time 是华为自研的调度标记,当值为 MaxUint64 时表示 Pod 已通过调度但尚未完成设备分配。如果 API Server 返回的数据丢失了这个注解,可能导致重复分配。此机制确保缓存中的正确值被同步回 API Server。
3.1.6 Pod缓存定期巡检 — PodInformerInspector()
FNV32(nodeName) % 600] -----------------------^ Expecting 'SQE', 'DOUBLECIRCLEEND', 'PE', '-)', 'STADIUMEND', 'SUBROUTINEEND', 'PIPE', 'CYLINDEREND', 'DIAMOND_STOP', 'TAGEND', 'TRAPEND', 'INVTRAPEND', 'UNICODE_TEXT', 'TEXT', 'TAGSTART', got 'PS'
文件: kube_cache.go L40-L75
func (ki *ClientK8s) PodInformerInspector(ctx context.Context) {
// 使用FNV哈希将节点名映射到0-599的数字
// 目的: 让不同节点的Device Plugin在不同的时间点开始巡检
// 避免所有节点同时向API Server发起大量查询(避免惊群效应)
hashVal := fnv.New32()
if _, err := hashVal.Write([]byte(ki.NodeName)); err != nil {
return
}
val := hashVal.Sum32() % periodicForStartCheck // periodicForStartCheck = 600
// 随机延迟启动
time.Sleep(time.Duration(val) * time.Second)
// 每10分钟执行一次巡检
// wait.Until会在ctx.Done()时自动停止
wait.Until(func() { ki.checkPodInCache(ctx) }, timeIntervalForCheckPod, ctx.Done())
// timeIntervalForCheckPod = 10 * time.Minute
}
func (ki *ClientK8s) checkPodInCache(ctx context.Context) {
lock.Lock()
defer lock.Unlock()
needDelete := make([]types.UID, 0) // 待删除的Pod UID列表
needRefresh := make([]types.UID, 0) // 待刷新的Pod UID列表
for uid, pi := range podCache {
// 检查1: 更新时间距今是否小于1小时
// podCacheTimeout = time.Hour
// 1小时内的Pod视为有效,跳过检查
if time.Since(pi.updateTime) < podCacheTimeout {
continue
}
// 超过1小时未更新的Pod,需要从API Server验证
pod, err := ki.getPod(ctx, pi.Namespace, pi.Name)
if err != nil {
if errors.IsNotFound(err) {
// Pod已被删除,从缓存中清除
needDelete = append(needDelete, uid)
continue
}
// 其他错误(如网络超时)跳过,下次巡检再检查
continue
}
// 检查2: Pod是否仍然调度在当前节点
// 可能Pod被重新调度到了其他节点
if pod.Spec.NodeName != ki.NodeName || pod.UID != uid {
needDelete = append(needDelete, uid)
continue
}
// Pod仍然有效,刷新更新时间
needRefresh = append(needRefresh, uid)
}
// 批量执行删除
for _, uid := range needDelete {
delete(podCache, uid)
}
// 批量执行刷新
for _, uid := range needRefresh {
podCache[uid].updateTime = time.Now()
}
}
3.1.7 ConfigMap 写入机制
3.1.7.1 设备信息写入 WriteDeviceInfoDataIntoCM()
文件: client_server.go L156-L207
逐行解析:
func (ki *ClientK8s) WriteDeviceInfoDataIntoCM(nodeDeviceData *common.NodeDeviceInfoCache,
manuallySeparateNPU string, switchInfo common.SwitchFaultInfo, dpuInfo common.DpuInfo,
reasonCm string) (*common.NodeDeviceInfoCache, error) {
// 步骤1: 计算数据校验码
// MakeDataHash 对 DeviceInfo 内容做哈希,用于数据完整性校验
// 下游消费者可以通过比较CheckCode判断数据是否变更
nodeDeviceData.CheckCode = common.MakeDataHash(nodeDeviceData.DeviceInfo)
// 步骤2: 序列化各数据块
var data, switchData, dpuData []byte
// 判断DPU功能是否开启:通过DeepEqual比较dpuInfo与零值
dpuOpen := !reflect.DeepEqual(dpuInfo, common.DpuInfo{})
if data = common.MarshalData(nodeDeviceData); len(data) == 0 {
return nil, fmt.Errorf("marshal nodeDeviceData failed")
}
if switchData = common.MarshalData(switchInfo); len(switchData) == 0 {
return nil, fmt.Errorf("marshal switchDeviceData failed")
}
// 步骤3: 构建ConfigMap对象
deviceInfoCM := &v1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Name: ki.DeviceInfoName, // "mindx-dl-deviceinfo-" + nodeName
Namespace: api.KubeNS, // 华为设备插件专用namespace
Labels: map[string]string{api.CIMCMLabelKey: common.CmConsumerValue},
},
}
// 步骤4: 按硬件型号设置不同数据结构
switch common.ParamOption.RealCardType {
case api.Ascend910A5:
// A5型号: 完整数据,包含交换机故障信息和可选的DPU信息
deviceInfoCM.Data = map[string]string{
api.DeviceInfoCMDataKey: string(data),
api.SwitchInfoCMDataKey: string(switchData),
common.DeviceInfoCmUpgradeFaultReasonKey: reasonCm,
common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
common.DescriptionKey: common.DescriptionValue,
}
if dpuOpen {
if dpuData = common.MarshalData(dpuInfo); len(dpuData) == 0 {
return nil, fmt.Errorf("marshal DpuDeviceData failed")
}
deviceInfoCM.Data[api.DpuInfoCMDataKey] = string(dpuData)
}
case api.Ascend910A3:
// A3型号: 有交换机信息但无DPU
deviceInfoCM.Data = map[string]string{
api.DeviceInfoCMDataKey: string(data),
api.SwitchInfoCMDataKey: string(switchData),
common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
common.DeviceInfoCmUpgradeFaultReasonKey: reasonCm,
common.DescriptionKey: common.DescriptionValue,
}
default:
// 其他型号: 只有基础设备信息,无交换机/DPU
deviceInfoCM.Data = map[string]string{
api.DeviceInfoCMDataKey: string(data),
common.DeviceInfoCMManuallySeparateNPUKey: manuallySeparateNPU,
common.DeviceInfoCmUpgradeFaultReasonKey: reasonCm,
common.DescriptionKey: common.DescriptionValue,
}
}
// 步骤5: 写入ConfigMap(先Update,不存在则Create)
if err := ki.createOrUpdateDeviceCM(deviceInfoCM); err != nil {
return nil, err
}
return nodeDeviceData, nil
}
3.1.7.2 重置信息写入 WriteResetInfoDataIntoCM()
文件: client_server.go L209-L260
关键逻辑解析:
// 检查隔离错误 + 新设备列表非空的情况
// 如果旧CM中已包含 IsolateError 标记,且新的任务信息中有设备
// 则拒绝更新,返回"task should be rescheduled"
// 设计意图: 隔离错误状态下不应该继续添加设备,需要重新调度
if strings.Contains(oldResetInfoData, common.IsolateError) && len(taskInfo.RankList) != 0 {
return nil, fmt.Errorf("task should be rescheduled")
}
// setNewTaskInfoWithHexString: 将错误码转换为十六进制字符串
// 设计意图: 上游消费者可能以字符串形式处理错误码,
// 十六进制表示更紧凑且便于人工阅读
func setNewTaskInfoWithHexString(taskInfo *common.TaskResetInfo) *common.TaskResetInfo {
var newTaskInfo common.TaskResetInfo
for _, deviceInfo := range taskInfo.RankList {
newDeviceInfo := *deviceInfo
// 转换为十六进制大写
newDeviceInfo.ErrorCodeHex = strings.ToUpper(common.Int64Tool.ToHexString(newDeviceInfo.ErrorCode))
newDeviceInfo.ErrorCode = []int64{} // 清空原始错误码,节省存储空间
newTaskInfo.RankList = append(newTaskInfo.RankList, &newDeviceInfo)
}
if newTaskInfo.RankList == nil {
newTaskInfo.RankList = make([]*common.TaskDevInfo, 0)
}
return &newTaskInfo
}
3.1.8 NPU使用情况查询
本模块提供两条路径查询节点上NPU的使用情况,形成主备容错:
文件: client_server.go L318-L372
// 主路径: 通过Kubelet直接获取(数据更新,无API Server延迟)
func (ki *ClientK8s) GetPodsUsedNPUByKlt() sets.String {
podList, err := ki.getPodsByKltPort()
if err != nil {
// Kubelet连接失败时,使用带限流的日志
// domainForKubeletConnectErr = "kubeletConnect"
// 防止Kubelet持续不可用时日志刷屏
hwlog.RunLog.ErrorfWithLimit(domainForKubeletConnectErr, os.Getenv(KubeletPortEnv),
"get pods used NPU failed: %v", err)
// 降级到缓存路径
return ki.GetPodsUsedNpuByCommon()
}
usedNPU := make([]string, 0)
for _, pod := range podList.Items {
// 安全检查: Pod名称和命名空间合法性
if err := common.CheckPodNameAndSpace(pod.GetName(), common.PodNameMaxLength); err != nil {
continue
}
if err := common.CheckPodNameAndSpace(pod.GetNamespace(), common.PodNameSpaceMaxLength); err != nil {
continue
}
// 过滤已结束的Pod
if pod.Status.Phase == v1.PodFailed || pod.Status.Phase == v1.PodSucceeded {
continue
}
// 读取 PodAnnotationAscendReal 注解
// 该注解格式: "Ascend910-0,Ascend910-1" 表示使用了设备0和1
realAllocTag := fmt.Sprintf("%s", api.PodAnnotationAscendReal)
tmpNPU, ok := pod.Annotations[realAllocTag]
if !ok || len(tmpNPU) == 0 || len(tmpNPU) > common.PodAnnotationMaxLength {
continue
}
// 分割设备列表
tmpNPUList := strings.Split(tmpNPU, common.CommaSepDev)
if len(tmpNPUList) == 0 || len(tmpNPUList) > common.MaxDevicesNum {
continue
}
usedNPU = append(usedNPU, tmpNPUList...)
}
// 使用sets.String去重
return sets.NewString(usedNPU...)
}
3.1.9 资源事件处理机制 — ResourceEventHandler()
文件: kubeclient.go L287-L312
func (ki *ClientK8s) ResourceEventHandler(res ResourceType, filter func(obj interface{}) bool) cache.
ResourceEventHandler {
// enqueue: 内部闭包,将事件加入WorkQueue
enqueue := func(obj interface{}, event EventType) {
// 关键优化: Pod的Update事件不入队
// 因为Pod更新通过第一个EventHandler(UpdatePodList)已处理
// 这里只需要Add和Delete事件入队
if res == PodResource && event == EventTypeUpdate {
return
}
// 生成对象的唯一Key: namespace/name
key, err := cache.MetaNamespaceKeyFunc(obj)
if err != nil {
return
}
// 将事件加入速率限制队列
// RateLimitingQueue提供:
// 1. 去重: 相同Key的事件会被合并
// 2. 限流: 失败重试有指数退避
ki.Queue.AddRateLimited(Event{
Resource: res,
Key: key,
Type: event,
})
}
// 返回FilteringResourceEventHandler
// 先通过filter函数过滤,再调用对应的处理函数
return cache.FilteringResourceEventHandler{
FilterFunc: filter, // 过滤函数: checkPod(检查HuaweiAscend910注解)
Handler: cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
enqueue(obj, EventTypeAdd)
},
UpdateFunc: func(oldObj, newObj interface{}) {
// DeepEqual优化: 只有对象实际变更才处理
if reflect.DeepEqual(oldObj, newObj) {
return
}
enqueue(newObj, EventTypeUpdate)
},
DeleteFunc: func(obj interface{}) {
enqueue(obj, EventTypeDelete)
},
},
}
}
3.1.10 节点资源清理 — RemoveOldResource()
文件: kubeclient.go L420-L452
func (c *ClientK8s) RemoveOldResource(oldResourceName string) error {
// 场景: 设备型号切换时,需要清理旧的资源名称
// 例如: 从 Ascend910-0 切换到 Ascend910A5-0,旧的资源名需要删除
node, err := c.Clientset.CoreV1().Nodes().Get(
context.TODO(),
c.NodeName,
metav1.GetOptions{ResourceVersion: "0"},
)
if err != nil {
return fmt.Errorf("get node failed: %v", err)
}
// 检查旧资源是否存在于节点容量中
if _, exists := node.Status.Capacity[v1.ResourceName(oldResourceName)]; !exists {
// 不存在则跳过,幂等操作
return nil
}
// 从Capacity和Allocatable中删除旧资源
delete(node.Status.Capacity, v1.ResourceName(oldResourceName))
delete(node.Status.Allocatable, v1.ResourceName(oldResourceName))
// 更新节点状态到API Server
updatedNode, err := c.Clientset.CoreV1().Nodes().UpdateStatus(
context.TODO(),
node,
metav1.UpdateOptions{},
)
if err != nil {
return fmt.Errorf("update node status failed: %v", err)
}
// 验证删除是否成功(双重确认)
if _, exists := updatedNode.Status.Capacity[v1.ResourceName(oldResourceName)]; exists {
return errors.New("failed to remove resource from capacity")
}
if _, exists := updatedNode.Status.Allocatable[v1.ResourceName(oldResourceName)]; exists {
return errors.New("failed to remove resource from allocatable")
}
return nil
}
3.1.11 容器运行时检测 — GetContainerRuntime()
文件: kubeclient.go L190-L207
func (ki *ClientK8s) GetContainerRuntime() (string, error) {
// 通过节点状态信息获取容器运行时
node, err := ki.GetNode()
if err != nil {
return "", fmt.Errorf("failed to get node: %w", err)
}
// 读取节点状态中的容器运行时版本字符串
// 格式: "docker://1.13.1" 或 "containerd://1.6.0"
runtimeVersion := node.Status.NodeInfo.ContainerRuntimeVersion
if runtimeVersion == "" {
return "", fmt.Errorf("container runtime version not found in node status")
}
// 通过URL Scheme前缀判断运行时类型
if strings.HasPrefix(runtimeVersion, "docker://") {
return DockerRuntime, nil // "docker"
}
if strings.HasPrefix(runtimeVersion, "containerd://") {
return ContainerdRuntime, nil // "containerd"
}
return "", fmt.Errorf("unknown container runtime: %s", runtimeVersion)
}
设计意图:此方法为 duplicatedetector 模块提供运行时类型判断,决定创建 Docker 客户端还是 containerd 客户端。
3.1.12 IsApiErr 错误标记机制
机制说明:
- 任何 K8s API 调用返回包含 “443”(ApiServerPort)的错误时,
IsApiErr被置为true - 443 是 K8s API Server 的默认端口,错误中包含此端口表示 API Server 不可达
- 当
IsApiErr=true时,下次读取缓存会触发全量刷新 - 刷新成功后重置为
false,形成自恢复闭环
鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。
更多推荐



所有评论(0)