Ascend Device Plugin — gRPC 服务层 超深度源码分析

分析目录:pkg/server/
涵盖文件:manager.goplugin.goserver.gotypes.gopod_resource.godpu.gomanager_v2.goplugin_v2.gonpu_base_v2.gotypes_v2.go


一、模块定位

1.1 业务职责

pkg/server 是 Ascend Device Plugin 的 gRPC 服务核心层,承担以下关键职责:

职责 说明
Kubernetes 设备插件 gRPC 服务 实现 Kubernetes Device Plugin v1beta1 接口(ListAndWatchAllocateGetDevicePluginOptionsPreStartContainer),向 kubelet 注册并上报设备状态
设备生命周期管理 管理华为昇腾 NPU(910A/910B/910A3/910A5/310/310P/310B)的发现、分类、健康监测、故障处理、热复位
Pod 资源查询 通过 kubelet Pod Resources API 获取 Pod 已分配的设备信息
节点标签与注解管理 向 Kubernetes Node 写入标签(label)和注解(annotation),包括芯片名称、服务器类型、拓扑信息、SuperPod/Rack ID 等
Volcano 调度集成 支持 Volcano 调度器场景下的设备分配映射(klt→real device 映射)
DPU 网卡健康监测 针对 A5(Ascend910A5)设备,周期性查询 DPU 网卡 operstate 并更新到 ConfigMap
故障订阅与热复位 订阅 NPU/交换机故障事件,执行芯片热复位(hot reset),支持训练/推理场景的容错
Rank Table 拓扑生成 为 A5 设备生成 HCCL Rank Table 所需的 LevelList 拓扑信息(URMA EID、RoCE 地址等)
vNPU 动态切分 支持 vNPU 的动态创建与销毁(仅 PresetVDevice=false 模式)
软共享设备 支持 soft-share device 场景的配额管理和配置文件写入

1.2 在系统中的位置

硬件层

Ascend Device Plugin 进程

Kubernetes Node (kubelet)

Kubernetes Master

gRPC Register

ListAndWatch

Allocate

Pod Resources API

DCMI

ethtool

Switch Fault

Patch Node / Annotation

Annotation

Pod Informer

kube-apiserver

Volcano Scheduler

kubelet

pod-resources.sock

device-plugin.sock

pkg/server (本模块)

pkg/common

pkg/device

ascend-common/devmanager (DCMI)

pkg/kubeclient

昇腾 NPU 芯片

DPU 网卡

训练组网交换机

1.3 文件职责总览

pkg/server 文件职责

manager.go
设备管理器核心
(1792行)

plugin.go
gRPC DevicePlugin 接口实现
(1326行)

pod_resource.go
Pod资源查询客户端
(257行)

dpu.go
DPU网卡健康监测
(180行)

manager_v2.go
A5扩展管理功能
(205行)

server.go
gRPC Server 启停与注册
(196行)

plugin_v2.go
HCCL拓扑环境变量
(65行)

npu_base_v2.go
NPU基础拓扑信息
(607行)

types_v2.go
V2类型定义
(44行)

types.go
类型定义
(77行)


二、模块整体结构

2.1 核心类型与接口定义

ServerMap

uses

uses via manager_v2

«interface»

InterfaceServer

+Start(*FileWatch) : error

+Stop()

+GetRestartFlag() : bool

+SetRestartFlag(bool)

+LastSendSuccess() : bool

PluginServer

-manager: DevManager

-grpcServer: *grpc.Server

-isRunning: *AtomicBool

-cachedDevices: []NpuDevice

-deviceType: string

-ascendRuntimeOptions: string

-defaultDevs: []string

-allocMapLock: RWMutex

-cachedLock: RWMutex

-reciChan: chan interface

-stop: chan interface

-klt2RealDevMap: map<string,string>

-restart: bool

-deviceSyncStat: *SendStats

-restartTimes: atomic.Uint64

-podLock: Mutex

+ListAndWatch() : error

+Allocate()(*AllocateResponse, error)

+Notify([]*NpuDevice) : bool

+Start(*FileWatch) : error

+Stop()

+GetRestartFlag() : bool

+SetRestartFlag(bool)

+LastSendSuccess() : bool

+GetDevicePluginOptions()

+PreStartContainer()

+GetPreferredAllocation()

HwDevManager

+SwitchDevManager: *SwitchDevManager

-groupDevice: map<string,[]*NpuDevice>

+ServerMap: map<string,InterfaceServer>

-allInfo: NpuAllInfo

-manager: DevManager

+RunMode: string

+WorkMode: string

-baseNPUInfo: map<string,*NpuBaseInfo>

-dpuManager: *DpuFilter

+ManagerLock: Mutex

+ContainerRuntime: string

+ListenDevice(ctx)

+Serve(ctx)

+SignCatch(cancel)

+UpdateNode() : error

+GetDevManager() : DevManager

PodResource

-conn: *grpc.ClientConn

-client: PodResourcesListerClient

+GetPodResource()(map<string,PodDevice>, error)

+IsPodMoveComplete() : bool

PodDevice

+ResourceName: string

+DeviceIds: []string

NpuBase

-productInfo: *ProductBase

-eidPortMap: map<string,[]string>

-portMapMutex: RWMutex

-urmaDevInfoMap: map<int32,[]UrmaDeviceInfo>

+SetUrmaDeviceInfoByHdm()

+getRankLevelInfoKeyArr() : []string

+GetPortListByEid()([]string, error)

ProductBase

+superPodSize: uint32

+superPodID: uint32

+serverIndex: uint32

+chassisID: uint32

+superPodType: uint8

+nodeInternalIP: string

+cardType: string

+topoInfo: *TopoInfo

+getID(level) : string

+isPodScene() : bool

+isServer() : bool

+isSuperServer() : bool

+isStandCard() : bool

+getTopoFileInfo()(*TopoInfo, error)

TopoInfo

+Version: string

+HardwareType: string

+PeerCount: int

+PeerList: []Peer

+EdgeList: []Edge

2.2 核心方法清单

HwDevManager 核心方法
方法 文件 作用
NewHwDevManager manager.go 构造函数,初始化设备管理器
setAscendManager manager.go 根据设备类型设置运行模式和 manager
setAllDeviceAndType manager.go 初始化 K8s 客户端并获取所有 NPU 设备
UpdateNode / updateNode manager.go 更新节点标签和注解
getNewNodeLabel manager.go 计算需要写入的新节点标签
getNewNodeAnnotation manager.go 计算需要写入的新节点注解(含 baseDevInfo、SuperPod 等)
initPluginServer manager.go 为每种设备类型创建 PluginServer 实例
ListenDevice manager.go 主循环:订阅故障、周期性更新设备信息、通知 kubelet
handleDeviceInfoUpdate manager.go 周期设备信息更新核心逻辑
Serve manager.go 管理 gRPC server 的启动与事件处理
notifyToK8s manager.go 设备状态变化后通知 kubelet(通过 PluginServer.Notify)
chipHotReset manager.go 芯片热复位入口(推理卡)
hotReset manager.go 执行芯片复位并轮询等待启动完成
SubscribeFaultEvent manager.go 订阅 NPU 和交换机故障事件
updatePodAnnotation manager.go 更新 Pod 的 AscendReal 注解
useVolcanoNotify manager.go 通过 Volcano ListAndWatch 上报设备信息
getSuperPodInfo manager.go 获取 SuperPod 拓扑信息
ListenDpu dpu.go 周期查询 DPU operstate
updateDpuHealthy dpu.go 更新 DPU 健康状态到 groupDevice
getLevelList manager_v2.go 生成 A5 Rank Table LevelList
getROCEAddrList manager_v2.go 获取 RoCE 地址列表
PluginServer 核心方法
方法 文件 作用
ListAndWatch plugin.go gRPC 接口:向 kubelet 持续上报设备列表
Allocate plugin.go gRPC 接口:为 Pod 分配设备并构造挂载响应
Notify plugin.go 设备状态变更通知,触发 ListAndWatch 推送
responseToKubelet plugin.go 构造 ListAndWatch 响应(含 volcano/非 volcano/动态切分三种模式)
checkAllocateRequest plugin.go 校验 Allocate 请求合法性
useVolcano / doWithVolcanoSchedule plugin.go Volcano 调度场景下的设备分配
GetKltAndRealAllocateDev plugin.go 获取 kubelet 分配设备与实际设备的映射
updateAllocMap plugin.go 更新 klt→real device 映射表
DestroyNotUsedVNPU plugin.go 销毁未使用的虚拟 NPU
setNPUDeviceMount plugin.go 设置设备挂载(ascend-docker 或原始模式)
setHcclTopoFilePathEnv plugin_v2.go 为 A5 设置 HCCL 拓扑文件路径环境变量
Start / Stop server.go gRPC Server 启停与 kubelet 注册
serve server.go 创建 gRPC server 并监听 unix socket
register server.go 向 kubelet 注册 Device Plugin
PodResource 核心方法
方法 文件 作用
GetPodResource pod_resource.go 调用 kubelet Pod Resources API 获取 Pod 设备分配
IsPodMoveComplete pod_resource.go 检查故障设备上的 Pod 是否已完成迁移
getContainerResource pod_resource.go 从容器资源中提取设备 ID 列表
NpuBase / ProductBase 核心方法
方法 文件 作用
getID npu_base_v2.go 根据 rank level 返回网络实例 ID
getRankLevelInfoKeyArr npu_base_v2.go 获取各层级的网络类型数组(UB/UBG/UBoE/RoCE)
getNetTypeAndFeIDListByRankLevel npu_base_v2.go 根据层级获取网络类型和功能实体 ID 列表
getRandAddrByFuncEntityID npu_base_v2.go 根据 phyID/feID/层级获取 RankAddr 列表
GetPortListByEid npu_base_v2.go 根据 EID 获取端口列表(含缓存)
getTopoFileInfo npu_base_v2.go 读取并解析拓扑文件
SetUrmaDeviceInfoByHdm npu_base_v2.go 通过 DCMI 获取 URMA 设备 EID 列表

2.3 内部调用关系

ticker

triggerTicker

reciChan

PollImmediate

NewHwDevManager()

setAscendManager()

setAllDeviceAndType()

checkSupportedProductType()

setSuperPodInfo()

UpdateNode() → updateNode()

initPluginServer()

NewPluginServer()

ListenDevice(ctx) 主循环

subscribeFaultEvent()

loadFaultCodeAndDeviceInfoCm()

Serve(ctx)

updateNodeAnnotations()

handleDeviceInfoUpdate()

parseTriggers()

updateAllInfo()

mendSubscribeFaultEvents()

updatePodAnnotation()

notifyToK8s()

checkNodeResetInfo()

useVolcanoNotify()

chipHotReset()

manager.UpdateHealth()

graceTolerance()

pluginNotify()

PluginServer.Notify()

PluginServer.ListAndWatch()

startAllServer()

PluginServer.Start()

serve() gRPC

register() kubelet

resetCommonInferCard()

resetDuoCard()

ResetWithoutHccsServer()

ResetHccsServer()

ResetServerForA3()

hotReset()

execResetChip()

GetDeviceBootStatus()

2.4 数据流入流出方式

数据输出

pkg/server 处理

数据输入

DCMI 接口
(设备发现/故障/复位)

K8s API Server
(Node/Pod/ConfigMap)

kubelet gRPC
(ListAndWatch/Allocate)

kubelet Pod Resources
(pod-resources.sock)

ethtool
(DPU operstate)

fsnotify
(socket 文件监听)

OS Signal
(SIGTERM/SIGHUP)

faultCode ConfigMap
(故障码配置)

HwDevManager

PluginServer

PodResource

NpuBase / ProductBase

kubelet
(设备列表/分配响应)

Node Labels
(芯片名/服务器类型/拓扑)

Node Annotations
(baseDevInfo/SuperPod/Rack)

Pod Annotations
(AscendReal)

设备复位
(SetDeviceReset)

Volcano ListAndWatch
(设备健康/内存)

DPU ConfigMap
(operstate)

容器环境变量
(HCCL_TOPO_FILE_PATH)

NPU配置文件
(soft-share)


三、核心业务逻辑深度解析


3.1 manager.go — HwDevManager 设备管理器核心

3.1.1 结构体定义
// HwDevManager 管理华为昇腾设备
type HwDevManager struct {
    SwitchDevManager *deviceswitch.SwitchDevManager  // 交换机设备管理器(A3 场景)
    groupDevice      map[string][]*common.NpuDevice  // 按设备类型分组的设备列表
    ServerMap        map[string]InterfaceServer      // 设备类型 → PluginServer 映射
    allInfo          common.NpuAllInfo                // 所有 NPU 设备信息
    manager          device.DevManager                // 设备管理器接口(910/310/310P)
    RunMode          string                           // 运行模式:Ascend310 / Ascend910 / Ascend310P
    WorkMode         string                           // 工作模式(910 的 SMP/AMP)
    baseNPUInfo      map[string]*common.NpuBaseInfo   // NPU 基础信息缓存(用于注解比对)
    dpuManager       *dpucontrol.DpuFilter            // DPU 过滤器(A5 场景)
    ManagerLock      sync.Mutex                       // 管理器锁
    ContainerRuntime string                           // 容器运行时(docker/containerd)
}
// shareDevResourceQuota 软共享设备资源配额
type shareDevResourceQuota struct {
    aicoreQuota      int  // AI Core 配额
    hbmQuota         int  // HBM 内存配额
    schedulingPolicy int  // 调度策略
}
3.1.2 NewHwDevManager 构造函数

NewHwDevManager(devM)

初始化 dpuManager

setAscendManager(devM)
根据设备类型设置 manager

setAllDeviceAndType()
初始化 K8s 客户端 + 获取所有 NPU

device.InitResetInfoMgr()
初始化复位信息管理器

checkSupportedProductType()
检查产品类型是否受支持

setSuperPodInfo()
获取并缓存 SuperPod 信息

UpdateNode()
更新节点标签和注解

删除旧资源名
(如 Huawei.com/Ascend910)

initPluginServer()
为每种设备类型创建 PluginServer

获取容器运行时

返回 *HwDevManager

逐行解析:

func NewHwDevManager(devM devmanager.DeviceInterface) *HwDevManager {
    var hdm HwDevManager
    // 初始化 DPU 过滤器实例,用于 A5 场景的 DPU 信息管理
    hdm.dpuManager = &dpucontrol.DpuFilter{}

    // 步骤1:根据设备接口类型,设置对应的 manager(910/310/310P)
    if err := hdm.setAscendManager(devM); err != nil {
        hwlog.RunLog.Errorf("init hw dev manager failed, err: %v", err)
        return nil  // 初始化失败返回 nil,调用方需检查
    }

    // 步骤2:创建 K8s 客户端,获取所有 NPU 设备并分类
    if err := hdm.setAllDeviceAndType(); err != nil {
        hwlog.RunLog.Errorf("set all device and type failed, err: %v", err)
        return nil
    }

    // 步骤3:初始化复位信息管理器,用于记录设备复位状态
    device.InitResetInfoMgr(hdm.manager.GetKubeClient())

    // 步骤4:检查产品类型是否支持(如 Atlas 300I Duo 不支持动态虚拟化)
    if err := hdm.checkSupportedProductType(); err != nil {
        hwlog.RunLog.Errorf("check supported product type failed, err: %v", err)
        return nil
    }

    // 步骤5:获取并缓存 SuperPod 信息(SuperPod ID、Server Index、Rack ID)
    hdm.setSuperPodInfo()

    // 步骤6:更新 K8s 节点标签和注解(芯片名、服务器类型、拓扑信息等)
    if err := hdm.UpdateNode(); err != nil {
        hwlog.RunLog.Errorf("update node label failed, err: %v", err)
        return nil
    }

    // 步骤7:如果设备类型从旧名称迁移到新名称(如 Ascend910 → huawei.com/npu),
    //         则删除旧资源名,避免 kubelet 上残留旧资源
    kubeClient := hdm.manager.GetKubeClient()
    if kubeClient != nil {
        deviceType := hdm.manager.GetDmgr().GetDevType()
        if !customname.IsOldDeviceType(deviceType) {
            hwlog.RunLog.Info("current device type changes to Huawei.com/npu, delete old resource name")
            err := kubeClient.RemoveOldResource(api.HuaweiAscend910)
            if err != nil {
                hwlog.RunLog.Errorf("failed to delete old resource name: %v", err)
                return nil
            }
        }
    }

    // 步骤8:为每种设备类型创建 PluginServer 实例
    if err := hdm.initPluginServer(); err != nil {
        hwlog.RunLog.Errorf("init plugin server failed, err: %v", err)
        return nil
    }

    // 步骤9:获取容器运行时类型(docker/containerd)
    if runtime, err := hdm.manager.GetKubeClient().GetContainerRuntime(); err == nil {
        hdm.ContainerRuntime = runtime
    }

    return &hdm
}
3.1.3 setAscendManager — 设备类型路由
func (hdm *HwDevManager) setAscendManager(dmgr devmanager.DeviceInterface) error {
    devType := dmgr.GetDevType()

    // PresetVDevice=false 仅支持 310P 和 910B(动态切分模式)
    if !common.ParamOption.PresetVDevice && devType != api.Ascend310P && devType != api.Ascend910B {
        return fmt.Errorf("only 310p and 910b support to set presetVirtualDevice false")
    }
    common.ParamOption.RealCardType = devType  // 缓存真实卡类型到全局参数

    // 根据设备类型选择对应的 manager 实现
    switch devType {
    case api.Ascend310, api.Ascend310B:
        hdm.RunMode = api.Ascend310
        hdm.manager = device.NewHwAscend310Manager()
    case api.Ascend910A, api.Ascend910B, api.Ascend910A3, api.Ascend910A5:
        hdm.RunMode = api.Ascend910
        hdm.manager = device.NewHwAscend910Manager()
        hdm.WorkMode = dmgr.GetNpuWorkMode()  // 910 需要额外的工作模式
    case api.Ascend310P:
        hdm.RunMode = api.Ascend310P
        hdm.manager = device.NewHwAscend310PManager()
    default:
        return fmt.Errorf("an unsupported device type")
    }

    hdm.manager.SetDmgr(dmgr)  // 注入设备管理接口

    // 获取所有产品类型并缓存
    productTypes, err := hdm.manager.GetDmgr().GetAllProductType()
    if err != nil {
        return err
    }
    common.ParamOption.ProductTypes = productTypes

    // 检查 310P 混插模式是否合法
    if err = common.CheckCardUsageMode(common.ParamOption.Use310PMixedInsert, productTypes); err != nil {
        return err
    }

    // 非边缘场景需要获取芯片 AI Core 数量
    if common.ParamOption.BuildScene != common.EdgeScene {
        aiCoreCount, err := hdm.manager.GetChipAiCoreCount()
        if err != nil {
            return err
        }
        common.ParamOption.AiCoreCount = aiCoreCount
    }
    return nil
}
3.1.4 ListenDevice — 主循环

这是整个 device plugin 的 核心运行循环,启动后持续运行直到收到停止信号。

ctx.Done

triggerTicker.C 每秒

ticker.C 每 ListAndWatchPeriod 秒

ListenDevice(ctx)

subscribeFaultEvent()
订阅 NPU + 交换机故障事件

如果 A3 + EnableSwitchFault
启动交换机故障轮询 goroutine

loadFaultCodeAndDeviceInfoCm(ctx)
加载故障码 + 设备信息 ConfigMap

go Serve(ctx)
启动 gRPC Server 管理 goroutine

如果 CheckCachedPods
启动 PodInformerInspector

go updateNodeAnnotations(ctx)
定期更新节点注解

go WriteFaultToEvent(ctx)
设备故障写入 K8s Event

initTime = time.Now()

主循环

return 停止

parseTriggers(ctx, initTime)
检查更新触发通道

handleDeviceInfoUpdate(ctx, initTime)
周期设备信息更新

逐行解析关键部分:

func (hdm *HwDevManager) ListenDevice(ctx context.Context) {
    hwlog.RunLog.Info("starting the listen device")

    // 1. 订阅故障事件(NPU + 交换机)
    hdm.subscribeFaultEvent()

    // 2. 如果是 A3 设备且启用了交换机故障检测,启动定期查询 goroutine
    if common.ParamOption.RealCardType == api.Ascend910A3 && common.ParamOption.EnableSwitchFault {
        go hdm.SwitchDevManager.GetSwitchFaultCodeByInterval(ctx,
            time.Second*common.GetSwitchFaultCodeInterval)
    }

    // 3. 加载故障码和设备信息 ConfigMap(阻塞,完成后启动后台轮询)
    hdm.loadFaultCodeAndDeviceInfoCm(ctx)

    // 4. 启动 gRPC Server 管理 goroutine(处理 socket 文件监听、重启等)
    go hdm.Serve(ctx)

    // 5. 如果需要检查缓存的 Pod,启动 Pod Informer 检查器
    if common.ParamOption.CheckCachedPods {
        go hdm.manager.GetKubeClient().PodInformerInspector(ctx)
    }

    // 6. 启动节点注解定期更新 goroutine(每 60 秒检查一次)
    go hdm.updateNodeAnnotations(ctx)

    // 7. 启动设备故障写入 K8s Event goroutine
    go hdm.manager.WriteFaultToEvent(ctx)

    // 8. 主循环:周期性设备信息更新 + 触发器检查
    initTime := time.Now()
    ticker := time.NewTicker(time.Duration(common.ParamOption.ListAndWatchPeriod) * time.Second)
    defer ticker.Stop()
    triggerTicker := time.NewTicker(time.Second)  // 每秒检查触发器
    defer triggerTicker.Stop()

    for {
        select {
        case _, ok := <-ctx.Done():
            // 收到上下文取消信号,退出主循环
            hwlog.RunLog.Info("listen device stop")
            return
        case <-triggerTicker.C:\n            // 每秒检查是否有更新触发信号(如故障事件触发的即时更新)\n            hdm.parseTriggers(ctx, initTime)\n        case <-ticker.C:\n            // 周期性设备信息更新(默认 5 秒)\n            hwlog.RunLog.Debug("Periodic device info update")
            hdm.handleDeviceInfoUpdate(ctx, &initTime)
        }
    }
}
3.1.5 handleDeviceInfoUpdate — 周期设备信息更新

handleDeviceInfoUpdate(ctx, initTime)

common.LockAllDeviceInfo()
加全局设备信息锁

updateAllInfo()
更新所有设备信息

mendSubscribeFaultEvents()
补充订阅方式无法上报的故障

updatePodAnnotation()
更新 Pod AscendReal 注解

updateDeviceUsedInfo()
标记设备是否被 Pod 使用

notifyToK8s(ctx, initTime)
通知 kubelet 设备状态变化

checkNodeResetInfo()
检查并清理已恢复设备的复位信息

useVolcanoNotify()
Volcano ListAndWatch 上报

chipHotReset()
推理卡热复位检查

DelOnceRecoverFault()
删除已恢复的故障记录

common.Synchronize = true
标记同步完成

common.UnlockAllDeviceInfo()

逐行解析:

func (hdm *HwDevManager) handleDeviceInfoUpdate(ctx context.Context, initTime *time.Time) {
    // 加全局锁,防止并发更新设备信息
    common.LockAllDeviceInfo()
    defer common.UnlockAllDeviceInfo()

    // 1. 更新所有设备信息(动态切分场景需先销毁未用 vNPU,重新获取设备列表)
    if err := hdm.updateAllInfo(); err != nil {
        hwlog.RunLog.Error(err)
        return
    }

    // 2. 补充订阅方式无法上报的故障码(如断卡、丢芯片等)
    hdm.mendSubscribeFaultEvents()

    // 3. 更新 Pod 注解(写入 AscendReal 真实分配设备信息)
    if err := hdm.updatePodAnnotation(); err != nil {
        hwlog.RunLog.Error(err)
    }

    // 4. 标记设备是否被 Pod 使用(从 kubelet 获取已分配设备列表)
    hdm.updateDeviceUsedInfo(hdm.groupDevice)

    // 5. 通知 kubelet 设备状态变化(触发 ListAndWatch 推送)
    hdm.notifyToK8s(ctx, initTime)

    // 6. 检查节点复位信息:如果 annotation 中有复位失败设备但设备已恢复,清除注解
    hdm.checkNodeResetInfo()

    // 7. Volcano ListAndWatch 上报(设备健康状态、内存信息)
    hdm.useVolcanoNotify()

    // 8. 推理卡热复位检查(仅 HotResetInfer 模式)
    hdm.chipHotReset()

    // 9. 删除已恢复的故障记录和频率故障记录
    common.DelOnceRecoverFault(hdm.groupDevice)
    common.DelOnceFrequencyFault()

    // 10. 标记同步完成,允许其他操作读取最新设备信息
    common.Synchronize = true
}
3.1.6 notifyToK8s — 通知 kubelet 设备状态
func (hdm *HwDevManager) notifyToK8s(ctx context.Context, initTime *time.Time) {
    // 检查是否支持训练容错
    hdm.isSupportGraceTolerance()

    // 深拷贝旧设备状态(用于变更检测)
    oldGroupDevice := deepCopyGroupDevice(hdm.groupDevice)

    // 更新设备健康状态(根据故障码判断 healthy/unhealthy)
    hdm.manager.UpdateHealth(hdm.groupDevice, hdm.allInfo.AICoreDevs, hdm.RunMode)

    // A5 设备:更新 DPU 健康状态
    if hdm.manager.GetDmgr().GetDevType() == api.Ascend910A5 {
        hdm.updateDpuHealthy(hdm.groupDevice)
    }

    // 训练容错处理(热复位场景下,正在复位的设备标记为 healthy)
    hdm.graceTolerance(ctx, hdm.groupDevice)

    // 检测设备状态是否变化
    isDevStateChange := hdm.manager.GetChange(hdm.groupDevice, oldGroupDevice)

    // 对每种设备类型,如果状态变化或需要重发,通知 PluginServer
    for devType, isChanged := range isDevStateChange {
        server := hdm.ServerMap[devType]
        if server == nil {
            continue
        }
        // 如果状态未变化,且满足以下条件之一则跳过:
        // - 启动 1 分钟内且上次发送成功,且启动 1 小时内
        if !isChanged &&
            (time.Now().Sub(*initTime) < time.Minute || server.LastSendSuccess()) &&
            time.Now().Sub(*initTime) < time.Hour {
            continue
        }
        *initTime = time.Now()

        // 非预置虚拟设备:使用 AICore 设备列表通知
        if !common.ParamOption.PresetVDevice {
            hdm.pluginNotify(hdm.allInfo.AICoreDevs, common.AiCoreResourceName)
            return
        }
        // 预置虚拟设备:按设备类型分别通知
        hdm.pluginNotify(hdm.groupDevice[devType], devType)
    }
}
3.1.7 hotReset — 芯片热复位

error

success

未启动完成

全部 BootStartFinish

hotReset(device, devices)

SetCardsInResetting(true)
标记设备正在复位

wait.PollImmediate(1s, 1min)

execResetChip()
执行复位

复位失败

遍历 devices
检查 GetDeviceBootStatus()

SetDeviceInit()
初始化设备状态

SetResetFailedTimes(0)
SetCardsInResetting(false)
复位成功

SetCardsInResetting(false)
SetResetFailedTimes(+1)

func (hdm *HwDevManager) hotReset(device *common.NpuDevice, devices []*common.NpuDevice) {
    hwlog.RunLog.Infof("will start to reset device %s", device.DeviceName)

    // 标记设备正在复位,防止重复复位
    hdm.manager.SetCardsInResetting(device.LogicID, true)

    var isResetExec = false
    successResetDevList := sets.NewInt32()

    // 轮询等待:每秒检查一次,最多等待 1 分钟
    if err := wait.PollImmediate(time.Second, time.Minute, func() (bool, error) {
        // 执行复位操作(只执行一次,通过 isResetExec 标志控制)
        if err := hdm.execResetChip(device.LogicID, &isResetExec); err != nil {
            return false, err
        }
        // 检查所有关联设备的启动状态
        for _, dev := range devices {
            if successResetDevList.Has(dev.LogicID) {
                continue  // 已确认启动完成的设备跳过
            }
            bootState, err := hdm.manager.GetDmgr().GetDeviceBootStatus(dev.LogicID)
            if err != nil {
                return false, err
            }
            if bootState != common.BootStartFinish {
                // 设备尚未启动完成,继续等待
                return false, nil
            }
            successResetDevList.Insert(dev.LogicID)
        }
        // 所有设备都启动完成
        common.SetDeviceInit(device.LogicID)
        return true, nil
    }); err != nil {
        // 超时或出错
        hwlog.RunLog.Warnf("hot reset failed, timeout or err: %v", err)
        hdm.manager.SetCardsInResetting(device.LogicID, false)
        hdm.manager.SetResetFailedTimes(device.LogicID, hdm.manager.GetResetFailedTimes(device.LogicID)+1)
        return
    }

    // 复位成功:清除复位失败计数,清除复位中标记
    hdm.manager.SetResetFailedTimes(device.LogicID, 0)
    hdm.manager.SetCardsInResetting(device.LogicID, false)
    hwlog.RunLog.Info("hot reset success")
}
3.1.8 Serve — gRPC Server 管理
func (hdm *HwDevManager) Serve(ctx context.Context) {
    hwlog.RunLog.Info("Serve start")

    // 创建 socket 文件监听器(监听 /var/lib/kubelet/device-plugin/ 目录)
    watcher, err := common.NewFileWatch()
    if err != nil {
        hwlog.RunLog.Error("createSocketWatcher error")
        return
    }
    defer func() {
        if watcher != nil {
            if err := watcher.FileWatcher.Close(); err != nil {
                hwlog.RunLog.Errorf("close file watcher, err: %v", err)
            }
        }
    }()

    // 创建重启信号监听(SIGHUP)
    restartSignal := common.NewSignWatcher(syscall.SIGHUP)

    for {
        // 启动所有标记需要重启的 PluginServer
        allSuccess := hdm.startAllServer(watcher)

        // 等待事件(上下文取消、重启信号、socket 文件变更)
        if hdm.handleEvents(ctx, restartSignal, watcher) {
            break  // 收到停止信号,退出循环
        }

        // 如果启动失败,等待一段时间后重试
        if !allSuccess {
            time.Sleep(common.SleepTime * time.Second)
        }
    }
}

handleEvents 处理三种事件:

  • ctx.Done():收到停止信号,返回 true 退出
  • restartSignal (SIGHUP):设置所有 server 重启标志
  • fsnotify 事件:socket 文件被删除时打印告警;kubelet.sock 创建时记录日志
3.1.9 updateNode — 节点标签与注解更新

Yes

No

失败

Yes

No

成功

UpdateNode()

BuildScene == EdgeScene?

return nil

InitPodInformer()

updateNode()

GetNode() 获取当前节点

SetNodeInternalIPInK8s()
缓存节点内部 IP

getNewNodeLabel()
计算新标签

getNewNodeAnnotation()
计算新注解

PatchNodeState()
打补丁到 K8s

重试 < RetryUpdateCount?

Sleep 1s → 重新 Patch

return error

return nil

getNewNodeLabel 关键逻辑:

  • updateChipNameToNode():获取芯片名称写入标签
  • 服务器类型标签:Ascend910-32ascend-32(新命名)
  • HBM 内存标签:NPUChipMemoryLabel(如 32G
  • A300IA2 推理卡标签:当 910B + 推理模式 + 特定 BoardId 时添加
  • 拓扑标签:A3/A5 添加 SuperPodId,A5 额外添加 RackId

getNewNodeAnnotation 关键逻辑:

  • cardType:板卡类型(如 A5300I)
  • BaseDevInfoAnno:JSON 序列化的 NPU 基础信息(设备名→IP/DeviceID/SuperDeviceID/LevelList)
  • SuperPodIDKey / serverIndexKey / serverTypeKey / RackIDKey:SuperPod 拓扑信息
3.1.10 故障码 ConfigMap 管理

interval 秒

error

success

No

Yes

Yes

No

loadFaultCodeAndDeviceInfoCm(ctx)

loadFaultCode()
从 CM 加载故障码,返回轮询间隔

LoadDeviceInfoCm(ctx)
加载设备信息 ConfigMap

loadDeviceFaultFromUpgradeReason()
从升级故障原因加载故障

go pollFaultCodeCM(ctx, interval)
后台轮询故障码 CM

UpdateHealth()
根据故障更新健康状态

SetDeviceInit()
标记设备已初始化

pollFaultCodeCM

loadFaultCode()
重新加载

loadFaultCode()

GetConfigMap(FaultCodeCMName)

initFaultInfoFromFile()
从本地文件加载

updateFaultConfigFromCm()

resourceVersion 变化?

return(无变化)

loadFaultCode()
加载故障码

A3 + EnableSwitchFault?

loadSwitchFaultCode()
加载交换机故障码

loadFaultCustomization()
加载故障自定义配置

3.1.11 chipHotReset — 推理卡热复位

No

Yes

Yes

No

Yes

No

Yes

No

No

Yes

无 HCCS

有 HCCS

Yes

No

No

Yes

chipHotReset()

HotReset == HotResetInfer?

return

NewPodResource()

遍历 groupDevice

IsVirtualDev?

IsContainAtlas300IDuo?

resetDuoCard()

resetCommonInferCard()

RealCardType == A3?

ResetServerForA3()

getServerUsageAndBoardId()

usage == Infer?

逐设备检查
isPodRemove + checkNoProc
→ hotReset()

boardId 类型判断

ResetWithoutHccsServer()
逐卡复位

ResetHccsServer()
整组复位

按 CardID 分组

isDuoCardChipHealthy()?

下一张卡

isDuoRemove()?
检查所有芯片 Pod 是否迁移

hotReset()


Logo

鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。

更多推荐